From d2f70e5c9767edeadabfbc7424a1eeddc712b2ae Mon Sep 17 00:00:00 2001 From: Noah Levitt Date: Thu, 6 Nov 2014 19:26:32 -0800 Subject: [PATCH] switch to org.apache.kafka.clients.producer.KafkaProducer to see if we can avoid the InterruptedException at the end of the crawl (and log warnings on failure to send message) --- .../postprocessor/KafkaCrawlLogFeed.java | 39 ++++++++++++------- 1 file changed, 25 insertions(+), 14 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 c1fd8ffc..afd84b94 100644 --- a/contrib/src/main/java/org/archive/modules/postprocessor/KafkaCrawlLogFeed.java +++ b/contrib/src/main/java/org/archive/modules/postprocessor/KafkaCrawlLogFeed.java @@ -23,11 +23,11 @@ import java.util.Map; import java.util.Properties; import java.util.logging.Logger; -import kafka.javaapi.producer.Producer; -import kafka.producer.KeyedMessage; -import kafka.producer.ProducerConfig; - import org.apache.commons.collections.Closure; +import org.apache.kafka.clients.producer.Callback; +import org.apache.kafka.clients.producer.KafkaProducer; +import org.apache.kafka.clients.producer.ProducerRecord; +import org.apache.kafka.clients.producer.RecordMetadata; import org.archive.crawler.framework.Frontier; import org.archive.crawler.frontier.AbstractFrontier; import org.archive.crawler.frontier.BdbFrontier; @@ -164,30 +164,41 @@ public class KafkaCrawlLogFeed extends Processor implements Lifecycle { super.stop(); } - transient protected Producer kafkaProducer; - protected Producer kafkaProducer() { + transient protected KafkaProducer kafkaProducer; + protected KafkaProducer kafkaProducer() { if (kafkaProducer == null) { synchronized (this) { if (kafkaProducer == null) { Properties props = new Properties(); props.put("metadata.broker.list", getBrokerList()); - // XXX in asynchronous mode, we have no way to handle delivery - // problems, so it doesn't matter if they are acked or not - props.put("request.required.acks", "0"); + props.put("request.required.acks", "1"); props.put("producer.type", "async"); - ProducerConfig config = new ProducerConfig(props); - kafkaProducer = new Producer(config); + kafkaProducer = new KafkaProducer(props); } } } return kafkaProducer; } + protected static class KafkaResultCallback implements Callback { + private CrawlURI curi; + + public KafkaResultCallback(CrawlURI curi) { + this.curi = curi; + } + + @Override + public void onCompletion(RecordMetadata metadata, Exception exception) { + if (exception != null) { + logger.warning("kafka delivery failed for " + curi + " - " + exception); + } + } + } + @Override protected void innerProcess(CrawlURI curi) throws InterruptedException { byte[] message = buildMessage(curi); - KeyedMessage keyedMessage = new KeyedMessage(getTopic(), message); - // XXX no error handling - kafkaProducer().send(keyedMessage); + ProducerRecord producerRecord = new ProducerRecord(getTopic(), message); + kafkaProducer().send(producerRecord, new KafkaResultCallback(curi)); } }