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
113 changes: 54 additions & 59 deletions core/src/main/java/org/apache/stormcrawler/bolt/FetcherBolt.java
Original file line number Diff line number Diff line change
Expand Up @@ -24,7 +24,6 @@
import java.time.Instant;
import java.util.Arrays;
import java.util.Map;
import java.util.Map.Entry;
import java.util.concurrent.atomic.AtomicInteger;
import org.apache.commons.lang3.StringUtils;
import org.apache.storm.Config;
Expand All @@ -38,7 +37,14 @@
import org.apache.storm.utils.Utils;
import org.apache.stormcrawler.Constants;
import org.apache.stormcrawler.Metadata;
import org.apache.stormcrawler.bolt.FetchItemQueues.FetchItemQueue;
import org.apache.stormcrawler.fetcher.CrawlDelayPolicy;
import org.apache.stormcrawler.fetcher.FetchItem;
import org.apache.stormcrawler.fetcher.FetchItemQueues;
import org.apache.stormcrawler.fetcher.FetchItemQueues.FetchItemQueue;
import org.apache.stormcrawler.fetcher.FetchOutcomes;
import org.apache.stormcrawler.fetcher.FetchTimeoutHelpers;
import org.apache.stormcrawler.fetcher.ProtocolMetrics;
import org.apache.stormcrawler.fetcher.RobotRulesLookup;
import org.apache.stormcrawler.metrics.CrawlerMetrics;
import org.apache.stormcrawler.metrics.ScopedCounter;
import org.apache.stormcrawler.metrics.ScopedReducedMetric;
Expand Down Expand Up @@ -152,22 +158,22 @@ public void run() {

activeThreads.incrementAndGet(); // count threads

beingFetched[threadNum] = fit.url;
beingFetched[threadNum] = fit.url();

LOG.debug(
"[Fetcher #{}] {} => activeThreads={}, spinWaiting={}, queueID={}",
taskId,
getName(),
activeThreads,
spinWaiting,
fit.queueId);
fit.queueId());

LOG.debug("[Fetcher #{}] {} : Fetching {}", taskId, getName(), fit.url);
LOG.debug("[Fetcher #{}] {} : Fetching {}", taskId, getName(), fit.url());

Metadata metadata = null;

