fatmanverse commented on PR #11541: URL: https://github.com/apache/seatunnel/pull/11541#issuecomment-5089996470
Thanks for the detailed re-review — the blocking finding was accurate and is fixed in `fb7ab542`, along with all four follow-ups. Two of them I implemented a little differently than suggested, so let me walk through the reasoning in case you see it another way. ### Issue 1 (blocking): async failure is now scoped to its transaction Confirmed — nothing ever cleared `asyncSendException`, so a transient send error would have turned into a permanent checkpoint-failure loop. It is now reset in `beginTransaction(String)`. Of the two options you offered I went with the reset rather than `getAndSet(null)` in `checkAsyncSendException()`, because of how it combines with Issue 3. Once `send()` also inspects the same reference, a consume-once read means the first `send()` throws and clears the flag, and the following `prepareCommit()` then sees `null`, treats the transaction as healthy, and commits a transaction that already lost a record. Keeping the failure until the next transaction begins avoids that interaction. For the same reason I left `abortTransaction()` untouched: every `beginTransaction()` builds a fresh producer through `getTransactionProducer()`, so callbacks from the previous producer can't reach the new transaction's state. With `beginTransaction()` as the single reset point, the failure's lifetime is exactly one transaction, which is what lets `send()` and `prepareCommit()` share the reference safely. The recovery path you asked for is covered by `senderRecoversAfterFailedTransaction`: begin → async failure → prepareCommit throws → abort → begin → successful send → commit succeeds. ### Issue 3: fail fast on the write path Good call. `send()` now calls `checkAsyncSendException()` before handing the record to the producer, so the error surfaces on the write path instead of a checkpoint interval later. `sendFailsFastAfterAsyncSendFailureIsRecorded` also asserts that no further record reaches the producer. ### Issue 4: reflection removed from the tests Added a package-private `@VisibleForTesting` constructor taking a `TransactionProducerFactory`, using the project's shaded Guava annotation as `KafkaSourceSplitEnumerator` does. The suite now drives the real `beginTransaction()` → `send()` → `prepareCommit()` lifecycle with no reflection, so the `recordNumInTransaction` reset semantics the new guard relies on are exercised rather than injected. ### Issue 5: secondary failures are logged The callback is now a named `onSendCompleted` method that logs at WARN with the transaction id when the CAS loses, so a multi-partition incident stays diagnosable. I went with logging rather than `Throwable#addSuppressed` here: the callback runs on the producer's sender thread while the task thread may already be reading that exception in `prepareCommit()`, so mutating it would introduce a data race. ### Issue 2: docs One thing worth flagging on the paths: I couldn't find `docs/en/connector-v2/Error-Quick-Reference.md` in the tree, and there doesn't seem to be an error-code quick reference anywhere under `docs/` — the Kafka sink doc sits at `docs/en/connectors/sink/Kafka.md`. The existing KAFKA-01 … KAFKA-07 codes are undocumented too, so building that reference looks like a larger piece of work than this PR should carry. Happy to open a separate PR for it if maintainers agree it's worth having, though it would need to cover well beyond Kafka. The operational note itself is in: the Semantics section of both `docs/en/connectors/sink/Kafka.md` and `docs/zh/connectors/sink/Kafka.md` now explains that under `EXACTLY_ONCE` an async send failure fails the checkpoint — naming both codes — instead of silently dropping records. ### Verification `./mvnw test -pl seatunnel-connectors-v2/connector-kafka` → 65 tests, 0 failures (`KafkaTransactionSenderTest` 4 → 6). `spotless:apply` clean. -- 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]
