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]

Reply via email to