dybyte commented on code in PR #11569:
URL: https://github.com/apache/seatunnel/pull/11569#discussion_r3870383845
##########
seatunnel-e2e/seatunnel-connector-v2-e2e/connector-jdbc-e2e/connector-jdbc-e2e-part-1/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/xa/XaGroupOpsImplIT.java:
##########
Review Comment:
Would it be worth adding a non-blocking E2E test covering partial XA commit
followed by Zeta restart and recovery? The unit tests cover the reconciliation
logic well, but this would validate the full checkpoint → failure → restore
path.
##########
seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/sink/JdbcSinkAggregatedCommitter.java:
##########
@@ -70,22 +80,76 @@ private void tryOpen() throws IOException {
@Override
public List<JdbcAggregatedCommitInfo> commit(
List<JdbcAggregatedCommitInfo> aggregatedCommitInfos) throws
IOException {
+ return commitPreparedTransactions(aggregatedCommitInfos);
+ }
+
+ /**
+ * Reconciles checkpoint XIDs with the resource manager using commit-order
evidence. Checkpoint
+ * XIDs from the first still-prepared transaction onward must all be
present in the recovery
+ * scan and are replayed strictly. An all-absent batch is treated as
already resolved, while an
+ * absent prefix before a still-prepared suffix is treated as already
resolved only after that
+ * suffix commits successfully.
+ */
+ @Override
+ public List<JdbcAggregatedCommitInfo> restoreCommit(
+ List<JdbcAggregatedCommitInfo> aggregatedCommitInfos) throws
IOException {
+ tryOpen();
+ for (JdbcAggregatedCommitInfo aggregatedCommitInfo :
aggregatedCommitInfos) {
+ // Refresh RM evidence for every batch because transactions may be
resolved concurrently
+ // during failover while earlier restored batches are being
replayed.
+ replayRecoveredCheckpoint(
+ aggregatedCommitInfo.getXidInfoList(),
recoverCheckpointTransactions());
+ }
Review Comment:
Could you clarify which actor can concurrently resolve these checkpoint XIDs
during failover? From the Zeta lifecycle, I see one aggregated committer task
per sink, pipeline restore happens after the old task group reaches a terminal
state, and master failover explicitly skips redeployment when the same
TaskGroupLocation is still active. If the concurrent resolver is external to
Zeta, a fresh recovery scan per batch only reduces stale evidence; it does not
prevent a transaction from being resolved between recover() and commit(). Is
per-batch recovery actually required for correctness here, or is it only a
defensive freshness measure?
--
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]