oscerd commented on code in PR #27142:
URL: https://github.com/apache/camel/pull/27142#discussion_r4154140956


##########
components/camel-kafka/src/main/java/org/apache/camel/component/kafka/KafkaProducer.java:
##########
@@ -506,16 +509,26 @@ private void doSend(Object key, ProducerRecord<Object, 
Object> record, KafkaProd
                     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) {

Review Comment:
   Fixed the comment in c8d9fad. Dropped the `max.block.ms` / buffer-exhaustion 
`TimeoutException` and `RecordTooLargeException` examples — those are 
`ApiException`s and, as you note, Kafka reports them through the callback where 
`onCompletion` balances the count — and replaced them with the cases that 
actually throw synchronously: a `SerializationException` (or other non-API 
`KafkaException`), a closed-producer `IllegalStateException`, and 
`InterruptException`. Thanks for the precise breakdown.
   
   _Claude Code on behalf of oscerd_
   



##########
components/camel-kafka/src/main/java/org/apache/camel/component/kafka/KafkaProducer.java:
##########
@@ -506,16 +509,26 @@ private void doSend(Object key, ProducerRecord<Object, 
Object> record, KafkaProd
                     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 (e.g. buffer exhaustion / 
max.block.ms timeout, a serialization error, or a
+            // closed producer), so no Kafka callback will ever fire for this 
record. 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 (CAMEL-24783).
+            cb.decrement();

Review Comment:
   Good catch. I've documented it in the `doSend` catch comment as a known, 
pre-existing exposure: with a transactional producer whose transaction is 
already in an error state, `transactionManager.maybeAddPartition()` can throw 
*after* the record's callback is registered, so decrementing here undercounts 
by one and routing can continue one callback early. It only arises for an 
already-failed transaction and matches the prior behaviour, so I left the code 
as-is with the comment rather than widening the fix here. Happy to file a 
follow-up JIRA to track it if you'd prefer.
   
   _Claude Code on behalf of oscerd_
   



##########
components/camel-kafka/src/test/java/org/apache/camel/component/kafka/KafkaProducerTest.java:
##########
@@ -191,6 +192,43 @@ public void processAsyncSendsMessageWithException() {
         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
+        Producer kp = producer.getKafkaProducer();
+        Future future = Mockito.mock(Future.class);
+        Mockito.when(kp.send(any(ProducerRecord.class), any(Callback.class)))
+                .thenReturn(future)
+                .thenThrow(new ApiException());

Review Comment:
   Applied in c8d9fad — the test now `.thenThrow(new 
SerializationException("boom"))` and asserts 
`isA(SerializationException.class)`, with the import added, so it exercises 
what `send(record, callback)` actually throws on this path. I kept the change 
scoped to this test; the pre-existing `processAsync*` tests share the 
ApiException-thrown pattern, but I left them out to keep the PR focused (they'd 
be a good small follow-up cleanup).
   
   _Claude Code on behalf of oscerd_
   



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to