DanielLeens commented on code in PR #11569:
URL: https://github.com/apache/seatunnel/pull/11569#discussion_r3870884724
##########
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:
Thanks for tracing the Zeta lifecycle so carefully — that reasoning is
correct, and this is not defending against another Zeta task.
Within a single running `JdbcSinkAggregatedCommitter`, there is no
legitimate SeaTunnel-internal actor that could race this scan: restore only
replays after the old task group reaches a terminal state, and master failover
skips redeploying a `TaskGroupLocation` that is still active, so two live
instances of the same committer are never mutating the same XIDs at once. The
comment is guarding against something external to Zeta: the resource-manager
side of XA resolution. A prepared branch can be resolved independently of the
job — a DBA issuing a manual `XA COMMIT`/`XA ROLLBACK` against an in-doubt
transaction, or RM-side in-doubt-transaction cleanup tooling that many
operators run against long-pending prepared branches (sometimes triggered by
exactly this kind of stalled committer retry loop).
That matters concretely here because `aggregatedCommitInfos` passed into
`restoreCommit()` is not always a single batch.
`SinkAggregatedCommitterTask.restoreState()` flattens every
`ActionSubtaskState`'s stored bytes across all not-yet-committed checkpoints
(`seatunnel-engine/seatunnel-engine-server/.../SinkAggregatedCommitterTask.java:260-275`),
so if several checkpoints completed their prepare phase without their commit
notification landing before the failure, restore can walk more than one
`JdbcAggregatedCommitInfo` in a single call. `commitXidInfos()`
(`JdbcSinkAggregatedCommitter.java:129`) can spend up to `maxCommitAttempts`
rounds x 1s backoff per batch, so with N batches the wall-clock gap between
"scan taken" and "commit attempted for the last batch" grows with every earlier
batch's retries if only one upfront scan is used.
To be precise about what the per-batch refresh actually buys: it narrows
that staleness window, it does not close it. There is still an irreducible gap
between `recover()` returning and the following `xaFacade.commit()` call for
the very same batch — two separate RM round-trips can't be made atomic without
RM support this connector doesn't rely on. So this is a defensive freshness
measure, not a race-free guarantee, and I'd rather say that plainly than
oversell it.
What actually protects correctness if that residual gap is hit: `XAER_NOTA`
is not in `TRANSIENT_ERR_CODES` (`XaFacadeImplAutoLoad.java:75-76`), so if a
checkpoint XID the fresh scan reported as still-prepared gets resolved
externally before the commit call reaches it, `xaFacade.commit()` throws a
permanent `JdbcConnectorException`, which surfaces through `result.failed()` ->
`throwIfAnyFailed("commit")` (`XaGroupOpsImpl.java:66-72`) and fails the
restore loudly rather than silently treating it as resolved. So the residual
window can only ever manifest as a rare, operator-visible hard failure, never
as silent data loss.
I'll tighten the Javadoc on `restoreCommit()` (currently just says
"concurrently") to say it means resolution by an external actor on the resource
manager, since as written it does read like it could mean another Zeta task,
which was a completely fair thing to question.
##########
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:
That's a reasonable ask, and I don't want to wave it off — but I'd treat it
as a valuable follow-up rather than a blocker for this PR.
What's already exercising the real failure/restore path today:
`XaGroupOpsImplIT` in this same file runs against a real MySQL testcontainer
(not a mock), and `testCommitFailurePropagatesThroughAggregatedCommitter` (line
126) drives an actual XA prepare -> forced commit failure -> propagation
through `JdbcSinkAggregatedCommitter`. On top of that,
`JdbcSinkAggregatedCommitterTest` has unit coverage for every reconciliation
branch this PR introduces: recovered-by-value matching, unrelated recovered
XIDs, already-resolved prefixes, no-evidence missing gaps failing closed, and
bounded transient-retry behavior for both commit and the recovery scan itself.
What a full "partial XA commit -> kill Zeta -> restart -> verify recovery"
E2E would add on top of that is genuine end-to-end confidence that the wiring
between `SinkAggregatedCommitterTask.restoreState()` and this committer is
correct under an actual cluster restart, which the IT/UT layer can't fully
substitute for. The reason I didn't reach for it in this PR is that this exact
class of test (kill a running Zeta cluster mid-checkpoint, actually restart it,
and assert final DB state) is exactly the shape of E2E the project has been
trying to get away from — controlling when the kill lands relative to the XA
prepare/commit boundary without a `Thread.sleep`-based timing guess is
nontrivial, and a flaky version of this test would be worse for the project
than not having it at all.
I'm not opposed to adding it as a separate, deliberately-scoped follow-up PR
with proper deterministic synchronization (for example gating the kill on an
observable committer-side signal instead of a timing guess) — I'd rather get
that design right in its own review than fold it into this one under time
pressure. Happy to file the follow-up issue if that's useful.
--
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]