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]

Reply via email to