From a1a72325300a29c425bb8cda65ef6676223df5fa Mon Sep 17 00:00:00 2001 From: Noah Levitt Date: Wed, 14 Jun 2017 18:31:39 -0700 Subject: [PATCH] update for new data model queued_url->uncrawled_url and some new fields --- .../postprocessor/TroughCrawlLogFeed.java | 120 +++++++++++------- 1 file changed, 72 insertions(+), 48 deletions(-) diff --git a/contrib/src/main/java/org/archive/modules/postprocessor/TroughCrawlLogFeed.java b/contrib/src/main/java/org/archive/modules/postprocessor/TroughCrawlLogFeed.java index 22a08da5..57151309 100644 --- a/contrib/src/main/java/org/archive/modules/postprocessor/TroughCrawlLogFeed.java +++ b/contrib/src/main/java/org/archive/modules/postprocessor/TroughCrawlLogFeed.java @@ -51,6 +51,7 @@ import org.springframework.context.Lifecycle; * timestamp DATETIME, * status_code INTEGER, * size BIGINT, + * payload_size BIGINT, * url VARCHAR(4000), * hop_path VARCHAR(255), * is_seed_redirect INTEGER(1), @@ -61,13 +62,15 @@ import org.springframework.context.Lifecycle; * is_duplicate INTEGER(1), * warc_filename VARCHAR(255), * warc_offset VARCHAR(255), + * warc_content_bytes BIGINT, * host VARCHAR(255)); - * - * CREATE TABLE queued_url( + * + * CREATE TABLE uncrawled_url( * id INTEGER PRIMARY KEY AUTOINCREMENT, * timestamp DATETIME, * url VARCHAR(4000), * hop_path VARCHAR(255), + * status_code INTEGER, * via VARCHAR(255), * seed VARCHAR(4000), * host VARCHAR(255)); @@ -94,8 +97,10 @@ public class TroughCrawlLogFeed extends Processor implements Lifecycle { protected static final int BATCH_MAX_TIME_MS = 20 * 1000; protected static final int BATCH_MAX_SIZE = 500; - protected List batch = new ArrayList(); - protected long batchLastTime = System.currentTimeMillis(); + protected List crawledBatch = new ArrayList(); + protected long crawledBatchLastTime = System.currentTimeMillis(); + protected List uncrawledBatch = new ArrayList(); + protected long uncrawledBatchLastTime = System.currentTimeMillis(); protected Frontier frontier; public Frontier getFrontier() { @@ -138,39 +143,22 @@ public class TroughCrawlLogFeed extends Processor implements Lifecycle { if (!isRunning) { return; } - if (!batch.isEmpty()) { - postBatch(); + if (!crawledBatch.isEmpty()) { + postCrawledBatch(); } if (frontier instanceof BdbFrontier) { Closure closure = new Closure() { public void execute(Object o) { - CrawlURI curi = (CrawlURI) o; - batch.add("(" - + sqlValue(new Date()) + ", " - + sqlValue(curi) + ", " - + sqlValue(curi.getPathFromSeed()) + ", " - + sqlValue(curi.getVia()) + ", " - + sqlValue(curi.getSourceTag()) + ", " - + sqlValue(serverCache.getHostFor(curi.getUURI()).getHostName()) + ")"); - - if (batch.size() >= BATCH_MAX_SIZE) { - String sql = "insert into queued_url (timestamp, url, hop_path, via, seed, host) values " - + StringUtils.join(batch, ", ") + ";"; - post(sql); - batch.clear(); + try { + innerProcess((CrawlURI) o); + } catch (InterruptedException e) { } } }; logger.info("dumping " + frontier.queuedUriCount() + " queued urls to trough feed"); ((BdbFrontier) frontier).forAllPendingDo(closure); - if (!batch.isEmpty()) { - String sql = "insert into queued_url (timestamp, url, hop_path, via, seed, host) values " - + StringUtils.join(batch, ", ") + ";"; - post(sql); - batch.clear(); - } logger.info("dumped " + frontier.queuedUriCount() + " queued urls to trough feed"); } else { logger.warning("frontier is not a BdbFrontier, cannot dump queued urls to trough feed"); @@ -221,36 +209,72 @@ public class TroughCrawlLogFeed extends Processor implements Lifecycle { @Override protected void innerProcess(CrawlURI curi) throws InterruptedException { - batch.add("(" + sqlValue(new Date(curi.getFetchBeginTime())) + ", " - + sqlValue(curi.getFetchStatus()) + ", " - + sqlValue(curi.getContentSize()) + ", " - + sqlValue(curi) + ", " - + sqlValue(curi.getPathFromSeed()) + ", " - + sqlValue(curi.isSeed() && !"".equals(curi.getPathFromSeed()) ? 1 : 0) + ", " - + sqlValue(curi.getVia()) + ", " - + sqlValue(MimetypeUtils.truncate(curi.getContentType())) + ", " - + sqlValue(curi.getContentDigestSchemeString()) + ", " - + sqlValue(curi.getSourceTag()) + ", " - + sqlValue(curi.isRevisit() ? 1 : 0) + ", " - + sqlValue(curi.getExtraInfo().opt("warcFilename")) + ", " - + sqlValue(curi.getExtraInfo().opt("warcOffset")) + ", " - + sqlValue(serverCache.getHostFor(curi.getUURI()).getHostName()) + ")"); + if (curi.getFetchStatus() > 0) { + // compute warcContentBytes + long warcContentBytes; + if (curi.getExtraInfo().opt("warcFilename") != null) { + if (curi.isRevisit()) { + warcContentBytes = curi.getContentSize() - curi.getContentLength(); + } else { + warcContentBytes = curi.getContentSize(); + } + } else { + warcContentBytes = 0; + } + crawledBatch.add("(" + sqlValue(new Date(curi.getFetchBeginTime())) + ", " + + sqlValue(curi.getFetchStatus()) + ", " + + sqlValue(curi.getContentSize()) + ", " + + sqlValue(curi.getContentLength()) + ", " + + sqlValue(curi) + ", " + + sqlValue(curi.getPathFromSeed()) + ", " + + sqlValue(curi.isSeed() && !"".equals(curi.getPathFromSeed()) ? 1 : 0) + ", " + + sqlValue(curi.getVia()) + ", " + + sqlValue(MimetypeUtils.truncate(curi.getContentType())) + ", " + + sqlValue(curi.getContentDigestSchemeString()) + ", " + + sqlValue(curi.getSourceTag()) + ", " + + sqlValue(curi.isRevisit() ? 1 : 0) + ", " + + sqlValue(curi.getExtraInfo().opt("warcFilename")) + ", " + + sqlValue(curi.getExtraInfo().opt("warcOffset")) + ", " + + sqlValue(warcContentBytes) + + sqlValue(serverCache.getHostFor(curi.getUURI()).getHostName()) + ")"); + if (crawledBatch.size() >= BATCH_MAX_SIZE || System.currentTimeMillis() - crawledBatchLastTime > BATCH_MAX_TIME_MS) { + postCrawledBatch(); + } + } else { + uncrawledBatch.add("(" + + sqlValue(new Date()) + ", " + + sqlValue(curi) + ", " + + sqlValue(curi.getPathFromSeed()) + ", " + + sqlValue(curi.getFetchStatus()) + ", " + + sqlValue(curi.getVia()) + ", " + + sqlValue(curi.getSourceTag()) + ", " + + sqlValue(serverCache.getHostFor(curi.getUURI()).getHostName()) + ")"); + if (uncrawledBatch.size() >= BATCH_MAX_SIZE || System.currentTimeMillis() - uncrawledBatchLastTime > BATCH_MAX_TIME_MS) { + postUncrawledBatch(); + } - if (batch.size() >= BATCH_MAX_SIZE || System.currentTimeMillis() - batchLastTime > BATCH_MAX_TIME_MS) { - postBatch(); } } - protected void postBatch() { + protected void postCrawledBatch() { String sql = "insert into crawled_url (" - + "timestamp, status_code, size, url, hop_path, is_seed_redirect, " + + "timestamp, status_code, size, payload_size url, hop_path, is_seed_redirect, " + "via, mimetype, content_digest, seed, is_duplicate, warc_filename, " - + "warc_offset, host) values " - + StringUtils.join(batch, ", ") + + "warc_offset, warc_content_bytes, host) values " + + StringUtils.join(crawledBatch, ", ") + ";"; post(sql); - batchLastTime = System.currentTimeMillis(); - batch.clear(); + crawledBatchLastTime = System.currentTimeMillis(); + crawledBatch.clear(); } + protected void postUncrawledBatch() { + String sql = "insert into uncrawled_url (" + + "timestamp, url, hop_path, status_code, via, seed, host) values " + + StringUtils.join(uncrawledBatch, ", ") + + ";"; + post(sql); + uncrawledBatchLastTime = System.currentTimeMillis(); + uncrawledBatch.clear(); + } }