Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 2 additions & 1 deletion src/java/org/apache/nutch/crawl/CrawlDb.java
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,7 @@
import org.apache.hadoop.util.Tool;
import org.apache.hadoop.util.ToolRunner;
import org.apache.nutch.metadata.Nutch;
import org.apache.nutch.metrics.NutchMetrics;
import org.apache.nutch.util.FSUtils;
import org.apache.nutch.util.HadoopFSUtil;
import org.apache.nutch.util.LockUtil;
Expand Down Expand Up @@ -145,7 +146,7 @@ public void update(Path crawlDb, Path[] segments, boolean normalize,

if (filter) {
long urlsFiltered = job.getCounters()
.findCounter("CrawlDB filter", "URLs filtered").getValue();
.findCounter(NutchMetrics.GROUP_CRAWLDB_FILTER, NutchMetrics.CRAWLDB_URLS_FILTERED_TOTAL).getValue();
LOG.info(
"CrawlDb update: Total number of existing URLs in CrawlDb rejected by URL filters: {}",
urlsFiltered);
Expand Down
10 changes: 4 additions & 6 deletions src/java/org/apache/nutch/crawl/DeduplicationJob.java
Original file line number Diff line number Diff line change
Expand Up @@ -335,12 +335,10 @@ public int run(String[] args) throws IOException {
fs.delete(tempDir, true);
throw new RuntimeException(message);
}
CounterGroup g = job.getCounters().getGroup("DeduplicationJobStatus");
if (g != null) {
Counter counter = g.findCounter("Documents marked as duplicate");
long dups = counter.getValue();
LOG.info("Deduplication: {} documents marked as duplicates", dups);
}
long dups = job.getCounters()
.findCounter(NutchMetrics.GROUP_DEDUP, NutchMetrics.DEDUP_DOCUMENTS_MARKED_DUPLICATE_TOTAL)
.getValue();
LOG.info("Deduplication: {} documents marked as duplicates", dups);
} catch (IOException | InterruptedException | ClassNotFoundException e) {
LOG.error("DeduplicationJob:", e);
fs.delete(tempDir, true);
Expand Down
34 changes: 26 additions & 8 deletions src/java/org/apache/nutch/fetcher/QueueFeeder.java
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@

import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.mapreduce.Counter;
import org.apache.nutch.crawl.CrawlDatum;
import org.apache.nutch.fetcher.FetchItemQueues.QueuingStatus;
import org.apache.nutch.fetcher.Fetcher.FetcherRun;
Expand All @@ -48,6 +49,12 @@ public class QueueFeeder extends Thread {
private URLNormalizers urlNormalizers = null;
private String urlNormalizerScope = URLNormalizers.SCOPE_DEFAULT;

// Cached counter references to avoid repeated lookups in hot paths
private Counter hitByTimeoutCounter;
private Counter hitByTimelimitCounter;
private Counter filteredCounter;
private Counter aboveExceptionThresholdCounter;

public QueueFeeder(FetcherRun.Context context,
FetchItemQueues queues, int size) {
this.context = context;
Expand All @@ -62,6 +69,21 @@ public QueueFeeder(FetcherRun.Context context,
if (conf.getBoolean("fetcher.normalize.urls", false)) {
urlNormalizers = new URLNormalizers(conf, urlNormalizerScope);
}
initCounters();
}

/**
* Initialize cached counter references to avoid repeated lookups in hot paths.
*/
private void initCounters() {
hitByTimeoutCounter = context.getCounter(
NutchMetrics.GROUP_FETCHER, NutchMetrics.FETCHER_HIT_BY_TIMEOUT_TOTAL);
hitByTimelimitCounter = context.getCounter(
NutchMetrics.GROUP_FETCHER, NutchMetrics.FETCHER_HIT_BY_TIMELIMIT_TOTAL);
filteredCounter = context.getCounter(
NutchMetrics.GROUP_FETCHER, NutchMetrics.FETCHER_FILTERED_TOTAL);
aboveExceptionThresholdCounter = context.getCounter(
NutchMetrics.GROUP_FETCHER, NutchMetrics.FETCHER_ABOVE_EXCEPTION_THRESHOLD_TOTAL);
}

/** Filter and normalize the url */
Expand Down Expand Up @@ -95,16 +117,14 @@ public void run() {
LOG.info("QueueFeeder stopping, timeout reached.");
}
queuingStatus[qstatus]++;
context.getCounter(NutchMetrics.GROUP_FETCHER,
NutchMetrics.FETCHER_HIT_BY_TIMEOUT_TOTAL).increment(1);
hitByTimeoutCounter.increment(1);
} else {
int qstatus = QueuingStatus.HIT_BY_TIMELIMIT.ordinal();
if (queuingStatus[qstatus] == 0) {
LOG.info("QueueFeeder stopping, timelimit exceeded.");
}
queuingStatus[qstatus]++;
context.getCounter(NutchMetrics.GROUP_FETCHER,
NutchMetrics.FETCHER_HIT_BY_TIMELIMIT_TOTAL).increment(1);
hitByTimelimitCounter.increment(1);
}
try {
hasMore = context.nextKeyValue();
Expand Down Expand Up @@ -136,8 +156,7 @@ public void run() {
String u = filterNormalize(url.toString());
if (u == null) {
// filtered or failed to normalize
context.getCounter(NutchMetrics.GROUP_FETCHER,
NutchMetrics.FETCHER_FILTERED_TOTAL).increment(1);
filteredCounter.increment(1);
continue;
}
url = new Text(u);
Expand All @@ -154,8 +173,7 @@ public void run() {
QueuingStatus status = queues.addFetchItem(url, datum);
queuingStatus[status.ordinal()]++;
if (status == QueuingStatus.ABOVE_EXCEPTION_THRESHOLD) {
context.getCounter(NutchMetrics.GROUP_FETCHER,
NutchMetrics.FETCHER_ABOVE_EXCEPTION_THRESHOLD_TOTAL).increment(1);
aboveExceptionThresholdCounter.increment(1);
}
cnt++;
feed--;
Expand Down
23 changes: 15 additions & 8 deletions src/java/org/apache/nutch/hostdb/UpdateHostDbMapper.java
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@
import org.apache.hadoop.io.FloatWritable;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.io.Writable;
import org.apache.hadoop.mapreduce.Counter;
import org.apache.hadoop.mapreduce.Mapper;
import org.apache.hadoop.conf.Configuration;

Expand Down Expand Up @@ -61,6 +62,10 @@ public class UpdateHostDbMapper
protected URLFilters filters = null;
protected URLNormalizers normalizers = null;

// Cached counter references to avoid repeated lookups in hot paths
protected Counter malformedUrlCounter;
protected Counter filteredRecordsCounter;

@Override
public void setup(Mapper<Text, Writable, Text, NutchWritable>.Context context) {
Configuration conf = context.getConfiguration();
Expand All @@ -72,6 +77,12 @@ public void setup(Mapper<Text, Writable, Text, NutchWritable>.Context context) {
filters = new URLFilters(conf);
if (normalize)
normalizers = new URLNormalizers(conf, URLNormalizers.SCOPE_DEFAULT);

// Initialize cached counter references
malformedUrlCounter = context.getCounter(
NutchMetrics.GROUP_HOSTDB, NutchMetrics.HOSTDB_MALFORMED_URL_TOTAL);
filteredRecordsCounter = context.getCounter(
NutchMetrics.GROUP_HOSTDB, NutchMetrics.HOSTDB_FILTERED_RECORDS_TOTAL);
}

/**
Expand Down Expand Up @@ -137,8 +148,7 @@ public void map(Text key, Writable value,
try {
url = new URL(keyStr);
} catch (MalformedURLException e) {
context.getCounter(NutchMetrics.GROUP_HOSTDB,
NutchMetrics.HOSTDB_MALFORMED_URL_TOTAL).increment(1);
malformedUrlCounter.increment(1);
return;
}
String hostName = URLUtil.getHost(url);
Expand All @@ -148,8 +158,7 @@ public void map(Text key, Writable value,

// Filtered out?
if (buffer == null) {
context.getCounter(NutchMetrics.GROUP_HOSTDB,
NutchMetrics.HOSTDB_FILTERED_RECORDS_TOTAL).increment(1);
filteredRecordsCounter.increment(1);
LOG.debug("UpdateHostDb: {} crawldatum has been filtered", hostName);
return;
}
Expand Down Expand Up @@ -222,8 +231,7 @@ public void map(Text key, Writable value,

// Filtered out?
if (buffer == null) {
context.getCounter(NutchMetrics.GROUP_HOSTDB,
NutchMetrics.HOSTDB_FILTERED_RECORDS_TOTAL).increment(1);
filteredRecordsCounter.increment(1);
LOG.debug("UpdateHostDb: {} hostdatum has been filtered", keyStr);
return;
}
Expand All @@ -247,8 +255,7 @@ public void map(Text key, Writable value,

// Filtered out?
if (buffer == null) {
context.getCounter(NutchMetrics.GROUP_HOSTDB,
NutchMetrics.HOSTDB_FILTERED_RECORDS_TOTAL).increment(1);
filteredRecordsCounter.increment(1);
LOG.debug("UpdateHostDb: {} score has been filtered", keyStr);
return;
}
Expand Down
23 changes: 17 additions & 6 deletions src/java/org/apache/nutch/hostdb/UpdateHostDbReducer.java
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,7 @@
import org.apache.hadoop.io.LongWritable;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.io.Writable;
import org.apache.hadoop.mapreduce.Counter;
import org.apache.hadoop.mapreduce.Reducer;
import org.apache.hadoop.util.StringUtils;

Expand Down Expand Up @@ -73,6 +74,11 @@ public class UpdateHostDbReducer
protected BlockingQueue<Runnable> queue = new SynchronousQueue<>();
protected ThreadPoolExecutor executor = null;

// Cached counter references to avoid repeated lookups in hot paths
protected Counter urlLimitNotReachedCounter;
protected Counter totalHostsCounter;
protected Counter skippedNotEligibleCounter;

/**
* Configures the thread pool and prestarts all resolver threads.
*/
Expand Down Expand Up @@ -146,6 +152,14 @@ public void setup(Reducer<Text, NutchWritable, Text, HostDatum>.Context context)
// Run all threads in the pool
executor.prestartAllCoreThreads();
}

// Initialize cached counter references
urlLimitNotReachedCounter = context.getCounter(
NutchMetrics.GROUP_HOSTDB, NutchMetrics.HOSTDB_URL_LIMIT_NOT_REACHED_TOTAL);
totalHostsCounter = context.getCounter(
NutchMetrics.GROUP_HOSTDB, NutchMetrics.HOSTDB_TOTAL_HOSTS_TOTAL);
skippedNotEligibleCounter = context.getCounter(
NutchMetrics.GROUP_HOSTDB, NutchMetrics.HOSTDB_SKIPPED_NOT_ELIGIBLE_TOTAL);
}

/**
Expand Down Expand Up @@ -380,14 +394,12 @@ else if (value instanceof FloatWritable) {
// Impose limits on minimum number of URLs?
if (urlLimit > -1l) {
if (hostDatum.numRecords() < urlLimit) {
context.getCounter(NutchMetrics.GROUP_HOSTDB,
NutchMetrics.HOSTDB_URL_LIMIT_NOT_REACHED_TOTAL).increment(1);
urlLimitNotReachedCounter.increment(1);
return;
}
}

context.getCounter(NutchMetrics.GROUP_HOSTDB,
NutchMetrics.HOSTDB_TOTAL_HOSTS_TOTAL).increment(1);
totalHostsCounter.increment(1);

// See if this record is to be checked
if (shouldCheck(hostDatum)) {
Expand All @@ -404,8 +416,7 @@ else if (value instanceof FloatWritable) {
// Do not progress, the datum will be written in the resolver thread
return;
} else if (checkAny) {
context.getCounter(NutchMetrics.GROUP_HOSTDB,
NutchMetrics.HOSTDB_SKIPPED_NOT_ELIGIBLE_TOTAL).increment(1);
skippedNotEligibleCounter.increment(1);
LOG.debug("UpdateHostDb: {}: skipped_not_eligible", key);
}

Expand Down
Loading