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]