This is an automated email from the ASF dual-hosted git repository.
oscerd pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/camel.git
The following commit(s) were added to refs/heads/main by this push:
new 1ba71e2bee86 CAMEL-24783: camel-kafka - don't continue routing while
batch sends are still in flight (#27142)
1ba71e2bee86 is described below
commit 1ba71e2bee86fe23e6b90141072d16ffa5f00310
Author: Andrea Cosentino <[email protected]>
AuthorDate: Thu Oct 1 14:15:28 2026 +0200
CAMEL-24783: camel-kafka - don't continue routing while batch sends are
still in flight (#27142)
In the async batch/iterator producer path
(KafkaProducer.processIterableAsync ->
doSend) records are dispatched one at a time. If a later element failed to
dispatch - kafkaProducer.send() throwing synchronously (a
SerializationException,
a closed-producer IllegalStateException or an InterruptException), or the
record
iterator throwing (bad CamelKafkaOverrideTimestamp, header serialization) -
process() recorded the exception and immediately completed the async
callback, so
routing continued while the records already dispatched still had in-flight
Kafka
callbacks. Those callbacks then ran setException(...) and
recordMetadataList.add(...)
on an exchange that had already continued down the route (and, with exchange
pooling, may have been reset and reused).
Fix: on any dispatch failure, arm completion the same way the success path
does
(producerCallBack.allSent()) instead of completing in place. allSent()
releases
the initial hold and, while sends are still in flight, defers done() to the
last
callback - so routing continues exactly once, after the in-flight sends have
settled. When nothing is in flight it completes immediately, unchanged.
doSend now undoes its increment if send() throws (no Kafka callback fires
in that
case), via the new KafkaProducerCallBack.decrement(), so the counter can
reach zero.
Related: CAMEL-24779 (single-message path), CAMEL-24780 (transaction begin).
Co-authored-by: Claude Opus 4.8 <[email protected]>
Signed-off-by: Andrea Cosentino <[email protected]>
---
.../camel/component/kafka/KafkaProducer.java | 44 ++++++++++++++++------
.../producer/support/KafkaProducerCallBack.java | 10 +++++
.../camel/component/kafka/KafkaProducerTest.java | 42 +++++++++++++++++++++
3 files changed, 85 insertions(+), 11 deletions(-)
diff --git
a/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/KafkaProducer.java
b/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/KafkaProducer.java
index 30c7fc744dfb..8bb4129de392 100755
---
a/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/KafkaProducer.java
+++
b/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/KafkaProducer.java
@@ -478,10 +478,13 @@ public class KafkaProducer extends DefaultAsyncProducer
implements RouteIdAware
return producerCallBack.allSent();
} catch (Exception e) {
exchange.setException(e);
+ // Do not continue routing immediately. In the batch/iterator
path, records already dispatched before the
+ // failure have in-flight Kafka callbacks, and completing the
exchange now would let those late callbacks
+ // mutate a continued - and, with exchange pooling, possibly
recycled - exchange. Arm completion instead:
+ // allSent() releases the initial hold and, when sends are still
in flight, defers done() to the last
+ // callback; when nothing is in flight (the common non-batch
failure) it completes here (CAMEL-24783).
+ return producerCallBack.allSent();
}
-
- callback.done(true);
- return true;
}
private void processIterableAsync(Exchange exchange, KafkaProducerCallBack
producerCallBack, Message message) {
@@ -508,16 +511,35 @@ public class KafkaProducer extends DefaultAsyncProducer
implements RouteIdAware
record.key());
}
- if (key != null) {
- KafkaProducerMetadataCallBack metadataCallBack = new
KafkaProducerMetadataCallBack(
- key, configuration.isRecordMetadata());
+ try {
+ if (key != null) {
+ KafkaProducerMetadataCallBack metadataCallBack = new
KafkaProducerMetadataCallBack(
+ key, configuration.isRecordMetadata());
- // make sure to cb is last in the order here
- DelegatingCallback delegatingCallback = new
DelegatingCallback(metadataCallBack, cb);
+ // make sure to cb is last in the order here
+ DelegatingCallback delegatingCallback = new
DelegatingCallback(metadataCallBack, cb);
- kafkaProducer.send(record, delegatingCallback);
- } else {
- kafkaProducer.send(record, cb);
+ kafkaProducer.send(record, delegatingCallback);
+ } else {
+ kafkaProducer.send(record, cb);
+ }
+ } catch (RuntimeException dispatchFailure) {
+ // send() threw synchronously rather than reporting through the
callback: a SerializationException (or other
+ // non-API KafkaException), an IllegalStateException from a closed
producer, or an InterruptException. In
+ // those cases no Kafka callback fires for this record. (An
ApiException - a max.block.ms/buffer-exhaustion
+ // TimeoutException, RecordTooLargeException, etc. - is NOT thrown
here: Kafka reports it through the
+ // callback, which balances the count via onCompletion.) Undo the
increment above to keep the completion
+ // counter accurate; otherwise a mid-batch failure would leave it
above zero and routing would never
+ // continue. The exception propagates to process(), which records
it and arms completion for the records
+ // already in flight.
+ //
+ // Narrow known exposure, left as-is: with a transactional
producer whose transaction is already in an
+ // error state, transactionManager.maybeAddPartition() can throw
*after* the record's callback has been
+ // registered, so that callback still fires later (on abort).
Decrementing here then undercounts by one and
+ // routing can continue one callback early. This matches the
pre-existing behaviour and only arises for an
+ // already-failed transaction (CAMEL-24783).
+ cb.decrement();
+ throw dispatchFailure;
}
}
diff --git
a/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/producer/support/KafkaProducerCallBack.java
b/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/producer/support/KafkaProducerCallBack.java
index c97ada3e8b92..1309dd4b4f23 100644
---
a/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/producer/support/KafkaProducerCallBack.java
+++
b/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/producer/support/KafkaProducerCallBack.java
@@ -59,6 +59,16 @@ public final class KafkaProducerCallBack implements Callback
{
count.incrementAndGet();
}
+ /**
+ * Undoes an {@link #increment()} for a record whose dispatch failed
synchronously, so that no Kafka callback will
+ * ever fire for it. Unlike {@link #onCompletion} it never continues
routing: it is only called while the initial
+ * hold is still in place (a mid-batch dispatch failure in {@code
KafkaProducer.doSend}), so the counter cannot
+ * reach zero here (CAMEL-24783).
+ */
+ public void decrement() {
+ count.decrementAndGet();
+ }
+
public boolean allSent() {
if (count.decrementAndGet() == 0) {
LOG.trace("All messages sent, continue routing.");
diff --git
a/components/camel-kafka/src/test/java/org/apache/camel/component/kafka/KafkaProducerTest.java
b/components/camel-kafka/src/test/java/org/apache/camel/component/kafka/KafkaProducerTest.java
index ffa317f46e1d..e842694e5edd 100755
---
a/components/camel-kafka/src/test/java/org/apache/camel/component/kafka/KafkaProducerTest.java
+++
b/components/camel-kafka/src/test/java/org/apache/camel/component/kafka/KafkaProducerTest.java
@@ -51,11 +51,13 @@ import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.clients.producer.RecordMetadata;
import org.apache.kafka.common.errors.ApiException;
+import org.apache.kafka.common.errors.SerializationException;
import org.junit.jupiter.api.Test;
import org.mockito.ArgumentCaptor;
import org.mockito.Mockito;
import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertInstanceOf;
import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertNull;
@@ -192,6 +194,46 @@ public class KafkaProducerTest {
assertRecordMetadataExists();
}
+ @Test
+ void
processAsyncMidBatchDispatchFailureDefersRoutingUntilInflightSendsComplete() {
+ // CAMEL-24783: when a later record in a batch fails to dispatch, the
records already dispatched still have
+ // in-flight Kafka callbacks. Routing must not continue until those
callbacks have run, otherwise they would
+ // mutate a continued - and, with exchange pooling, possibly recycled
- exchange.
+ endpoint.getConfiguration().setTopic("sometopic");
+ Mockito.when(exchange.getIn()).thenReturn(in);
+ Mockito.when(exchange.getMessage()).thenReturn(in);
+
+ // the first record is accepted (its callback stays in flight), the
second fails to dispatch mid-batch.
+ // Kafka's send(record, callback) throws synchronously for a
SerializationException (and a closed-producer
+ // IllegalStateException or an InterruptException), whereas an
ApiException is reported through the callback
+ // instead - so a synchronously thrown SerializationException is what
this path actually sees.
+ Producer kp = producer.getKafkaProducer();
+ Future future = Mockito.mock(Future.class);
+ Mockito.when(kp.send(any(ProducerRecord.class), any(Callback.class)))
+ .thenReturn(future)
+ .thenThrow(new SerializationException("boom"));
+
+ ArrayNode node = JsonNodeFactory.instance.arrayNode();
+ node.add(1);
+ node.add(2);
+ in.setBody(node);
+
+ boolean sync = producer.process(exchange, callback);
+
+ // the dispatch failure is recorded, but routing is deferred while the
first send is still in flight
+
Mockito.verify(exchange).setException(isA(SerializationException.class));
+ assertFalse(sync);
+ Mockito.verify(callback, Mockito.never()).done(Mockito.anyBoolean());
+
+ // the first record now completes on the Kafka sender thread; only now
may routing continue, exactly once
+ ArgumentCaptor<Callback> callBackCaptor =
ArgumentCaptor.forClass(Callback.class);
+ Mockito.verify(kp, Mockito.times(2)).send(any(ProducerRecord.class),
callBackCaptor.capture());
+ callBackCaptor.getAllValues().get(0).onCompletion(new
RecordMetadata(null, 0, 0, 0, 0, 0), null);
+
+ // done() is delivered from the worker pool
+ Mockito.verify(callback, Mockito.timeout(2000)).done(eq(false));
+ }
+
@Test
public void processAsyncCompletesCallbackWhenBeginTransactionFails()
throws Exception {
// CAMEL-24780: a failure to begin the transaction must set the
exception and complete the async