if (fit.tuple.contains("metadata")) {
metadata = (Metadata) fit.tuple.getValueByField("metadata");
if (fit.tuple().contains("metadata")) {
metadata = (Metadata) fit.tuple().getValueByField("metadata");
}
if (metadata == null) {
metadata = new Metadata();
Expand All @@ -184,54 +190,54 @@ public void run() {
boolean saturated = false;

try {
URL url = URLUtil.toURL(fit.url);
URL url = URLUtil.toURL(fit.url());
Protocol protocol = protocolFactory.getProtocol(url);

if (protocol == null) {
throw new RuntimeException(
"No protocol implementation found for " + fit.url);
"No protocol implementation found for " + fit.url());
}

RobotRulesLookup.Result robots =
robotsLookup.lookup(protocol, fit.url, metadata);
robotsLookup.lookup(protocol, fit.url(), metadata);
BaseRobotRules rules = robots.rules();

discoverSitemaps(fit.tuple, url, metadata, robots);
discoverSitemaps(fit.tuple(), url, metadata, robots);

if (!rules.isAllowed(fit.url)) {
LOG.info("Denied by robots.txt: {}", fit.url);
if (!rules.isAllowed(fit.url())) {
LOG.info("Denied by robots.txt: {}", fit.url());
// pass the info about denied by robots
metadata.setValue(Constants.STATUS_ERROR_CAUSE, "robots.txt");
collector.emit(
org.apache.stormcrawler.Constants.StatusStreamName,
fit.tuple,
new Values(fit.url, metadata, Status.ERROR));
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;
}
FetchItemQueue fiq = fetchQueues.getFetchItemQueue(fit.queueId, metadata);
FetchItemQueue fiq = fetchQueues.getFetchItemQueue(fit.queueId(), metadata);
CrawlDelayPolicy.Decision decision =
crawlDelayPolicy.decide(
fit.url,
fit.queueId,
fit.url(),
fit.queueId(),
rules.getCrawlDelay(),
fiq.crawlDelay.get());
fiq.getCrawlDelay());
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));
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());
fiq.setCrawlDelay(decision.delay());
robotsCrawlDelaySecs = decision.robotsCrawlDelaySecs();
if (robotsCrawlDelaySecs != null) {
metadata.setValue(
Expand All @@ -240,13 +246,15 @@ public void run() {
}

long start = System.currentTimeMillis();
long timeInQueues = start - fit.creationTime;
long timeInQueues = start - fit.creationTime();

// waited longer than fetcher.timeout.queue: not fetched, acked without a
// status (see finally)
if (timeoutInQueues != -1 && timeInQueues > timeoutInQueues * 1000) {
LOG.info(
"[Fetcher #{}] Waited in queue for too long - {}", taskId, fit.url);
"[Fetcher #{}] Waited in queue for too long - {}",
taskId,
fit.url());
eventCounter.scope("queue.timeout").incrBy(1);
// no need to wait next time as we won't request from
// that site
Expand All @@ -258,9 +266,9 @@ public void run() {
ProtocolResponse response;
response =
fetchHelpers.call(
() -> protocol.getProtocolOutput(fit.url, fetchMetadata),
() -> protocol.getProtocolOutput(fit.url(), fetchMetadata),
protocol,
fit.url,
fit.url(),
fetchMetadata);

long timeFetching = System.currentTimeMillis() - start;
Expand All @@ -281,7 +289,7 @@ public void run() {
LOG.info(
"[Fetcher #{}] Fetched {} with status {} in msec {}",
taskId,
fit.url,
fit.url(),
response.getStatusCode(),
timeFetching);

Expand All @@ -296,25 +304,25 @@ public void run() {
// send content for parsing
collector.emit(
Utils.DEFAULT_STREAM_ID,
fit.tuple,
new Values(fit.url, response.getContent(), outcome.metadata()));
fit.tuple(),
new Values(fit.url(), response.getContent(), outcome.metadata()));
} else {
// allowRedirs() can be overridden: it is asked for every redirect and
// nothing else (#954)
if (outcome.status() == Status.REDIRECTION
&& allowRedirs()
&& outcome.redirectTarget() != null) {
emitOutlink(
fit.tuple,
fit.tuple(),
url,
outcome.redirectTarget(),
outcome.metadata(),
outcome.redirectKeyValues());
}
collector.emit(
Constants.StatusStreamName,
fit.tuple,
new Values(fit.url, outcome.metadata(), outcome.status()));
fit.tuple(),
new Values(fit.url(), outcome.metadata(), outcome.status()));
}

} catch (FetchTimeoutHelpers.SaturatedException e) {
Expand All @@ -333,25 +341,25 @@ && allowRedirs()
saturated = true;
} catch (Exception exece) {
FetchOutcomes.Outcome outcome =
fetchOutcomes.ofFailure(exece, fit.url, metadata);
fetchOutcomes.ofFailure(exece, fit.url(), metadata);
collector.emit(
Constants.StatusStreamName,
fit.tuple,
new Values(fit.url, outcome.metadata(), outcome.status()));
fit.tuple(),
new Values(fit.url(), outcome.metadata(), outcome.status()));
} finally {
if (saturated) {
long delay = fetchQueues.backOffFetchItem(fit);
LOG.debug(
"[Fetcher #{}] queue {} backed off for {} ms",
taskId,
fit.queueId,
fit.queueId(),
delay);
} else {
fetchQueues.finishFetchItem(fit, asap);
}
activeThreads.decrementAndGet(); // count threads
// ack it whatever happens
collector.ack(fit.tuple);
collector.ack(fit.tuple());
beingFetched[threadNum] = "";
}
}
Expand Down Expand Up @@ -443,14 +451,14 @@ public void prepare(
context,
stormConf,
"in_queues",
() -> fetchQueues.inQueues.get(),
() -> fetchQueues.numQueuedItems(),
metricsTimeBucketSecs);

CrawlerMetrics.registerGauge(
context,
stormConf,
"num_queues",
() -> fetchQueues.queues.size(),
() -> fetchQueues.numQueues(),
metricsTimeBucketSecs);

