oscerd opened a new pull request, #27142:
URL: https://github.com/apache/camel/pull/27142
## Problem
In the asynchronous batch/iterator producer path
(`KafkaProducer.processIterableAsync` → `doSend`),
records are dispatched one at a time. If a **later** element fails to
dispatch — `kafkaProducer.send()`
throwing synchronously (buffer exhaustion / `max.block.ms` timeout, a
serialization error, a closed
producer), or the record iterator's `next()` throwing (bad
`CamelKafkaOverrideTimestamp` conversion,
or header serialization when `batchWithIndividualHeaders=true`) —
`process()`'s catch recorded the
exception and **immediately** completed the async callback, so routing
continued.
The records already dispatched still had in-flight Kafka callbacks, though.
When those complete on the
Kafka sender thread, `KafkaProducerCallBack.onCompletion` runs
`setException(...)` and
`recordMetadataList.add(...)` on an exchange that has **already continued**
down the route — and, with
exchange pooling, may have been reset and reused. That is a data race /
use-after-continue: a late
callback can set an exception on, or mutate the `KafkaRecordMeta` header
list of, a continued or
recycled exchange.
(The completion counter never reached zero on the old failure path —
`allSent()` was skipped — so
there was no double `done()`; the mutation race is the bug.)
## 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. While sends are still in flight it
returns `false` and the
**last** in-flight callback continues routing (from the worker pool), so
the exchange is never
continued while a callback might still touch it.
- When nothing is in flight — the common non-batch failure (body conversion,
`beginTransaction()`, a
single `send()` throw) — the count reaches zero and it completes in place,
exactly as before.
For the counter to reach zero, a synchronous `send()` failure must not leave
a phantom count. `doSend`
now undoes its `increment()` if `send()` throws (no Kafka callback fires in
that case), via a new
`KafkaProducerCallBack.decrement()`.
Because a mid-batch failure now waits for the already-dispatched sends to
settle before the exchange
fails, a transactional producer also aborts cleanly — the in-flight records
are accounted for before
the unit of work rolls back.
## Tests
`KafkaProducerTest.processAsyncMidBatchDispatchFailureDefersRoutingUntilInflightSendsComplete`:
a
two-record batch where the first `send()` is accepted (its callback stays in
flight) and the second
throws; asserts the exception is recorded, `sync == false`, and `done()` is
**not** called until the
first record's callback completes — then exactly `done(false)`.
Revert-to-red verified (restoring the
old immediate `done(true)` fails the deferral assertion).
The existing async / exception / transaction producer tests (single-`send()`
throw,
`beginTransaction()` failure) still pass unchanged: they hit the "nothing in
flight" branch, which
completes in place as before.
Related: CAMEL-24779 (single-message path), CAMEL-24780 (transaction begin)
— both merged.
🤖 Generated with [Claude Code](https://claude.com/claude-code)
--
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]