fatmanverse opened a new pull request, #11541:
URL: https://github.com/apache/seatunnel/pull/11541

   ## Purpose of this pull request
   
   Fixes #11534.
   
   Under `semantics = EXACTLY_ONCE`, the Kafka sink could intermittently lose
   the first record of a transaction during checkpointing (e.g.
   `KafkaIT#testKafkaToKafkaExactlyOnceOnStreaming` receiving only 9 of 10
   records).
   
   ### Root cause
   
   `KafkaSinkWriter` sends records via the producer's asynchronous `send()`,
   which returns immediately. Kafka only marks a transaction as *started*
   (`transactionStarted = true`) once the `AddPartitionsToTxn` request — issued
   asynchronously by the producer's sender thread — is acknowledged by the 
broker.
   
   `KafkaTransactionSender#prepareCommit()` is called by the engine before
   `snapshotState()`. If a checkpoint reached `prepareCommit()` before the first
   record's transaction registration completed, `isTxnStarted()` returned a 
stale
   `false`, which was then persisted into `KafkaCommitInfo`.
   
   On checkpoint completion, `KafkaSinkCommitter` resumed the transaction with
   `txnStarted = false`, and Kafka's `TransactionManager` skips `EndTxn` in that
   state. The pending broker-side transaction was never committed and eventually
   timed out / was aborted, leaving that record permanently invisible to
   `read_committed` consumers.
   
   ### Fix
   
   In `prepareCommit()`:
   
   1. `flush()` pending sends **before** capturing the transaction state, so the
      `AddPartitionsToTxn` registration is reflected in `isTxnStarted()`.
   2. Fail fast (`TRANSACTION_NOT_STARTED`, `KAFKA-08`) when the transaction has
      records but is still reported as not started after flushing — this means 
the
      registration never completed or a send failed asynchronously. Committing
      with `txnStarted = false` would silently drop those records, so we abort 
the
      transaction via the checkpoint instead of producing a lossy commit info.
   
   The empty-transaction path (`recordNumInTransaction == 0`) is unchanged and
   still legitimately reports `txnStarted = false`.
   
   ## Does this PR introduce _any_ user-facing change?
   
   No. No config option, public API, or SPI contract is changed. Under a slow
   broker, checkpoints may block slightly longer on `flush()`, which is the
   correct trade-off for exactly-once (a checkpoint timeout is preferable to
   silent data loss).
   
   ## How was this patch tested?
   
   - **Unit** — new `KafkaTransactionSenderTest` (3 cases):
     - flush happens before the transaction state is captured (race modeled by
       flipping `isTxnStarted()` inside the `flush()` stub, not by invocation
       counting);
     - fail-fast when records were sent but the transaction is not started after
       flushing;
     - empty transaction still allowed with `txnStarted = false`.
   - **E2E** — `KafkaIT#testKafkaToKafkaExactlyOnceOnStreaming` is the existing
     reproducer for this race; its `checkData` assertion already fails on both a
     lost record (matched < 10) and a duplicate (matched > 10). Added a
     regression-guard javadoc documenting that intent.
   - `spotless:apply` clean; module unit tests pass.
   
   ## Check list
   
   - [x] Code changed are covered with tests, or it does not need tests for 
reason in the PR description.
   - [x] If any new Jar binary package adding in your PR, please add License 
Notice according [New License 
Guide](https://github.com/apache/seatunnel/blob/dev/docs/en/contribution/new-license.md)
   - [x] If necessary, please update the documentation to describe the new 
feature. https://github.com/apache/seatunnel/tree/dev/docs
   


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