From c611c46c8a931edc295e8ceae4deed846cc8cd58 Mon Sep 17 00:00:00 2001 From: Nicholas Clarke Date: Tue, 21 Nov 2017 13:44:54 +0100 Subject: [PATCH] Merged frontier-management with upgrade to bdb 7 --- .../archive/crawler/framework/CrawlJob.java | 53 +++++++++++++++++++ .../archive/crawler/framework/Frontier.java | 14 +++++ .../archive/crawler/frontier/BdbFrontier.java | 27 ++++++++++ .../frontier/BdbMultipleWorkQueues.java | 27 ++++++++++ .../crawler/prefetch/QuotaEnforcerTest.java | 19 +++++++ 5 files changed, 140 insertions(+) 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 c9aa8bb8..02bdbe2e 100644 --- a/engine/src/main/java/org/archive/crawler/framework/CrawlJob.java +++ b/engine/src/main/java/org/archive/crawler/framework/CrawlJob.java @@ -20,18 +20,25 @@ package org.archive.crawler.framework; import java.io.BufferedReader; +import java.io.BufferedWriter; import java.io.File; import java.io.FileInputStream; +import java.io.FileOutputStream; import java.io.IOException; import java.io.InputStreamReader; +import java.io.OutputStreamWriter; import java.io.PrintWriter; +import java.nio.charset.StandardCharsets; import java.util.Collections; import java.util.HashMap; import java.util.LinkedHashMap; import java.util.LinkedList; import java.util.List; import java.util.Map; +import java.util.Set; import java.util.TreeMap; +import java.util.concurrent.BlockingQueue; +import java.util.concurrent.Semaphore; import java.util.logging.FileHandler; import java.util.logging.Formatter; import java.util.logging.Handler; @@ -51,6 +58,7 @@ import org.apache.commons.io.FileUtils; import org.apache.commons.lang.StringUtils; import org.archive.crawler.event.CrawlStateEvent; import org.archive.crawler.framework.CrawlController.StopCompleteEvent; +import org.archive.crawler.frontier.WorkQueue; import org.archive.crawler.reporting.AlertThreadGroup; import org.archive.crawler.reporting.CrawlStatSnapshot; import org.archive.crawler.reporting.StatisticsTracker; @@ -58,6 +66,7 @@ import org.archive.spring.ConfigPath; import org.archive.spring.ConfigPathConfigurer; import org.archive.spring.PathSharingContext; import org.archive.util.ArchiveUtils; +import org.archive.util.ObjectIdentityCache; import org.archive.util.TextUtils; import org.joda.time.DateTime; import org.springframework.beans.BeanWrapperImpl; @@ -970,4 +979,48 @@ public class CrawlJob implements Comparable, ApplicationListener getAllQueues(); + + public BlockingQueue getReadyClassQueues(); + + public Set getInProcessQueues(); + } 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 5806df75..d80ffd16 100644 --- a/engine/src/main/java/org/archive/crawler/frontier/BdbFrontier.java +++ b/engine/src/main/java/org/archive/crawler/frontier/BdbFrontier.java @@ -23,7 +23,9 @@ import java.io.IOException; import java.io.PrintWriter; import java.util.Map.Entry; import java.util.Queue; +import java.util.Set; import java.util.SortedMap; +import java.util.concurrent.BlockingQueue; import java.util.concurrent.ConcurrentSkipListMap; import java.util.concurrent.DelayQueue; import java.util.concurrent.LinkedBlockingQueue; @@ -42,6 +44,7 @@ import org.archive.checkpointing.Checkpoint; import org.archive.checkpointing.Checkpointable; import org.archive.modules.CrawlURI; import org.archive.util.ArchiveUtils; +import org.archive.util.ObjectIdentityCache; import org.archive.util.Supplier; import org.json.JSONArray; import org.json.JSONException; @@ -468,4 +471,28 @@ implements Checkpointable, BeanNameAware { queueSummaries.put(key, val); } } + + @Override + public long exportPendingUris(PrintWriter writer) { + if (pendingUris == null) { + return -5L; + } + return pendingUris.exportPendingUris(writer); + } + + @Override + public ObjectIdentityCache getAllQueues() { + return allQueues; + } + + @Override + public BlockingQueue getReadyClassQueues() { + return readyClassQueues; + } + + @Override + public Set getInProcessQueues() { + return inProcessQueues; + } + } diff --git a/engine/src/main/java/org/archive/crawler/frontier/BdbMultipleWorkQueues.java b/engine/src/main/java/org/archive/crawler/frontier/BdbMultipleWorkQueues.java index 1f480c71..a45712dc 100644 --- a/engine/src/main/java/org/archive/crawler/frontier/BdbMultipleWorkQueues.java +++ b/engine/src/main/java/org/archive/crawler/frontier/BdbMultipleWorkQueues.java @@ -19,6 +19,7 @@ package org.archive.crawler.frontier; import java.io.IOException; +import java.io.PrintWriter; import java.io.UnsupportedEncodingException; import java.math.BigInteger; import java.util.ArrayList; @@ -558,4 +559,30 @@ public class BdbMultipleWorkQueues { } cursor.close(); } + + /** + * Run through all uris in the pending uris database and write them to the writer. + * @param writer destination writer for writting all the uris + * @return number of uris written to the writer + */ + public long exportPendingUris(PrintWriter writer) { + if (this.pendingUrisDB == null) { + return -6L; + } + sync(); + DatabaseEntry key = new DatabaseEntry(); + DatabaseEntry value = new DatabaseEntry(); + long uris = 0L; + Cursor cursor = pendingUrisDB.openCursor(null, null); + while (cursor.getNext(key, value, null) == OperationStatus.SUCCESS) { + if (value.getData().length == 0) { + continue; + } + CrawlURI item = (CrawlURI) crawlUriBinding.entryToObject(value); + writer.println(item.toString()); + ++uris; + } + cursor.close(); + return uris; + } } diff --git a/engine/src/test/java/org/archive/crawler/prefetch/QuotaEnforcerTest.java b/engine/src/test/java/org/archive/crawler/prefetch/QuotaEnforcerTest.java index 1db16743..8f209968 100644 --- a/engine/src/test/java/org/archive/crawler/prefetch/QuotaEnforcerTest.java +++ b/engine/src/test/java/org/archive/crawler/prefetch/QuotaEnforcerTest.java @@ -24,6 +24,8 @@ import java.io.IOException; import java.io.PrintWriter; import java.util.HashMap; import java.util.Map; +import java.util.Set; +import java.util.concurrent.BlockingQueue; import javax.management.openmbean.CompositeData; @@ -32,6 +34,7 @@ import org.archive.crawler.framework.CrawlerProcessorTestBase; import org.archive.crawler.framework.Frontier; import org.archive.crawler.framework.Frontier.FrontierGroup; import org.archive.crawler.frontier.FrontierJournal; +import org.archive.crawler.frontier.WorkQueue; import org.archive.modules.CoreAttributeConstants; import org.archive.modules.CrawlURI; import org.archive.modules.ProcessResult; @@ -289,6 +292,22 @@ public class QuotaEnforcerTest extends CrawlerProcessorTestBase { @Override public void endDisposition() { } + @Override + public long exportPendingUris(PrintWriter writer) { + return 0; + } + @Override + public ObjectIdentityCache getAllQueues() { + return null; + } + @Override + public BlockingQueue getReadyClassQueues() { + return null; + } + @Override + public Set getInProcessQueues() { + return null; + } } // separate methods to make it easier to know what failed