From 30e512e5ce4beddd34879c79614aec576f877e74 Mon Sep 17 00:00:00 2001 From: Drew Date: Mon, 15 Jul 2013 11:06:02 -0700 Subject: [PATCH] Modified the async producer so it re-queues failed batches. This will cause the queue to fill up under backpressure. --- .../scala/kafka/producer/async/ProducerSendThread.scala | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) diff --git a/core/src/main/scala/kafka/producer/async/ProducerSendThread.scala b/core/src/main/scala/kafka/producer/async/ProducerSendThread.scala index 2b41a4996ce4f..f0d4d120b9488 100644 --- a/core/src/main/scala/kafka/producer/async/ProducerSendThread.scala +++ b/core/src/main/scala/kafka/producer/async/ProducerSendThread.scala @@ -103,7 +103,11 @@ class ProducerSendThread[K,V](val threadName: String, if(size > 0) handler.handle(events) }catch { - case e => error("Error in handling batch of " + size + " events", e) + case e => + for (message <- events) { + queue.offer(message) + } + error("Re-queued batch of " + size + " events. Queue depth " + queue.size) } }