From a0f70d0865198908e2fa700772ca0694ee079a96 Mon Sep 17 00:00:00 2001 From: Noah Levitt Date: Mon, 10 Nov 2014 13:55:13 -0800 Subject: [PATCH] Revert "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)" This reverts commit d2f70e5c9767edeadabfbc7424a1eeddc712b2ae. Conflicts: contrib/src/main/java/org/archive/modules/postprocessor/KafkaCrawlLogFeed.java --- .../postprocessor/KafkaCrawlLogFeed.java | 41 +++++++------------ 1 file changed, 15 insertions(+), 26 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 fe20f2f4..c1fd8ffc 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,41 +164,30 @@ public class KafkaCrawlLogFeed extends Processor implements Lifecycle { super.stop(); } - transient protected KafkaProducer kafkaProducer; - protected KafkaProducer kafkaProducer() { + transient protected Producer kafkaProducer; + protected Producer kafkaProducer() { if (kafkaProducer == null) { synchronized (this) { if (kafkaProducer == null) { Properties props = new Properties(); - props.put("bootstrap.servers", getBrokerList()); - props.put("acks", "1"); + 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("producer.type", "async"); - kafkaProducer = new KafkaProducer(props); + ProducerConfig config = new ProducerConfig(props); + kafkaProducer = new Producer(config); } } } 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); - ProducerRecord producerRecord = new ProducerRecord(getTopic(), message); - kafkaProducer().send(producerRecord, new KafkaResultCallback(curi)); + KeyedMessage keyedMessage = new KeyedMessage(getTopic(), message); + // XXX no error handling + kafkaProducer().send(keyedMessage); } }