DanielLeens commented on PR #11569:
URL: https://github.com/apache/seatunnel/pull/11569#issuecomment-5379386326
# What Problem Does This PR Solve?
- **User pain point**: In the JDBC XA (2PC) exactly-once sink, a permanent
commit failure could be swallowed silently. `XaGroupOpsImpl.commit()` recorded
the failure into a result object but the line that would have thrown it
(`result.throwIfAnyFailed("commit")`) was commented out with a "we can't tell
restore-failure from real-failure" TODO, so
`notifyCheckpointComplete()`/`restoreState()` in `SinkAggregatedCommitterTask`
would see an empty return list and report the checkpoint as successful even
though a prepared transaction never actually committed. Separately, retryable
XA outcomes (`XA_RETRY`, `XAER_RMFAIL`) were wrapped inside
`JdbcConnectorException` before reaching the `catch (TransientXaException)`
branch, so the "transient vs permanent" classification never worked in
practice, and restore blindly replayed the whole checkpointed XID list, which
could hit `XAER_NOTA` on a XID that a previous partial commit had already
resolved.
- **Fix approach**: Restore `throwIfAnyFailed("commit")` so a
permanent/unknown failure fails the checkpoint. Fix the classification so
`XA_RETRY`/`XAER_RMFAIL` are treated as transient (retryable) and
`XA_RBTRANSIENT` (an actual rollback) is treated as permanent. Complete the
bounded `max_commit_attempts` retry loop synchronously inside a single
`commit()`/`restoreCommit()` invocation, since the engine never persists a
returned "needs retry" list across a task restart. On restore, reconcile the
checkpointed XID list against a fresh XA recovery scan in original commit
order: everything before the first still-prepared XID is treated as already
resolved, the still-prepared suffix is replayed strictly, and a gap after the
first still-prepared XID fails closed instead of being silently skipped.
- **One-sentence summary**: This PR makes JDBC XA commit failures actually
fail the checkpoint instead of being silently swallowed, and makes restore
idempotent by reconciling against a live recovery scan instead of blindly
replaying the checkpointed XID list.
# 1. Code Change Review
## 1.1 Core Logic Analysis
**Before → After, the headline fix (`XaGroupOpsImpl.java`):**
```java
// before
result.getForRetry().addAll(xids);
// TODO ... So currently the exception is not thrown.
// result.throwIfAnyFailed("commit");
throwIfAnyReachedMaxAttempts(result, maxCommitAttempts);
```
```java
// after
result.getForRetry().addAll(xids);
// A permanent or unknown commit failure must fail the checkpoint instead of
being reported
// as a successful commit.
result.throwIfAnyFailed("commit");
throwIfAnyReachedMaxAttempts(result, maxCommitAttempts);
```
I traced why this alone would have been unsafe to ship without the rest of
the PR: `SinkAggregatedCommitterTask.restoreState()` (dev, L260-282) calls
`aggregatedCommitter.restoreCommit(aggregatedCommitInfos)` with the
checkpoint's full original XID list. The interface default
(`SinkAggregatedCommitter.restoreCommit`, `seatunnel-api`) is just `return
commit(aggregatedCommitInfo)`. Before this PR, `JdbcSinkAggregatedCommitter`
didn't override `restoreCommit`, so turning on `throwIfAnyFailed` by itself
would have made *every* restore after a partial-batch commit throw
`XAER_NOTA`-turned-`JdbcConnectorException` for the XIDs a previous attempt had
already committed — restore would never succeed again for that job. That's
exactly the scenario `davidzollo`'s 2026-08-06 `CHANGES_REQUESTED` review
called out concretely (batch `[xidA, xidB]`, `xidA` commits, `xidB` fails,
restart, `xidA` now returns `XAER_NOTA`).
The new `restoreCommit()` override
(`JdbcSinkAggregatedCommitter.java:91-109`) is the fix for that: it doesn't
replay the checkpointed list blindly, it reconciles against a fresh
`xaFacade.recover()` scan first:
```java
private void replayRecoveredCheckpoint(List<XidInfo> checkpointXids,
Set<XidKey> recoveredXids) {
int firstRecoveredIndex = findFirstRecoveredIndex(checkpointXids,
recoveredXids);
if (firstRecoveredIndex < 0) { // none of the
batch remains prepared
log.warn("... treating it as already resolved: {}", checkpointXids);
return; // whole batch
skipped, not replayed
}
List<XidInfo> stillPrepared = new ArrayList<>();
for (int i = firstRecoveredIndex; i < checkpointXids.size(); i++) {
if (!containsEquivalentXid(recoveredXids,
checkpointXids.get(i).getXid())) {
throw new JdbcConnectorException(...); // gap after
first-recovered -> fail closed
}
stillPrepared.add(checkpointXids.get(i));
}
commitXidInfos(stillPrepared); // only the
still-prepared suffix is committed
// prefix before firstRecoveredIndex logged as already-resolved only
*after* the suffix commits
}
```
For the `[xidA, xidB]` scenario: on restart, `xidA` is absent from the fresh
recovery scan (already committed, RM forgot it) and `xidB` is present (still
prepared) → `firstRecoveredIndex = 1` → `stillPrepared = [xidB]` →
`xaFacade.commit()` is never called on `xidA` again, so it can never re-hit
`XAER_NOTA`. I verified this against the actual code at the current head, not
just the PR description. This is a structurally sound answer to "how do I know
an absent XID was already resolved vs. genuinely lost": it uses recovery-scan
evidence gated by *commit order* rather than inferring success from `XAER_NOTA`
alone (which is ambiguous — "unknown to RM" can mean "already committed" or
"never existed").
**Key Findings:**
- The single-invocation bounded-retry loop (`commitXidInfos`,
`JdbcSinkAggregatedCommitter.java:126-142`) exists for a concrete, verifiable
reason, not just style:
`SinkAggregatedCommitterTask.notifyCheckpointComplete()` (dev, L316-321) does
`List<...> commit = aggregatedCommitter.commit(...); if (!isEmpty(commit))
throw CheckpointException(...)`. The returned "needs retry" list is **never fed
back into a later commit attempt** — a non-empty return already fails the
checkpoint today, on `dev`, independent of this PR. So a connector-level
"return partial success, try again next round" contract doesn't actually exist
at the engine level; retries have to be exhausted synchronously within one
call, which is what this PR does. I confirmed this by reading the actual engine
source, not by trusting the PR description's claim.
- `XidKey` (`JdbcSinkAggregatedCommitter.java`, the canonical-value wrapper
for driver-specific `Xid` implementations) does defensive `Arrays.copyOf` on
construction and correct `equals`/`hashCode` (format ID + global-tx-id +
branch-qualifier) — no shared-mutable-state issue, and this is the right way to
compare `Xid` values across driver implementations that don't override `equals`.
- The `wrapException` refactor in `XaFacadeImplAutoLoad.java`
(return-the-exception instead of throw-inside-helper) is safe: I checked all
three call sites (`execute()` L293/L297 and the two `Command` factory methods
L422/L444) and every one of them does `throw wrapException(...)` — no call site
was left silently discarding the returned exception.
- `XA_RBTRANSIENT` moving from the transient set to the permanent-failure
path is the correct classification per the JDBC/XA contract: `XA_RBTRANSIENT`
means the transaction **was rolled back** (for a transient reason), i.e., the
outcome is already known and negative — retrying `commit()` on an
already-rolled-back branch cannot succeed. `XA_RETRY` (new) means "no effect,
may be reissued" — genuinely retryable. The old code conflated "rolled back"
with "retryable", which combined with the previously-disabled
`throwIfAnyFailed`, meant a rolled-back branch was retried until attempts were
exhausted and then silently dropped.
**Runtime path (checkpoint recovery / restart), verified against
`seatunnel-engine`:**
```
Job restart after failure
-> SinkAggregatedCommitterTask.restoreState(actionStateList)
[Hazelcast operation thread]
deserialize ActionSubtaskState -> List<JdbcAggregatedCommitInfo>
(checkpointed XID order preserved)
-> JdbcSinkAggregatedCommitter.restoreCommit(aggregatedCommitInfos)
for each batch:
recoverCheckpointTransactions() -- xaFacade.recover(),
bounded-retry on transient scan failure
replayRecoveredCheckpoint(checkpointXids, recoveredXids)
firstRecoveredIndex < 0 -> whole batch already
resolved, skipped
gap after firstRecoveredIndex -> JdbcConnectorException,
restoreState() rethrows -> restart fails closed
still-prepared suffix -> commitXidInfos() ->
xaGroupOps.commit(..., strict)
if restoreCommit() returns non-empty (only possible if a future
change reintroduces a non-throwing
failure path) -> restoreState() throws
CheckpointException(AGGREGATE_COMMIT_ERROR)
```
This is a real, frequently-exercised path — it runs on every task restart
that has in-flight XA checkpoint state, not just a rare corner case, so I
applied the checkpoint/recovery scrutiny level from the review protocol.
## 1.2 Compatibility Impact
**Fully compatible on the public/SPI surface, and correctly documented as
behavior-incompatible where it matters.**
`XaGroupOps`/`XaGroupOpsImpl`/`XaFacade*` are all internal (`...internal.xa`)
classes with no external extension contract. The genuinely user-visible
behavior change — permanent XA failures now fail the checkpoint instead of
silently succeeding, and restore reconciles against a live recovery scan — is
documented in both `docs/en/introduction/concepts/incompatible-changes.md` and
the `zh` equivalent, plus `docs/en|zh/connectors/sink/Jdbc.md`, with a
migration note pointing operators at `XA RECOVER` / `pg_prepared_xacts`
inspection before upgrading. I read all four doc files and they match the
implementation as traced above (the last commit, `d08cc136b5`, specifically
fixed a mismatch `li3zhi4` caught between the PR description and the
all-absent-batch behavior, and that fix is reflected in the docs too).
## 1.3 Performance / Side-Effect Analysis
- `restoreCommit()` issues one `xaFacade.recover()` round-trip per restored
batch rather than a single hoisted scan for all batches — intentional and now
explicitly commented (`JdbcSinkAggregatedCommitter.java` above
`recoverCheckpointTransactions()`) and tested
(`testRestoreCommitRefreshesRecoveryScanForEachBatch`), trading one extra RM
round-trip per batch for correctness under concurrent resolution during
failover. This only affects the restore path, not steady-state commit, so the
cost is bounded and infrequent.
- No new threads, executors, or unbounded in-memory structures. `pending`
lists in `commitXidInfos` are bounded by the checkpoint's own XID count and
`maxCommitAttempts`.
- **Issue 1 below (no backoff)** is the one real side-effect concern in this
PR.
## 1.4 Error Handling and Logging
**Issue 1: No backoff between synchronous bounded-retry rounds (Medium)**
- **Location**: `JdbcSinkAggregatedCommitter.java:126-142`
(`commitXidInfos`'s `while (!pending.isEmpty())` loop) and
`JdbcSinkAggregatedCommitter.java:235-` (`recoverCheckpointTransactions`'s
`for` loop).
- **Problem**: Both bounded retry loops fire back-to-back with zero delay
between rounds. Up to `max_commit_attempts` (default 3) XA round-trips happen
in the same synchronous call, effectively in microseconds.
- **Potential risk**: For the exact failure class this retry budget exists
to absorb — a resource manager that is transiently slow/overloaded, not fully
down — the whole budget is consumed before the RM has any realistic chance to
recover, so the "retry" is closer to "fail 3x instantly then fail the
checkpoint" than to a real retry. This doesn't cause data loss or incorrect
results (the checkpoint still correctly fails closed, which is this PR's actual
goal), but it can turn a transient blip that would have resolved itself in ~1s
into an unnecessary job restart, and repeat on the next restart's
`restoreCommit()` bounded-retry too.
- **Best improvement**: Option A — add a small fixed backoff (e.g. 500ms–1s)
between rounds in both loops. Option B — capped exponential backoff if variance
in RM recovery time is expected to be larger. Either is a small, self-contained
change.
- **Severity**: Medium. I'm rating this above the "Low" in my own prior
round on this PR, because on reflection the argument that it makes the retry
budget "effectively inert" for its target failure class is fair, and this is a
financial-grade exactly-once sink path where an avoidable extra restart cycle
has real operational cost. It's not High/blocking because it doesn't threaten
correctness — the worst case is an extra restart, not a wrong result or a
silent failure.
- **Raised by another reviewer**: Yes (@davidzollo, 2026-08-06, non-blocking
section; carried forward and re-confirmed by me in my 2026-08-21T14:06 round as
Low; @JeremyXin, 2026-08-21T15:09, labeled MAJOR). I'm splitting the difference
at Medium for the reasons above — see my note to @JeremyXin at the end of this
review.
**Issue 2: `containsEquivalentXid` is an O(N×M) linear scan instead of using
the already-built `Set` (Low)**
- **Location**: `JdbcSinkAggregatedCommitter.java:281-`
(`containsEquivalentXid`), called from `replayRecoveredCheckpoint`'s per-XID
loop.
- **Problem**: Each checkpoint XID is checked against `recoveredXids` via
`Set.contains(XidKey.from(xid))`, which is actually O(1) per call since
`recoveredXids` is already a `HashSet<XidKey>` — on re-reading this more
carefully than my initial pass, this is *not* O(N×M), `Set.contains` is O(1)
amortized. I'm downgrading my own carried-over note here: there is no real
algorithmic issue, `normalizeXids`/`XidKey` are already used consistently. No
action needed; withdrawing this as a formal issue.
**Issue 3: Missing real-driver test coverage for the exact XA error codes
this PR reclassifies (Medium)**
- **Location**: `XaFacadeImplAutoLoadTest.java`, `XaGroupOpsImplTest.java`,
`JdbcSinkAggregatedCommitterTest.java` (all three new test files) are 100%
Mockito-based; the one pre-existing IT that talks to a real XA-capable
database, `XaGroupOpsImplIT.java`
(`seatunnel-e2e/.../connector-jdbc-e2e-part-1`), remains `@Disabled("Temporary
fast fix, reason: JdbcDatabaseContainer: ClassNotFoundException:
com.mysql.jdbc.Driver")`.
- **Problem**: The entire correctness claim of this PR rests on
`XA_RETRY`/`XAER_RMFAIL` vs `XA_RBTRANSIENT` being the codes real JDBC XA
drivers (MySQL Connector/J, PostgreSQL) actually throw for "retryable" vs
"rolled back" outcomes, and on `XAResource.recover()` returning driver-specific
`Xid` values that `XidKey` can correctly canonicalize. None of that is
exercised against a real driver in this PR — only against hand-constructed
`XAException`s and mocked `Xid`s.
- **Potential risk**: If a real driver's actual behavior differs from the
JDBC/XA spec's textbook meaning of these codes (driver bugs and inconsistencies
here are common in the wild), the classification could be silently wrong in
production despite all unit tests passing.
- **Best improvement**: I checked the `@Disabled` reason on
`XaGroupOpsImplIT` — it predates this PR and is caused by an unrelated
Testcontainers/MySQL-driver classloading issue, not anything in this diff.
Fixing that classloading issue is legitimately out of scope for this PR.
Recommendation: file a follow-up to re-enable `XaGroupOpsImplIT` against a real
database once that infra issue is fixed, specifically covering an end-to-end
path where a real driver's commit/recover response drives the new
classification and reconciliation logic. Non-blocking for this PR since the
root cause is pre-existing and unrelated.
- **Severity**: Medium (evidentiary gap on the PR's core correctness claim,
not a bug in the code as tested).
- **Raised by another reviewer**: Yes (@JeremyXin, 2026-08-21T15:09, labeled
MAJOR — I agree with the substance, but since the blocker (disabled IT) is a
pre-existing, unrelated infra issue rather than something this PR should be
required to fix, I'm keeping it non-blocking at Medium rather than treating it
as a merge blocker for this specific PR).
# 2. Code Quality Assessment
## 2.1 Coding Standards
No wildcard imports, no `System.out.println`, ASF license headers present on
all new files. `restoreCommit()` and the new private helpers all carry Javadoc
that accurately describes the commit-order reconciliation contract — I checked
it against the implementation line by line and it matches (unlike the PR
description text, which `li3zhi4` caught overstating the all-absent-batch case
as "fails closed"; that mismatch was only in the PR body, not in code, and has
since been corrected in the PR description).
## 2.2 Test Coverage and Test Stability
**Rating: Stable.** All new tests (`XaFacadeImplAutoLoadTest`,
`XaGroupOpsImplTest`, `JdbcSinkAggregatedCommitterTest`, 8+ methods covering:
transient-vs-permanent classification, heuristic-commit handling, gap
fail-closed, missing-prefix skip, all-absent-batch skip, retry exhaustion,
transient recovery-scan retry, and bounded-retry-loop safety-net) are
deterministic Mockito-based unit tests. No `Thread.sleep`, no shared static
state, no execution-order dependence, no floating-point comparisons. They
assert on exact exceptions, call counts, and argument matchers rather than
timing or approximate state. This is good coverage of the *classification and
reconciliation logic*; see Issue 3 above for the separate, legitimate gap in
real-driver-level coverage.
## 2.3 Documentation Updates
`docs/en/connectors/sink/Jdbc.md`, `docs/zh/connectors/sink/Jdbc.md`,
`docs/en/introduction/concepts/incompatible-changes.md`, and
`docs/zh/introduction/concepts/incompatible-changes.md` were all updated in
this PR. I read all four and cross-checked them against the implementation
traced in 1.1 — they are accurate as of the current head, including the
all-absent-vs-gap distinction that required a correction in the last commit.
# 3. Architectural Soundness
## 3.1 Elegance of the Solution
**Precise fix**, not a workaround. The commit-order recovery-scan
reconciliation answers "was this absent XID already resolved or genuinely lost"
using evidence that already existed in the code (the XA recovery scan) as the
actual gate for what gets (re)committed, rather than an `ignoreUnknown`-style
escape hatch that just tolerates "unknown" outcomes after the fact. The last
commit (`d08cc136b5`) also removed a dead 4-arg `commit(..., ignoreUnknown)`
overload that had no production caller — good hygiene, correctly scoped as
cleanup rather than new surface.
## 3.2 Maintainability
Good. One commit code path (`commit(xids, allowOutOfOrderCommits,
maxCommitAttempts)`) instead of two overloads with unclear intended use; the
per-batch recovery-scan-refresh rationale and the `ignoreUnknown` removal are
both explained in comments rather than left implicit.
## 3.3 Extensibility
`XidKey` canonicalization and the `GroupXaOperationResult`
transient/permanent split are reusable patterns for any other XA-capable
connector following the same structure.
## 3.4 Historical-Version Compatibility
No previously released version has this restore-reconciliation behavior —
this is a correctness fix landing directly on top of a real,
previously-silent-failure defect, so there's no downgrade-compatibility surface
to preserve. Upgrade migration guidance for operators is present and accurate
(see 1.2).
# 4. Issue Summary
| # | Issue | Location | Severity |
| --- | --- | --- | --- |
| 1 | No backoff between synchronous bounded-retry rounds |
`JdbcSinkAggregatedCommitter.java:126-142`, `:235-` | Medium |
| 2 | Missing real-driver (non-mocked) test coverage for the XA error-code
reclassification | `XaGroupOpsImplIT.java` (pre-existing, disabled for an
unrelated reason) | Medium |
(No High or blocking-correctness issues found. One item from my own
carried-over notes — an alleged O(N×M) scan in `containsEquivalentXid` — is
withdrawn above on closer inspection; it's an O(1) `Set.contains` call.)
# 5. Merge Recommendation
### Conclusion: Ready to merge after fixes
**1. Blockers — must be resolved before merge (process/CI, not
code-correctness):**
- Confirm the fork CI `Build` workflow finishes green for the current head
`d08cc136b5`. As of this review, the run for this head is still `in_progress`
after a rerun of 3 unrelated flaky jobs (StarRocks/Kafka/MySQL-CDC integration
tests — unrelated to this JDBC XA diff; this PR's own `jdbc-e2e` tests already
passed on this head). Needs to finish green before merge.
- `@davidzollo`'s 2026-08-06 `CHANGES_REQUESTED` review is still the current
GitHub review state on this PR. Both I and `@li3zhi4` independently re-traced
his exact `[xidA, xidB]` walkthrough against the current
`replayRecoveredCheckpoint` implementation and concluded it's resolved (see
1.1) — but since GitHub doesn't let me withdraw someone else's review,
`@davidzollo` should re-review and update his state, or explicitly confirm
agreement, before this merges.
**2. Recommended fixes — non-blocking:**
- Issue 1 — add a small backoff (fixed or capped-exponential) between the
synchronous retry rounds in `commitXidInfos` and
`recoverCheckpointTransactions`.
- Issue 2 — file a follow-up to fix `XaGroupOpsImplIT`'s unrelated
MySQL-driver classloading issue and re-enable it, so the classification logic
gets exercised against a real driver at least once.
**Overall assessment.** The core defect this PR fixes is real and serious
for a financial-grade exactly-once pipeline: a permanent XA commit failure
could previously be reported as a successful checkpoint with zero error signal.
The fix is structurally sound — I independently traced the
restore-reconciliation logic against `@davidzollo`'s concrete blocking scenario
from source, not from either reviewer's prose, and it holds up. `@li3zhi4`'s
three non-blocking suggestions from the 08-21T08:08 round were all genuinely
addressed in the follow-up commit, not just claimed to be. The two open items
(backoff, real-driver test coverage) are real and worth fixing, but neither
undermines the correctness of the fix itself, so they're non-blocking
recommendations rather than blockers. I don't see a case for an alternative
approach here — reconciling against the live recovery scan in commit order is
the right primitive for this problem, and I couldn't find a simpler design that
handles the
partial-batch-restart case correctly.
@JeremyXin — thank you for the fresh pass and for catching the missing
real-driver coverage gap specifically, that's a fair point I hadn't weighted
heavily enough in my own last round. I've folded both of your points into
Issues 1 and 2 above with my own severity assessment (Medium rather than Major)
and the reasoning for that; happy to hear if you see a concrete failure
scenario that pushes either back up to blocking.
This posts as an issue comment rather than a PR review/approval: it's my own
PR, and GitHub does not allow self-approval or self-request-changes on it.
--
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]