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: 3 additions & 0 deletions ivy/ivy.xml
Original file line number Diff line number Diff line change
Expand Up @@ -122,6 +122,9 @@
<!-- Required for JUnit 5 (Jupiter) test execution -->
<dependency org="org.junit.jupiter" name="junit-jupiter-engine" rev="5.14.1" conf="test->default"/>
<dependency org="org.junit.jupiter" name="junit-jupiter-api" rev="5.14.1" conf="test->default"/>
<!-- Mockito for mocking in tests -->
<dependency org="org.mockito" name="mockito-core" rev="5.18.0" conf="test->default"/>
<dependency org="org.mockito" name="mockito-junit-jupiter" rev="5.18.0" conf="test->default"/>

<!-- Jetty used to serve test pages for unit tests, but is also provided as dependency of Hadoop -->
<dependency org="org.eclipse.jetty" name="jetty-server" rev="12.1.5" conf="test->default">
Expand Down
7 changes: 7 additions & 0 deletions src/java/org/apache/nutch/crawl/CrawlDbReducer.java
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,7 @@
import org.apache.hadoop.io.Writable;
import org.apache.hadoop.util.PriorityQueue;
import org.apache.nutch.metadata.Nutch;
import org.apache.nutch.metrics.ErrorTracker;
import org.apache.nutch.metrics.NutchMetrics;
import org.apache.nutch.scoring.ScoringFilterException;
import org.apache.nutch.scoring.ScoringFilters;
Expand All @@ -49,6 +50,7 @@ public class CrawlDbReducer extends
private boolean additionsAllowed;
private int maxInterval;
private FetchSchedule schedule;
private ErrorTracker errorTracker;

@Override
public void setup(Reducer<Text, CrawlDatum, Text, CrawlDatum>.Context context) {
Expand All @@ -60,6 +62,8 @@ public void setup(Reducer<Text, CrawlDatum, Text, CrawlDatum>.Context context) {
schedule = FetchScheduleFactory.getFetchSchedule(conf);
int maxLinks = conf.getInt("db.update.max.inlinks", 10000);
linked = new InlinkPriorityQueue(maxLinks);
// Initialize error tracker with cached counters
errorTracker = new ErrorTracker(NutchMetrics.GROUP_CRAWLDB, context);
}

@Override
Expand Down Expand Up @@ -162,6 +166,7 @@ public void reduce(Text key, Iterable<CrawlDatum> values,
scfilters.orphanedScore(key, old);
} catch (ScoringFilterException e) {
LOG.warn("Couldn't update orphaned score, key={}: {}", key, e);
errorTracker.incrementCounters(e);
}
context.write(key, old);
// Dynamic counter based on status name
Expand Down Expand Up @@ -208,6 +213,7 @@ public void reduce(Text key, Iterable<CrawlDatum> values,
} catch (ScoringFilterException e) {
LOG.warn("Cannot filter init score for url {}, using default: {}",
key, e.getMessage());
errorTracker.incrementCounters(e);
result.setScore(0.0f);
}
}
Expand Down Expand Up @@ -317,6 +323,7 @@ public void reduce(Text key, Iterable<CrawlDatum> values,
scfilters.updateDbScore(key, oldSet ? old : null, result, linkList);
} catch (Exception e) {
LOG.warn("Couldn't update score, key={}: {}", key, e);
errorTracker.incrementCounters(e);
}
// remove generation time, if any
result.getMetaData().remove(Nutch.WRITABLE_GENERATE_TIME_KEY);
Expand Down
14 changes: 10 additions & 4 deletions src/java/org/apache/nutch/crawl/Generator.java
Original file line number Diff line number Diff line change
Expand Up @@ -67,6 +67,7 @@
import org.apache.hadoop.io.WritableComparator;
import org.apache.nutch.hostdb.HostDatum;
import org.apache.nutch.metadata.Nutch;
import org.apache.nutch.metrics.ErrorTracker;
import org.apache.nutch.metrics.NutchMetrics;
import org.apache.nutch.net.URLFilterException;
import org.apache.nutch.net.URLFilters;
Expand Down Expand Up @@ -191,6 +192,7 @@ public static class SelectorMapper
private int intervalThreshold = -1;
private byte restrictStatus = -1;
private JexlScript expr = null;
private ErrorTracker errorTracker;

@Override
public void setup(
Expand All @@ -215,6 +217,8 @@ public void setup(
restrictStatus = CrawlDatum.getStatusByName(restrictStatusString);
}
expr = JexlUtil.parseExpression(conf.get(GENERATOR_EXPR, null));
// Initialize error tracker with cached counters
errorTracker = new ErrorTracker(NutchMetrics.GROUP_GENERATOR, context);
}

