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]