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


##########
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:
   A narrow edge case, not a blocker. With a transactional producer, 
`transactionManager.maybeAddPartition()` runs *after* `accumulator.append()` 
inside Kafka's `doSend`. If it throws (`maybeFailWithError`, or 
`IllegalStateException` when the transaction isn't `IN_TRANSACTION`), the 
record's callback is already registered and will still fire later, typically 
when the batch is aborted. Here we would `decrement()` anyway, so the counter 
undercounts by one: routing can continue one callback early, and that late 
callback then touches a continued exchange.
   
   This only happens when the transaction is already in an error state, and the 
old code had the same exposure, so it's fine to handle it separately. A code 
comment or a follow-up JIRA would be enough.



##########
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:
   Small factual point about the comment below (and the PR description). In 
kafka-clients 4.3.1, `KafkaProducer.doSend` catches `ApiException` itself: it 
calls the callback directly and returns a `FutureFailure` without throwing. 
That covers `TimeoutException` from `max.block.ms`/buffer exhaustion and 
`RecordTooLargeException`. Only non-API failures are thrown to the caller: 
`KafkaException` such as `SerializationException`, `IllegalStateException` for 
a closed producer, and `InterruptException`.
   
   The fix handles both cases correctly. In the callback case the count is 
balanced by `onCompletion`. It would still be good to drop "buffer exhaustion / 
max.block.ms timeout" from the examples here so the next reader isn't misled.



##########
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:
   Real Kafka never throws an `ApiException` from `send(record, callback)`; it 
calls the callback instead (see the comment on `doSend`). The existing tests 
use the same pattern, but a synchronously thrown `SerializationException` is 
what this path actually sees in production:
   
   ```suggestion
                   .thenThrow(new SerializationException("boom"));
   ```
   
   You'd also need `import 
org.apache.kafka.common.errors.SerializationException;` and to change the 
`isA(ApiException.class)` check to `isA(SerializationException.class)`.



-- 
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