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
39 changes: 27 additions & 12 deletions src/java/org/apache/nutch/crawl/AdaptiveFetchSchedule.java
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,7 @@
import java.lang.invoke.MethodHandles;
import java.net.URI;
import java.net.URISyntaxException;
import java.time.Duration;

/**
* This class implements an adaptive re-fetch algorithm. This works as follows:
Expand Down Expand Up @@ -219,15 +220,15 @@ private void setHostSpecificIntervals(String fileName,
// The custom intervals should respect the boundaries of the default values.
if (m < defaultMin) {
LOG.error(
"Min. interval out of bounds on line {} in the config. file: `{}`",
lineNo, line);
"Min. interval out of bounds ({}) on line {} in the config. file: `{}`",
defaultMin, lineNo, line);
continue;
}

if (M > defaultMax) {
LOG.error(
"Max. interval out of bounds on line {} in the config. file: `{}`",
lineNo, line);
"Max. interval out of bounds ({}) on line {} in the config. file: `{}`",
defaultMax, lineNo, line);
continue;
}

Expand Down Expand Up @@ -332,17 +333,30 @@ public CrawlDatum setFetchSchedule(Text url, CrawlDatum datum,
case FetchSchedule.STATUS_UNKNOWN:
break;
}
if (SYNC_DELTA) {
// try to synchronize with the time of change
long delta = (fetchTime - modifiedTime) / 1000L;
if (delta > interval)
interval = delta;
refTime = fetchTime - Math.round(delta * SYNC_DELTA_RATE * 1000);
}

// Ensure the interval does not fall outside of bounds
float minInterval = (getCustomMinInterval(url) != null) ? getCustomMinInterval(url) : MIN_INTERVAL;
float maxInterval = (getCustomMaxInterval(url) != null) ? getCustomMaxInterval(url) : MAX_INTERVAL;

if (SYNC_DELTA) {
// try to synchronize with the time of change
long delta = (fetchTime - modifiedTime);
if (delta > (interval * 1000))
interval = delta / 1000L;
// offset: a fraction (sync_delta_rate) of the difference between the last modification time, and the last fetch time.
long offset = Math.round(delta * SYNC_DELTA_RATE);
long maxIntervalMillis = (long) maxInterval * 1000L;
if (LOG.isTraceEnabled()) {
LOG.trace("delta (days): {}; offset (days): {}; maxInterval (days): {}",
Duration.ofMillis(delta).toDays(), Duration.ofMillis(offset).toDays(), Duration.ofMillis(maxIntervalMillis).toDays());
}
// convert the offset to a ratio of max interval: avoid next fetchTime in the past, and mimic fetches within max interval
if (delta > 0 && offset > maxIntervalMillis) {
offset = offset / delta * maxIntervalMillis; // ex: 9/30*7 = 2.1
}
refTime = fetchTime - offset;
}

if (interval < minInterval) {
interval = minInterval;
} else if (interval > maxInterval) {
Expand Down Expand Up @@ -389,7 +403,8 @@ public static void main(String[] args) throws Exception {
(p.getFetchInterval() / SECONDS_PER_DAY), miss);
if (p.getFetchTime() <= curTime) {
fetchCnt++;
fs.setFetchSchedule(new Text("http://www.example.com"), p, p
// Text (url) required by the API, but not relevant here.
fs.setFetchSchedule(new Text(), p, p
.getFetchTime(), p.getModifiedTime(), curTime, lastModified,
changed ? FetchSchedule.STATUS_MODIFIED
: FetchSchedule.STATUS_NOTMODIFIED);
Expand Down
92 changes: 92 additions & 0 deletions src/test/org/apache/nutch/crawl/TestAdaptiveFetchSchedule.java
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,12 @@

import static org.hamcrest.CoreMatchers.is;
import static org.hamcrest.MatcherAssert.assertThat;
import static org.junit.jupiter.api.Assertions.assertTrue;

import java.time.Duration;
import java.time.Instant;
import java.util.Date;
import java.util.Properties;

/**
* Test cases for AdaptiveFetchSchedule.
Expand Down Expand Up @@ -117,5 +123,91 @@ private void validateFetchInterval(int changed, int getInterval) {
}

}

/**
* Test https://issues.apache.org/jira/browse/NUTCH-1564
*/
@Test
public void testSetFetchSchedule1() {
// db.fetch.schedule.adaptive.sync_delta_rate = 0.3 (default)
// db.fetch.interval.default = 172800 (2 days)
// db.fetch.schedule.adaptive.min_interval = 86400 (1 day)
// db.fetch.schedule.adaptive.max_interval = 604800 (7 days)
// db.fetch.interval.max = 604800 (7 days)
// 3-days cycle
// 30 days since last modified
doTestSetFetchSchedule(0.3, 2, 1, 7, 7, 3, 30);
}

