diff --git a/engine/src/main/java/org/archive/crawler/framework/CrawlJob.java b/engine/src/main/java/org/archive/crawler/framework/CrawlJob.java index b248a76c..11751013 100644 --- a/engine/src/main/java/org/archive/crawler/framework/CrawlJob.java +++ b/engine/src/main/java/org/archive/crawler/framework/CrawlJob.java @@ -710,7 +710,13 @@ public class CrawlJob implements Comparable, ApplicationListener { .append(" queued = ") .append(total) .append(" total"); - return sb.toString(); + if(snapshot.futureUriCount>0) { + sb + .append(" (") + .append(snapshot.futureUriCount) + .append(" future)"); + } + return sb.toString(); } public String sizeTotalsReport() { diff --git a/engine/src/main/java/org/archive/crawler/framework/Frontier.java b/engine/src/main/java/org/archive/crawler/framework/Frontier.java index 16a523b2..5280ce75 100644 --- a/engine/src/main/java/org/archive/crawler/framework/Frontier.java +++ b/engine/src/main/java/org/archive/crawler/framework/Frontier.java @@ -183,6 +183,12 @@ public interface Frontier extends Lifecycle, MultiReporter { */ public long queuedUriCount(); + + /** + * @return Number of URIs not currently queued/eligible but scheduled for future + */ + public long futureUriCount(); + /** * Ordinal position of the 'deepest' URI eligible * for crawling. Essentially, the length of the longest diff --git a/engine/src/main/java/org/archive/crawler/frontier/AbstractFrontier.java b/engine/src/main/java/org/archive/crawler/frontier/AbstractFrontier.java index 1333002a..194f89ed 100644 --- a/engine/src/main/java/org/archive/crawler/frontier/AbstractFrontier.java +++ b/engine/src/main/java/org/archive/crawler/frontier/AbstractFrontier.java @@ -147,7 +147,7 @@ public abstract class AbstractFrontier /** size of the 'outbound' mediation queue between manager thread * and toethreads */ - int outboundQueueCapacity = 50; + int outboundQueueCapacity = 200; public int getOutboundQueueCapacity() { return this.outboundQueueCapacity; } @@ -156,7 +156,7 @@ public abstract class AbstractFrontier } /** size of the inbound queue as multiple of the outbound queue */ - int inboundQueueMultiple = 3; + int inboundQueueMultiple = 10; public int getInboundQueueMultiple() { return this.inboundQueueMultiple; } @@ -252,6 +252,8 @@ public abstract class AbstractFrontier /** total URIs queued to be visited */ protected AtomicLong queuedUriCount = new AtomicLong(0); + protected AtomicLong futureUriCount = new AtomicLong(0); + protected AtomicLong succeededFetchCount = new AtomicLong(0); protected AtomicLong failedFetchCount = new AtomicLong(0); @@ -446,7 +448,7 @@ public abstract class AbstractFrontier } else { break; } - } + } } /** @@ -759,7 +761,14 @@ public abstract class AbstractFrontier public long queuedUriCount() { return queuedUriCount.get(); } - + + /* (non-Javadoc) + * @see org.archive.crawler.framework.Frontier#futureUriCount() + */ + public long futureUriCount() { + return futureUriCount.get(); + } + /** * (non-Javadoc) * diff --git a/engine/src/main/java/org/archive/crawler/frontier/BdbFrontier.java b/engine/src/main/java/org/archive/crawler/frontier/BdbFrontier.java index 4ff104ab..657cd8b7 100644 --- a/engine/src/main/java/org/archive/crawler/frontier/BdbFrontier.java +++ b/engine/src/main/java/org/archive/crawler/frontier/BdbFrontier.java @@ -216,6 +216,7 @@ implements Checkpointable, BeanNameAware { JSONObject json = new JSONObject(); try { json.put("queuedUriCount", queuedUriCount.get()); + json.put("futureUriCount", futureUriCount.get()); json.put("succeededFetchCount", succeededFetchCount.get()); json.put("failedFetchCount", failedFetchCount.get()); json.put("disregardedUriCount", disregardedUriCount.get()); @@ -243,6 +244,7 @@ implements Checkpointable, BeanNameAware { JSONObject json = recoveryCheckpoint.loadJson(beanName); try { queuedUriCount.set(json.getLong("queuedUriCount")); + futureUriCount.set(json.getLong("futureUriCount")); succeededFetchCount.set(json.getLong("succeededFetchCount")); failedFetchCount.set(json.getLong("failedFetchCount")); disregardedUriCount.set(json.getLong("disregardedUriCount")); @@ -293,6 +295,9 @@ implements Checkpointable, BeanNameAware { // may exist snoozedClassQueues = new DelayQueue(); + this.futureUris = bdb.getStoredMap( + "futureUris", Long.class, CrawlURI.class, true, recoveryCheckpoint!=null); + // initialize master map in which other queues live this.pendingUris = createMultipleWorkQueues(); } diff --git a/engine/src/main/java/org/archive/crawler/frontier/WorkQueueFrontier.java b/engine/src/main/java/org/archive/crawler/frontier/WorkQueueFrontier.java index 4309d0bf..e942f518 100644 --- a/engine/src/main/java/org/archive/crawler/frontier/WorkQueueFrontier.java +++ b/engine/src/main/java/org/archive/crawler/frontier/WorkQueueFrontier.java @@ -68,6 +68,7 @@ import org.springframework.context.ApplicationContext; import org.springframework.context.ApplicationContextAware; import org.springframework.context.support.AbstractApplicationContext; +import com.sleepycat.collections.StoredSortedMap; import com.sleepycat.je.DatabaseException; /** @@ -211,6 +212,8 @@ implements Closeable, */ transient protected DelayQueue snoozedClassQueues; + protected StoredSortedMap futureUris; + transient protected WorkQueue longestActiveQueue = null; protected int highestPrecedenceWaiting = Integer.MAX_VALUE; @@ -584,6 +587,8 @@ implements Closeable, assert Thread.currentThread() == managerThread; // wake any snoozed queues wakeQueues(); + // consider rescheduled URIS + checkFutures(); // activate enough inactive queues to fill outbound int activationsWanted = outbound.remainingCapacity() - readyClassQueues.size(); @@ -671,7 +676,6 @@ implements Closeable, continue findauri; } } - } if(inProcessQueues.size()==0) { @@ -689,7 +693,23 @@ implements Closeable, return null; } - + /** + * Check for any future-scheduled URIs now eligible for reenqueuing + */ + protected void checkFutures() { + assert Thread.currentThread() == managerThread; + // TODO: consider only checking this every set interval + Iterator iter = + futureUris.headMap(System.currentTimeMillis()) + .values().iterator(); + while(iter.hasNext()) { + CrawlURI curi = iter.next(); + curi.setRescheduleTime(-1); // unless again set elsewhere + iter.remove(); + futureUriCount.decrementAndGet(); + receive(curi); + } + } /** * Activate an inactive queue, if any are available. @@ -981,12 +1001,18 @@ implements Closeable, reenqueueQueue(wq); } - curi.stripToMinimal(); - curi.processingCleanup(); - + if(curi.getRescheduleTime()>0) { + // marked up for forced-revisit at a set time + curi.processingCleanup(); + futureUris.put(curi.getRescheduleTime(),curi); + futureUriCount.incrementAndGet(); + } else { + curi.stripToMinimal(); + curi.processingCleanup(); + } } - private boolean includesRetireDirective(CrawlURI curi) { + protected boolean includesRetireDirective(CrawlURI curi) { return curi.containsDataKey(A_FORCE_RETIRE) && (Boolean)curi.getData().get(A_FORCE_RETIRE); } @@ -1519,7 +1545,8 @@ implements Closeable, public boolean isEmpty() { return queuedUriCount.get() == 0 && (uriUniqFilter == null || uriUniqFilter.pending() == 0) - && (inbound == null || inbound.isEmpty()); + && (inbound == null || inbound.isEmpty()) + && futureUriCount.get() == 0; } /* (non-Javadoc) diff --git a/engine/src/main/java/org/archive/crawler/postprocessor/ReschedulingProcessor.java b/engine/src/main/java/org/archive/crawler/postprocessor/ReschedulingProcessor.java new file mode 100644 index 00000000..13631c30 --- /dev/null +++ b/engine/src/main/java/org/archive/crawler/postprocessor/ReschedulingProcessor.java @@ -0,0 +1,73 @@ +/* + * This file is part of the Heritrix web crawler (crawler.archive.org). + * + * Licensed to the Internet Archive (IA) by one or more individual + * contributors. + * + * The IA 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.archive.crawler.postprocessor; + +import org.archive.modules.CrawlURI; +import org.archive.modules.Processor; + +/** + * The most simple forced-rescheduling step possible: use a local + * setting (perhaps overlaid to vary based on the URI) to set an exact + * future reschedule time, as a delay from now. Unless the + * reschedulDelaySeconds value is changed from its default, URIs + * are not rescheduled. + * + * @author gojomo + * @version $Date: 2009-11-16 22:10:42 -0800 (Mon, 16 Nov 2009) $, $Revision: 6665 $ + */ +public class ReschedulingProcessor extends Processor { + /** + * amount of time to wait before forcing a URI to be rescheduled + * default of -1 means "don't reschedule" + */ + { + setRescheduleDelaySeconds(-1L); + } + public long getRescheduleDelaySeconds() { + return (Long) kp.get("rescheduleDelaySeconds"); + } + public void setRescheduleDelaySeconds(long rescheduleDelaySeconds) { + kp.put("rescheduleDelaySeconds",rescheduleDelaySeconds); + } + + public ReschedulingProcessor() { + super(); + } + + @Override + protected boolean shouldProcess(CrawlURI curi) { + return true; + } + + @Override + protected void innerProcess(CrawlURI curi) { + if(curi.isPrerequisite()) { + // never resched prereqs; they get rescheduled as needed + curi.setRescheduleTime(-1); + return; + } + long rds = getRescheduleDelaySeconds(); + if(rds>0) { + curi.setRescheduleTime(System.currentTimeMillis()+(1000*rds)); + } else { + curi.setRescheduleTime(-1); + } + } +} diff --git a/engine/src/main/java/org/archive/crawler/reporting/CrawlStatSnapshot.java b/engine/src/main/java/org/archive/crawler/reporting/CrawlStatSnapshot.java index d65373df..cb1ec95b 100644 --- a/engine/src/main/java/org/archive/crawler/reporting/CrawlStatSnapshot.java +++ b/engine/src/main/java/org/archive/crawler/reporting/CrawlStatSnapshot.java @@ -37,6 +37,7 @@ public class CrawlStatSnapshot { public long discoveredUriCount; public long queuedUriCount; + public long futureUriCount; public long finishedUriCount; public long downloadedUriCount; public long downloadFailures; @@ -74,6 +75,7 @@ public class CrawlStatSnapshot { discoveredUriCount = controller.getFrontier().discoveredUriCount(); finishedUriCount = controller.getFrontier().finishedUriCount(); queuedUriCount = controller.getFrontier().queuedUriCount(); + futureUriCount = controller.getFrontier().futureUriCount(); downloadFailures = controller.getFrontier().failedFetchCount(); downloadDisregards = controller.getFrontier().disregardedUriCount(); diff --git a/modules/src/main/java/org/archive/modules/CrawlURI.java b/modules/src/main/java/org/archive/modules/CrawlURI.java index e91fedde..708c10da 100644 --- a/modules/src/main/java/org/archive/modules/CrawlURI.java +++ b/modules/src/main/java/org/archive/modules/CrawlURI.java @@ -1826,4 +1826,15 @@ implements MultiReporter, Serializable, OverlayContext { public CrawlURI getFullVia() { return fullVia; } + + /** + * A future time at which this CrawlURI should be reenqueued. + */ + protected long rescheduleTime = -1; + public void setRescheduleTime(long time) { + this.rescheduleTime = time; + } + public long getRescheduleTime() { + return this.rescheduleTime; + } }