diff --git a/core/src/main/java/org/apache/stormcrawler/bolt/CrawlDelayPolicy.java b/core/src/main/java/org/apache/stormcrawler/bolt/CrawlDelayPolicy.java new file mode 100644 index 000000000..448b311f9 --- /dev/null +++ b/core/src/main/java/org/apache/stormcrawler/bolt/CrawlDelayPolicy.java @@ -0,0 +1,114 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to you under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.stormcrawler.bolt; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * Decides what the crawl delay of robots.txt does to the fetch queue of a URL, following + * fetcher.max.crawl.delay, fetcher.max.crawl.delay.force, fetcher.server.delay and + * fetcher.server.delay.force. Delays are in milliseconds. + */ +final class CrawlDelayPolicy { + + // the bolt's category: log configurations for FetcherBolt keep covering these lines + private static final Logger LOG = LoggerFactory.getLogger(FetcherBolt.class); + + enum Action { + /** + * Leave the fetch queue as it is: robots.txt has no delay, or the one the queue already + * has, even above the cap. The caller must not write back the delay it read, which could + * undo a concurrent raise. + */ + UNCHANGED, + /** Do not fetch the URL: robots.txt asks for a longer delay than accepted. */ + SKIP, + APPLY + } + + /** + * @param delay the delay to set on the fetch queue, for {@link Action#APPLY} + * @param robotsCrawlDelaySecs the robots.txt delay in seconds, rounded up, when it was longer + * than fetcher.max.crawl.delay and capped; null otherwise + */ + record Decision(Action action, long delay, String robotsCrawlDelaySecs) { + static final Decision UNCHANGED = new Decision(Action.UNCHANGED, 0, null); + static final Decision SKIP = new Decision(Action.SKIP, 0, null); + } + + // max. delay accepted from robots.txt, negative for no limit + private final long maxCrawlDelay; + // whether maxCrawlDelay overwrites the longer value in robots.txt + // (otherwise URLs in this queue are skipped) + private final boolean maxCrawlDelayForce; + private final long serverDelay; + // whether the default delay is used even if the robots.txt + // specifies a shorter crawl-delay + private final boolean serverDelayForce; + + CrawlDelayPolicy( + long maxCrawlDelay, + boolean maxCrawlDelayForce, + long serverDelay, + boolean serverDelayForce) { + this.maxCrawlDelay = maxCrawlDelay; + this.maxCrawlDelayForce = maxCrawlDelayForce; + this.serverDelay = serverDelay; + this.serverDelayForce = serverDelayForce; + } + + /** + * @param url the URL, for the log + * @param queueId the ID of the fetch queue, for the log + * @param robotsDelay the crawl delay of robots.txt, not positive when there is none + */ + Decision decide(String url, String queueId, long robotsDelay, long queueDelay) { + if (robotsDelay <= 0 || robotsDelay == queueDelay) { + return Decision.UNCHANGED; + } + if (robotsDelay > maxCrawlDelay && maxCrawlDelay >= 0) { + if (!maxCrawlDelayForce) { + LOG.info("Crawl-Delay for {} too long ({}), skipping", url, robotsDelay); + return Decision.SKIP; + } + LOG.info( + "Crawl-Delay for {} too long ({}), using value of fetcher.max.crawl.delay" + + " instead", + url, + robotsDelay); + // report the delay the fetcher is not holding, so a frontier-side + // consumer can enforce it at the source (#867) + return new Decision( + Action.APPLY, maxCrawlDelay, Long.toString(1L + ((robotsDelay - 1L) / 1000L))); + } + if (robotsDelay < serverDelay && serverDelayForce) { + LOG.info( + "Crawl delay for {} too short ({}), set to fetcher.server.delay", + url, + robotsDelay); + return new Decision(Action.APPLY, serverDelay, null); + } + LOG.info( + "Crawl delay for queue: {} is set to {} as per robots.txt. url: {}", + queueId, + robotsDelay, + url); + return new Decision(Action.APPLY, robotsDelay, null); + } +} diff --git a/core/src/main/java/org/apache/stormcrawler/bolt/FetcherBolt.java b/core/src/main/java/org/apache/stormcrawler/bolt/FetcherBolt.java index 92aa08ad1..5763692bd 100644 --- a/core/src/main/java/org/apache/stormcrawler/bolt/FetcherBolt.java +++ b/core/src/main/java/org/apache/stormcrawler/bolt/FetcherBolt.java @@ -103,6 +103,8 @@ public class FetcherBolt extends StatusEmitterBolt { private RobotRulesLookup robotsLookup; + private CrawlDelayPolicy crawlDelayPolicy; + /** Largest number of helper threads ever alive; for tests. */ int helperPoolSize() { return fetchHelpers == null ? 0 : fetchHelpers.largestPoolSize(); @@ -119,14 +121,6 @@ public Map getComponentConfiguration() { /** This class picks items from queues and fetches the pages. */ private class FetcherThread extends Thread { - // max. delay accepted from robots.txt - private final long maxCrawlDelay; - // whether maxCrawlDelay overwrites the longer value in robots.txt - // (otherwise URLs in this queue are skipped) - private final boolean maxCrawlDelayForce; - // whether the default delay is used even if the robots.txt - // specifies a shorter crawl-delay - private final boolean crawlDelayForce; private final int threadNum; private long timeoutInQueues = -1; @@ -138,10 +132,6 @@ public FetcherThread(Config conf, int num) { this.setDaemon(true); // don't hang JVM on exit this.setName("FetcherThread #" + num); // use an informative name - this.maxCrawlDelay = ConfUtils.getInt(conf, "fetcher.max.crawl.delay", 30) * 1000L; - this.maxCrawlDelayForce = - ConfUtils.getBoolean(conf, "fetcher.max.crawl.delay.force", false); - this.crawlDelayForce = ConfUtils.getBoolean(conf, "fetcher.server.delay.force", false); this.threadNum = num; timeoutInQueues = ConfUtils.getLong(conf, QUEUED_TIMEOUT_PARAM_KEY, timeoutInQueues); protocolMetadataPrefix = @@ -231,56 +221,30 @@ public void run() { continue; } FetchItemQueue fiq = fetchQueues.getFetchItemQueue(fit.queueId, metadata); - if (rules.getCrawlDelay() > 0 - && rules.getCrawlDelay() != fiq.crawlDelay.get()) { - if (rules.getCrawlDelay() > maxCrawlDelay && maxCrawlDelay >= 0) { - boolean force = false; - String msg = "skipping"; - if (maxCrawlDelayForce) { - force = true; - msg = "using value of fetcher.max.crawl.delay instead"; - } - LOG.info( - "Crawl-Delay for {} too long ({}), {}", + CrawlDelayPolicy.Decision decision = + crawlDelayPolicy.decide( fit.url, - rules.getCrawlDelay(), - msg); - if (force) { - fiq.crawlDelay.set(maxCrawlDelay); - // report the delay the fetcher is not holding, so a frontier-side - // consumer can enforce it at the source (#867) - robotsCrawlDelaySecs = - Long.toString(1L + ((rules.getCrawlDelay() - 1L) / 1000L)); - metadata.setValue( - Constants.ROBOTS_CRAWL_DELAY_KEY, robotsCrawlDelaySecs); - } else { - // pass the info about crawl delay - metadata.setValue(Constants.STATUS_ERROR_CAUSE, "crawl_delay"); - collector.emit( - org.apache.stormcrawler.Constants.StatusStreamName, - fit.tuple, - new Values(fit.url, metadata, Status.ERROR)); - // no need to wait next time as we won't request - // from that site - asap = true; - continue; - } - } else if (rules.getCrawlDelay() < fetchQueues.crawlDelay - && crawlDelayForce) { - fiq.crawlDelay.set(fetchQueues.crawlDelay); - LOG.info( - "Crawl delay for {} too short ({}), " - + "set to fetcher.server.delay", - fit.url, - rules.getCrawlDelay()); - } else { - fiq.crawlDelay.set(rules.getCrawlDelay()); - LOG.info( - "Crawl delay for queue: {} is set to {} " - + "as per robots.txt. url: {}", fit.queueId, - fiq.crawlDelay.get(), - fit.url); + rules.getCrawlDelay(), + fiq.crawlDelay.get()); + if (decision.action() == CrawlDelayPolicy.Action.SKIP) { + // pass the info about crawl delay + metadata.setValue(Constants.STATUS_ERROR_CAUSE, "crawl_delay"); + collector.emit( + org.apache.stormcrawler.Constants.StatusStreamName, + fit.tuple, + new Values(fit.url, metadata, Status.ERROR)); + // no need to wait next time as we won't request + // from that site + asap = true; + continue; + } + if (decision.action() == CrawlDelayPolicy.Action.APPLY) { + fiq.crawlDelay.set(decision.delay()); + robotsCrawlDelaySecs = decision.robotsCrawlDelaySecs(); + if (robotsCrawlDelaySecs != null) { + metadata.setValue( + Constants.ROBOTS_CRAWL_DELAY_KEY, robotsCrawlDelaySecs); } } @@ -609,6 +573,13 @@ public void prepare( this.fetchQueues = new FetchItemQueues(conf); + crawlDelayPolicy = + new CrawlDelayPolicy( + ConfUtils.getInt(conf, "fetcher.max.crawl.delay", 30) * 1000L, + ConfUtils.getBoolean(conf, "fetcher.max.crawl.delay.force", false), + fetchQueues.crawlDelay, + ConfUtils.getBoolean(conf, "fetcher.server.delay.force", false)); + this.taskId = context.getThisTaskId(); int threadCount = ConfUtils.getInt(conf, "fetcher.threads.number", 10); diff --git a/core/src/test/java/org/apache/stormcrawler/bolt/CrawlDelayPolicyTest.java b/core/src/test/java/org/apache/stormcrawler/bolt/CrawlDelayPolicyTest.java new file mode 100644 index 000000000..b7d94d761 --- /dev/null +++ b/core/src/test/java/org/apache/stormcrawler/bolt/CrawlDelayPolicyTest.java @@ -0,0 +1,77 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to you under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.stormcrawler.bolt; + +import static org.junit.jupiter.api.Assertions.assertEquals; + +import org.apache.stormcrawler.bolt.CrawlDelayPolicy.Action; +import org.apache.stormcrawler.bolt.CrawlDelayPolicy.Decision; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.CsvSource; + +class CrawlDelayPolicyTest { + + /** + * Delays in milliseconds. The columns: robots.txt delay, current delay of the fetch queue, + * fetcher.max.crawl.delay, fetcher.max.crawl.delay.force, fetcher.server.delay, + * fetcher.server.delay.force, then the expected decision: action, delay, robots crawl delay + * reported in seconds. + */ + @ParameterizedTest + @CsvSource( + textBlock = + """ + # no delay in robots.txt, or the one the fetch queue already has + -9223372036854775808, 1000, 30000, false, 1000, false, UNCHANGED, 0, + 0, 1000, 30000, false, 1000, false, UNCHANGED, 0, + 5000, 5000, 30000, false, 1000, false, UNCHANGED, 0, + 60000, 60000, 30000, false, 1000, false, UNCHANGED, 0, + # longer than fetcher.max.crawl.delay, which applies before the server delay + 60000, 1000, 30000, false, 1000, false, SKIP, 0, + 60000, 1000, 30000, true, 1000, false, APPLY, 30000, 60 + 30500, 1000, 30000, true, 1000, false, APPLY, 30000, 31 + 60000, 1000, 30000, false, 1000, true, SKIP, 0, + 60000, 1000, 30000, true, 1000, true, APPLY, 30000, 60 + 5000, 1000, 0, false, 1000, false, SKIP, 0, + # within fetcher.max.crawl.delay, or no limit when negative + 30000, 1000, 30000, false, 1000, false, APPLY, 30000, + 120000, 1000, -1, false, 1000, false, APPLY, 120000, + 5000, 1000, 30000, false, 1000, true, APPLY, 5000, + # shorter than fetcher.server.delay + 500, 0, 30000, false, 1000, true, APPLY, 1000, + 500, 0, 30000, false, 1000, false, APPLY, 500, + """) + void decision( + long robotsDelay, + long queueDelay, + long maxCrawlDelay, + boolean maxCrawlDelayForce, + long serverDelay, + boolean serverDelayForce, + Action action, + long delay, + String reportedSecs) { + CrawlDelayPolicy policy = + new CrawlDelayPolicy( + maxCrawlDelay, maxCrawlDelayForce, serverDelay, serverDelayForce); + + assertEquals( + new Decision(action, delay, reportedSecs), + policy.decide("http://a.test/", "a.test", robotsDelay, queueDelay)); + } +}