@Test
public void testSetFetchSchedule2() {
// db.fetch.schedule.adaptive.sync_delta_rate = 0.3 (default)
// db.fetch.interval.default = 86400 (1 day)
// db.fetch.schedule.adaptive.min_interval = 86400 (1 day)
// db.fetch.schedule.adaptive.max_interval = 172800 (2 days)
// db.fetch.interval.max = 604800 (7 days)
// 1-day cycle
// 10 days since last modified
doTestSetFetchSchedule(0.3, 1, 1, 2, 7, 1, 10);
}

@Test
public void testSetFetchSchedule3() {
// db.fetch.schedule.adaptive.sync_delta_rate = 0.3 (default)
// db.fetch.interval.default = 172800 (2 days)
// db.fetch.schedule.adaptive.min_interval = 86400 (1 day)
// db.fetch.schedule.adaptive.max_interval = 864000 (10 days)
// db.fetch.interval.max = 864000 (10 days)
// 3-days cycle
// 180 days since last modified
doTestSetFetchSchedule(0.3, 2, 1, 10, 10, 3, 180);
}

private void doTestSetFetchSchedule(double deltaRate, int intervalDefaultDays,
int minIntervalDays, int maxIntervalDays, int intervalMaxDays,
int previousFetchTimeDays, int modifiedTimeDays) {
// need to properly override defaults
Properties props = new Properties();
props.setProperty("db.fetch.schedule.class", "org.apache.nutch.crawl.AdaptiveFetchSchedule");
props.setProperty("db.fetch.schedule.adaptive.sync_delta", "true"); // default
props.setProperty("db.fetch.schedule.adaptive.sync_delta_rate", String.valueOf(deltaRate));
props.setProperty("db.fetch.interval.default", String.valueOf(FetchSchedule.SECONDS_PER_DAY * intervalDefaultDays));
props.setProperty("db.fetch.schedule.adaptive.min_interval", String.valueOf(FetchSchedule.SECONDS_PER_DAY * minIntervalDays));
props.setProperty("db.fetch.schedule.adaptive.max_interval", String.valueOf(FetchSchedule.SECONDS_PER_DAY * maxIntervalDays));
props.setProperty("db.fetch.interval.max", String.valueOf(FetchSchedule.SECONDS_PER_DAY * intervalMaxDays));

conf = NutchConfiguration.create(true, props);
inc_rate = conf.getFloat("db.fetch.schedule.adaptive.inc_rate", 0.2f); // default
dec_rate = conf.getFloat("db.fetch.schedule.adaptive.dec_rate", 0.2f); // default

// ignore adaptive-host-specific-intervals.txt
Text url = new Text("http://www.example2.com");

AdaptiveFetchSchedule fs = new AdaptiveFetchSchedule();
fs.setConf(conf);

CrawlDatum datum = prepareCrawlDatum();
Date fetchTime = Date.from(Instant.now());
Date previousFetchTime = Date.from(Instant.now().minus(Duration.ofDays(previousFetchTimeDays)));
Date modifiedTime = Date.from(Instant.now().minus(Duration.ofDays(modifiedTimeDays)));
datum.setStatus(CrawlDatum.STATUS_FETCH_SUCCESS);
datum.setRetriesSinceFetch(0);
datum.setModifiedTime(modifiedTime.getTime());
datum.setFetchTime(fetchTime.getTime());

System.out.println("CrawlDatum fetchTime: " + fetchTime + "; modifiedTime: " + modifiedTime);

fs.setFetchSchedule(url, datum, previousFetchTime.getTime(), modifiedTime.getTime(),
fetchTime.getTime(), modifiedTime.getTime(), CrawlDatum.STATUS_DB_NOTMODIFIED);

Date nextFetchTime = new Date(datum.getFetchTime());
System.out.println("CrawlDatum next fetchTime: " + nextFetchTime);

assertTrue(nextFetchTime.after(fetchTime));
// adapt milliseconds to seconds
long fetchTimeDiff = (nextFetchTime.getTime() - fetchTime.getTime()) / 1000L ;
assertTrue(fetchTimeDiff >= FetchSchedule.SECONDS_PER_DAY * minIntervalDays);
assertTrue(fetchTimeDiff <= FetchSchedule.SECONDS_PER_DAY * maxIntervalDays);
}

}
Original file line number Diff line number Diff line change
Expand Up @@ -29,7 +29,7 @@
import static org.apache.nutch.crawl.CrawlDatum.*;
import static org.junit.jupiter.api.Assertions.fail;

public class TODOTestCrawlDbStates extends TestCrawlDbStates {
public class TestCrawlDbStatesExtended extends TestCrawlDbStates {

private static final Logger LOG = LoggerFactory
.getLogger(MethodHandles.lookup().lookupClass());
Expand Down