DanielLeens commented on PR #11541:
URL: https://github.com/apache/seatunnel/pull/11541#issuecomment-5279185807
@goutamadwant Thanks for the extra pass and the doc feedback — you're right,
and it's worth tightening.
I checked the actual abort semantics: `SinkWriter#abortPrepare()` is
documented as **only invoked by the Spark engine** on a `prepareCommit()`
failure (`seatunnel-api/.../sink/SinkWriter.java`, javadoc on
`abortPrepare()`). `KafkaSinkWriter#abortPrepare()` does call
`kafkaProducerSender.abortTransaction()`, but for Zeta/Flink that path isn't
triggered when `KAFKA-08`/`KAFKA-09` is thrown — the checkpoint just fails, and
the in-flight transaction is only cleaned up later, when the writer restarts
from the last successful checkpoint and `KafkaSinkWriter`'s constructor calls
`abortTransaction(checkpointId + 1)` on the stale transaction (or the broker
times it out via `transaction.timeout.ms` in the meantime).
So "both errors abort the current transaction" is accurate for Spark but
overstates it for Zeta/Flink. Your proposed wording ("fail the checkpoint and
prevent the transaction from being committed, then records are replayed from
the last completed checkpoint after recovery") matches the code for all three
engines and is the better description.
This is doc-precision only, no behavior change, so I agree it doesn't block
merge — just flagging so the wording can be nudged if there's another doc pass.
--
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]