From bd6607739d02f0847d56273faa0950eb301c0297 Mon Sep 17 00:00:00 2001 From: Noah Levitt Date: Mon, 13 Feb 2017 17:29:13 -0800 Subject: [PATCH] KafkaCrawlLogFeed had been using lots of heap because each callback instance was keeping around a CrawlURI and in some cases there could be many of thousands of callbacks waiting around; instead, just keep some stats and log them periodically --- .../postprocessor/KafkaCrawlLogFeed.java | 23 ++++++++++++------- 1 file changed, 15 insertions(+), 8 deletions(-) diff --git a/contrib/src/main/java/org/archive/modules/postprocessor/KafkaCrawlLogFeed.java b/contrib/src/main/java/org/archive/modules/postprocessor/KafkaCrawlLogFeed.java index 690e6689..90caf163 100644 --- a/contrib/src/main/java/org/archive/modules/postprocessor/KafkaCrawlLogFeed.java +++ b/contrib/src/main/java/org/archive/modules/postprocessor/KafkaCrawlLogFeed.java @@ -164,6 +164,9 @@ public class KafkaCrawlLogFeed extends Processor implements Lifecycle { } } + String rateStr = String.format("%0.01f", 0.01 * stats.errors / stats.total); + logger.info("final error count: " + stats.errors + "/" + stats.total + " (" + rateStr + "%)"); + if (kafkaProducer != null) { kafkaProducer.close(); kafkaProducer = null; @@ -224,25 +227,29 @@ public class KafkaCrawlLogFeed extends Processor implements Lifecycle { return kafkaProducer; } - protected static class KafkaResultCallback implements Callback { - private CrawlURI curi; - - public KafkaResultCallback(CrawlURI curi) { - this.curi = curi; - } + protected final class StatsCallback implements Callback { + public long errors = 0l; + public long total = 0l; @Override public void onCompletion(RecordMetadata metadata, Exception exception) { + total++; if (exception != null) { - logger.warning("kafka delivery failed for " + curi + " - " + exception); + errors++; + } + + if (total % 10000 == 0) { + String rateStr = String.format("%0.01f", 0.01 * errors / total); + logger.info("error count so far: " + errors + "/" + total + " (" + rateStr + "%)"); } } } + protected StatsCallback stats = new StatsCallback(); @Override protected void innerProcess(CrawlURI curi) throws InterruptedException { byte[] message = buildMessage(curi); ProducerRecord producerRecord = new ProducerRecord(getTopic(), message); - kafkaProducer().send(producerRecord, new KafkaResultCallback(curi)); + kafkaProducer().send(producerRecord, stats); } }