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

Reply via email to