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]