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]

Reply via email to