this.averagedMetrics =
Expand All @@ -474,7 +482,7 @@ public void prepare(
new CrawlDelayPolicy(
ConfUtils.getInt(conf, "fetcher.max.crawl.delay", 30) * 1000L,
ConfUtils.getBoolean(conf, "fetcher.max.crawl.delay.force", false),
fetchQueues.crawlDelay,
fetchQueues.defaultCrawlDelay(),
ConfUtils.getBoolean(conf, "fetcher.server.delay.force", false));

this.taskId = context.getThisTaskId();
Expand Down Expand Up @@ -555,7 +563,7 @@ public void execute(Tuple input) {
}

if (this.maxNumberUrlsInQueues != -1) {
while (this.activeThreads.get() + this.fetchQueues.inQueues.get()
while (this.activeThreads.get() + this.fetchQueues.numQueuedItems()
>= maxNumberUrlsInQueues) {
try {
Thread.sleep(500);
Expand All @@ -567,8 +575,8 @@ public void execute(Tuple input) {
"[Fetcher #{}] Threads : {}\tqueues : {}\tin_queues : {}",
taskId,
this.activeThreads.get(),
this.fetchQueues.queues.size(),
this.fetchQueues.inQueues.get());
this.fetchQueues.numQueues(),
this.fetchQueues.numQueuedItems());
}
}

Expand Down Expand Up @@ -610,24 +618,11 @@ public void execute(Tuple input) {
}

/**
* Logs the content of the queues and the URLs being fetched. Not synchronized: the queues and
* their items are iterated with weakly consistent iterators while the fetcher threads keep
* working, so the dump is a smear over the time it takes to produce it rather than a
* point-in-time snapshot. Sizes and item lists may not add up exactly.
* Logs the content of the queues and the URLs being fetched. The URLs being fetched are read
* without locking, like the queues, so the two lists may not match exactly.
*/
private void logQueuesContent() {
StringBuilder sb = new StringBuilder();
sb.append("\nNum queues : ").append(fetchQueues.queues.size());
for (Entry<String, FetchItemQueue> entry : fetchQueues.queues.entrySet()) {
sb.append("\nQueue ID : ").append(entry.getKey());
FetchItemQueue fiq = entry.getValue();
sb.append("\t size : ").append(fiq.getQueueSize());
sb.append("\t in progress : ").append(fiq.getInProgressSize());
for (FetchItem fetchItem : fiq.queue) {
sb.append("\n\t").append(fetchItem.url);
}
}
LOG.info("Dumping queue content {}", sb.toString());
LOG.info("Dumping queue content {}", fetchQueues.dump());

StringBuilder sb2 = new StringBuilder("\n");
// dump the list of URLs being fetched
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,8 @@
import org.apache.storm.utils.Utils;
import org.apache.stormcrawler.Constants;
import org.apache.stormcrawler.Metadata;
import org.apache.stormcrawler.fetcher.FetchTimeoutHelpers;
import org.apache.stormcrawler.fetcher.ProtocolMetrics;
import org.apache.stormcrawler.metrics.CrawlerMetrics;
import org.apache.stormcrawler.metrics.ScopedCounter;
import org.apache.stormcrawler.metrics.ScopedReducedMetric;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -15,22 +15,25 @@
* limitations under the License.
*/

package org.apache.stormcrawler.bolt;
package org.apache.stormcrawler.fetcher;

import org.apache.stormcrawler.bolt.FetcherBolt;
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.
*
* <p>For internal use only. Not part of StormCrawler's public API.
*/
final class CrawlDelayPolicy {
public 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 {
public 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
Expand All @@ -47,7 +50,7 @@ enum Action {
* @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) {
public 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);
}
Expand All @@ -62,7 +65,7 @@ record Decision(Action action, long delay, String robotsCrawlDelaySecs) {
// specifies a shorter crawl-delay
private final boolean serverDelayForce;

CrawlDelayPolicy(
public CrawlDelayPolicy(
long maxCrawlDelay,
boolean maxCrawlDelayForce,
long serverDelay,
Expand All @@ -78,7 +81,7 @@ record Decision(Action action, long delay, String robotsCrawlDelaySecs) {
* @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) {
public Decision decide(String url, String queueId, long robotsDelay, long queueDelay) {
if (robotsDelay <= 0 || robotsDelay == queueDelay) {
return Decision.UNCHANGED;
}
Expand Down
Loading
Loading