@Override
Expand All @@ -231,8 +235,7 @@ public void map(Text key, CrawlDatum value, Context context)
return;
}
} catch (URLFilterException e) {
context.getCounter(NutchMetrics.GROUP_GENERATOR,
NutchMetrics.GENERATOR_URL_FILTER_EXCEPTION_TOTAL).increment(1);
errorTracker.incrementCounters(e);
LOG.warn("Couldn't filter url: {} ({})", url, e.getMessage());
}
}
Expand Down Expand Up @@ -261,6 +264,7 @@ public void map(Text key, CrawlDatum value, Context context)
try {
sort = scfilters.generatorSortValue(key, crawlDatum, sort);
} catch (ScoringFilterException sfe) {
errorTracker.incrementCounters(sfe);
LOG.warn("Couldn't filter generatorSortValue for {}: {}", key, sfe);
}

Expand Down Expand Up @@ -326,6 +330,7 @@ public static class SelectorReducer extends
private JexlScript maxCountExpr = null;
private JexlScript fetchDelayExpr = null;
private Map<String, HostDatum> hostDatumCache = new HashMap<>();
private ErrorTracker errorTracker;

public void readHostDb() throws IOException {
if (conf.get(GENERATOR_HOSTDB) == null) {
Expand Down Expand Up @@ -419,6 +424,8 @@ public void setup(Context context) throws IOException {
fetchDelayExpr = JexlUtil
.parseExpression(conf.get(GENERATOR_FETCH_DELAY_EXPR, null));
}
// Initialize error tracker with cached counters
errorTracker = new ErrorTracker(NutchMetrics.GROUP_GENERATOR, context);

readHostDb();
}
Expand Down Expand Up @@ -516,8 +523,7 @@ public void reduce(FloatWritable key, Iterable<SelectorEntry> values,
} catch (MalformedURLException e) {
LOG.warn("Malformed URL: '{}', skipping ({})", urlString,
StringUtils.stringifyException(e));
context.getCounter(NutchMetrics.GROUP_GENERATOR,
NutchMetrics.GENERATOR_MALFORMED_URL_TOTAL).increment(1);
errorTracker.incrementCounters(e);
continue;
}

Expand Down
5 changes: 5 additions & 0 deletions src/java/org/apache/nutch/crawl/Injector.java
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,7 @@
import org.apache.hadoop.util.ToolRunner;

import org.apache.nutch.metadata.Nutch;
import org.apache.nutch.metrics.ErrorTracker;
import org.apache.nutch.metrics.NutchMetrics;
import org.apache.nutch.net.URLFilters;
import org.apache.nutch.net.URLNormalizers;
Expand Down Expand Up @@ -127,6 +128,7 @@ public static class InjectMapper
private boolean url404Purging;
private String scope;
private boolean filterNormalizeAll = false;
private ErrorTracker errorTracker;

@Override
public void setup(Context context) {
Expand All @@ -147,6 +149,8 @@ public void setup(Context context) {
curTime = conf.getLong("injector.current.time",
System.currentTimeMillis());
url404Purging = conf.getBoolean(CrawlDb.CRAWLDB_PURGE_404, false);
// Initialize error tracker with cached counters
errorTracker = new ErrorTracker(NutchMetrics.GROUP_INJECTOR, context);
}

/* Filter and normalize the input url */
Expand Down Expand Up @@ -239,6 +243,7 @@ public void map(Text key, Writable value, Context context)
LOG.warn(
"Cannot filter injected score for url {}, using default ({})",
url, e.getMessage());
errorTracker.incrementCounters(e);
}
context.getCounter(NutchMetrics.GROUP_INJECTOR,
NutchMetrics.INJECTOR_URLS_INJECTED_TOTAL).increment(1);
Expand Down
28 changes: 19 additions & 9 deletions src/java/org/apache/nutch/fetcher/FetcherThread.java
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,7 @@
import org.apache.nutch.crawl.SignatureFactory;
import org.apache.nutch.fetcher.Fetcher.FetcherRun;
import org.apache.nutch.fetcher.FetcherThreadEvent.PublishEventType;
import org.apache.nutch.metrics.ErrorTracker;
import org.apache.nutch.metrics.LatencyTracker;
import org.apache.nutch.metrics.NutchMetrics;
import org.apache.nutch.metadata.Metadata;
Expand Down Expand Up @@ -170,6 +171,9 @@ public class FetcherThread extends Thread {
// Latency tracker for fetch timing metrics
private LatencyTracker fetchLatencyTracker;

// Error tracker for categorized error metrics
private ErrorTracker errorTracker;

public FetcherThread(Configuration conf, AtomicInteger activeThreads, FetchItemQueues fetchQueues,
QueueFeeder feeder, AtomicInteger spinWaiting, AtomicLong lastRequestStart, FetcherRun.Context context,
AtomicInteger errors, String segmentName, boolean parsing, boolean storingContent,
Expand Down Expand Up @@ -292,6 +296,9 @@ private void initCounters() {
// Initialize latency tracker for fetch timing
fetchLatencyTracker = new LatencyTracker(
NutchMetrics.GROUP_FETCHER, NutchMetrics.FETCHER_LATENCY);

// Initialize error tracker for categorized error metrics
errorTracker = new ErrorTracker(NutchMetrics.GROUP_FETCHER);
}

@Override
Expand Down Expand Up @@ -548,15 +555,7 @@ public void run() {
} catch (Throwable t) { // unexpected exception
// unblock
fetchQueues.finishFetchItem(fit);
String message;
if (LOG.isDebugEnabled()) {
message = StringUtils.stringifyException(t);
} else if (logUtil.logShort(t)) {
message = t.getClass().getName();
} else {
message = StringUtils.stringifyException(t);
}
logError(fit.url, message);
logError(fit.url, t);
output(fit.url, fit.datum, null, ProtocolStatus.STATUS_FAILED,
CrawlDatum.STATUS_FETCH_RETRY);
}
Expand All @@ -570,6 +569,8 @@ public void run() {
}
// Emit fetch latency metrics
fetchLatencyTracker.emitCounters(context);
// Emit error metrics
errorTracker.emitCounters(context);
activeThreads.decrementAndGet(); // count threads
LOG.info("{} {} -finishing thread {}, activeThreads={}", getName(),
Thread.currentThread().getId(), getName(), activeThreads);
Expand Down Expand Up @@ -689,10 +690,19 @@ private FetchItem queueRedirect(Text redirUrl, FetchItem fit)
return fit;
}

private void logError(Text url, Throwable t) {
String message = t.getClass().getName() + ": " + t.getMessage();
LOG.info("{} {} fetch of {} failed with: {}", getName(),
Thread.currentThread().getId(), url, message);
errors.incrementAndGet();
errorTracker.recordError(t);
}

private void logError(Text url, String message) {
LOG.info("{} {} fetch of {} failed with: {}", getName(),
Thread.currentThread().getId(), url, message);
errors.incrementAndGet();
errorTracker.recordError(ErrorTracker.ErrorType.OTHER);
}

private ParseStatus output(Text key, CrawlDatum datum, Content content,
Expand Down
14 changes: 14 additions & 0 deletions src/java/org/apache/nutch/hostdb/ResolverThread.java
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@
import org.apache.hadoop.mapreduce.Reducer.Context;
import org.apache.hadoop.util.StringUtils;

import org.apache.nutch.metrics.ErrorTracker;
import org.apache.nutch.metrics.NutchMetrics;

import org.slf4j.Logger;
Expand Down Expand Up @@ -124,11 +125,24 @@ public void run() {

// Dynamic counter based on failure count - can't cache
context.getCounter(NutchMetrics.GROUP_HOSTDB, createFailureCounterLabel(datum)).increment(1);
// Common error counters for consistency
context.getCounter(NutchMetrics.GROUP_HOSTDB,
NutchMetrics.ERROR_TOTAL).increment(1);
context.getCounter(NutchMetrics.GROUP_HOSTDB,
NutchMetrics.ERROR_NETWORK_TOTAL).increment(1);
} catch (Exception ioe) {
LOG.warn(StringUtils.stringifyException(ioe));
context.getCounter(NutchMetrics.GROUP_HOSTDB,
NutchMetrics.ERROR_TOTAL).increment(1);
context.getCounter(NutchMetrics.GROUP_HOSTDB,
ErrorTracker.getCounterName(ioe)).increment(1);
}
} catch (Exception e) {
LOG.warn(StringUtils.stringifyException(e));
context.getCounter(NutchMetrics.GROUP_HOSTDB,
NutchMetrics.ERROR_TOTAL).increment(1);
context.getCounter(NutchMetrics.GROUP_HOSTDB,
ErrorTracker.getCounterName(e)).increment(1);
}

context.getCounter(NutchMetrics.GROUP_HOSTDB,
Expand Down
9 changes: 5 additions & 4 deletions src/java/org/apache/nutch/hostdb/UpdateHostDbMapper.java
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,7 @@
import org.apache.nutch.crawl.CrawlDatum;
import org.apache.nutch.crawl.NutchWritable;
import org.apache.nutch.metadata.Nutch;
import org.apache.nutch.metrics.ErrorTracker;
import org.apache.nutch.metrics.NutchMetrics;
import org.apache.nutch.net.URLFilters;
import org.apache.nutch.net.URLNormalizers;
Expand Down Expand Up @@ -63,8 +64,8 @@ public class UpdateHostDbMapper
protected URLNormalizers normalizers = null;

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

@Override
public void setup(Mapper<Text, Writable, Text, NutchWritable>.Context context) {
Expand All @@ -79,10 +80,10 @@ public void setup(Mapper<Text, Writable, Text, NutchWritable>.Context context) {
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);
// Initialize error tracker with cached counters
errorTracker = new ErrorTracker(NutchMetrics.GROUP_HOSTDB, context);
}

/**
Expand Down Expand Up @@ -148,7 +149,7 @@ public void map(Text key, Writable value,
try {
url = new URL(keyStr);
} catch (MalformedURLException e) {
malformedUrlCounter.increment(1);
errorTracker.incrementCounters(e);
return;
}
String hostName = URLUtil.getHost(url);
Expand Down
16 changes: 8 additions & 8 deletions src/java/org/apache/nutch/indexer/IndexerMapReduce.java
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,7 @@
import org.apache.nutch.crawl.Inlinks;
import org.apache.nutch.crawl.LinkDb;
import org.apache.nutch.crawl.NutchWritable;
import org.apache.nutch.metrics.ErrorTracker;
import org.apache.nutch.metrics.LatencyTracker;
import org.apache.nutch.metrics.NutchMetrics;
import org.apache.nutch.metadata.Metadata;
Expand Down Expand Up @@ -226,11 +227,12 @@ public static class IndexerReducer extends
private Counter deletedRedirectsCounter;
private Counter deletedDuplicatesCounter;
private Counter skippedNotModifiedCounter;
private Counter errorsScoringFilterCounter;
private Counter errorsIndexingFilterCounter;
private Counter deletedByIndexingFilterCounter;
private Counter skippedByIndexingFilterCounter;
private Counter indexedCounter;

// Error tracker with cached counters
private ErrorTracker errorTracker;

@Override
public void setup(Reducer<Text, NutchWritable, Text, NutchIndexAction>.Context context) {
Expand Down Expand Up @@ -279,16 +281,14 @@ private void initCounters(Reducer<Text, NutchWritable, Text, NutchIndexAction>.C
NutchMetrics.GROUP_INDEXER, NutchMetrics.INDEXER_DELETED_DUPLICATES_TOTAL);
skippedNotModifiedCounter = context.getCounter(
NutchMetrics.GROUP_INDEXER, NutchMetrics.INDEXER_SKIPPED_NOT_MODIFIED_TOTAL);
errorsScoringFilterCounter = context.getCounter(
NutchMetrics.GROUP_INDEXER, NutchMetrics.INDEXER_ERRORS_SCORING_FILTER_TOTAL);
errorsIndexingFilterCounter = context.getCounter(
NutchMetrics.GROUP_INDEXER, NutchMetrics.INDEXER_ERRORS_INDEXING_FILTER_TOTAL);
deletedByIndexingFilterCounter = context.getCounter(
NutchMetrics.GROUP_INDEXER, NutchMetrics.INDEXER_DELETED_BY_INDEXING_FILTER_TOTAL);
skippedByIndexingFilterCounter = context.getCounter(
NutchMetrics.GROUP_INDEXER, NutchMetrics.INDEXER_SKIPPED_BY_INDEXING_FILTER_TOTAL);
indexedCounter = context.getCounter(
NutchMetrics.GROUP_INDEXER, NutchMetrics.INDEXER_INDEXED_TOTAL);
// Initialize error tracker with cached counters
errorTracker = new ErrorTracker(NutchMetrics.GROUP_INDEXER, context);
}

@Override
Expand Down Expand Up @@ -416,7 +416,7 @@ public void reduce(Text key, Iterable<NutchWritable> values,
boost = scfilters.indexerScore(key, doc, dbDatum, fetchDatum, parse,
inlinks, boost);
} catch (final ScoringFilterException e) {
errorsScoringFilterCounter.increment(1);
errorTracker.incrementCounters(e);
LOG.warn("Error calculating score {}: {}", key, e);
return;
}
Expand Down Expand Up @@ -451,7 +451,7 @@ public void reduce(Text key, Iterable<NutchWritable> values,
doc = filters.filter(doc, parse, key, fetchDatum, inlinks);
} catch (final IndexingException e) {
LOG.warn("Error indexing {}: ", key, e);
errorsIndexingFilterCounter.increment(1);
errorTracker.incrementCounters(e);
return;
}

Expand Down
Loading