DanielLeens commented on PR #11569:
URL: https://github.com/apache/seatunnel/pull/11569#issuecomment-5508239093
Self-review note: this is my own PR, so GitHub blocks a formal review
submission on it — posting as a plain comment, consistent with every prior
round on this thread.
One new commit landed after my last review (`e59b7276c92` -> `360a1c81aae`,
`[Fix][Connector-V2] Close RocketMQ consumers before interrupt`). No one has
posted a new comment since my last round, so there is no pending thread to fold
in here. I re-verified from scratch rather than assuming it was safe to skip:
the diff touches only `RocketMqConsumerThread.java` (+8 lines, a new `close()`
method that calls `consumer.shutdown()`) and `RocketMqSourceReader.java` (+2
lines, `close()` now sets `running = false` and calls
`consumerThreads.values().forEach(RocketMqConsumerThread::close)` before
`executorService.shutdownNow()`). This is not part of this PR's JDBC XA scope,
and unlike every prior rider on this branch, it is the first one that changes
non-test, non-XA production source rather than only test/CI-stabilization code
— so I read it in full rather than waving it through as another test-only
rider. The JDBC XA core files (`XaFacade.java`, `XaFacadeImplAutoLoad.java`,
`XaGroupOps.
java`, `XaGroupOpsImpl.java`, `JdbcSinkAggregatedCommitter.java`) are
byte-for-byte unchanged since my last review, so this round's re-verification
of the core fix is a re-confirmation of the same source, not new logic.
# What Problem Does This PR Solve?
- User pain point: before this PR, a permanent JDBC XA commit failure could
be silently swallowed. `GroupXaOperationResult.throwIfAnyFailed("commit")` was
commented out in `XaGroupOpsImpl.commit()`, and
`XaFacadeImplAutoLoad.wrapException()` always pre-wrapped
`TransientXaException` inside a `JdbcConnectorException` before throwing it,
which made every `catch (XaFacade.TransientXaException e)` block dead code (the
concrete runtime type on the stack was always `JdbcConnectorException`, never
`TransientXaException`). A prepared XA transaction could fail to commit and be
neither retried, reported, nor rolled back, while the checkpoint was still
reported complete — real data loss for a two-phase-commit sink.
- Fix approach: restore failure propagation, correct transient-vs-permanent
XA error classification, add bounded synchronous retry with backoff, and add a
`restoreCommit()` path that reconciles restored checkpoint XIDs against a live
`xaFacade.recover()` scan in commit order before replaying still-in-doubt
transactions.
- One-sentence summary: the JDBC XA fix itself is unchanged and remains
correct on this round; what's new is one more off-scope rider commit, and for
the first time it touches production (not just test) code outside this PR's
stated purpose.
# 1. Code Change Review
## 1.1 Core Logic Analysis
**Before** (`XaFacadeImplAutoLoad.wrapException`, pre-fix):
```java
private static RuntimeException wrapException(String action, Optional<Xid>
xid, Exception ex) {
if (ex instanceof XAException) {
XAException xa = (XAException) ex;
if (TRANSIENT_ERR_CODES.contains(xa.errorCode)) {
throw new JdbcConnectorException(
JdbcConnectorErrorCode.XA_OPERATION_FAILED, new
TransientXaException(xa));
} else {
throw new JdbcConnectorException(...);
}
} else {
throw new JdbcConnectorException(...);
}
}
```
`TRANSIENT_ERR_CODES` was `{XA_RBTRANSIENT, XAER_RMFAIL}`, and
`XaGroupOpsImpl.commit()` had `result.throwIfAnyFailed("commit")` commented out
with a `// TODO ... exception is not thrown` note.
**After** (current head, unchanged since my last review):
```java
private static final Set<Integer> TRANSIENT_ERR_CODES =
new HashSet<>(Arrays.asList(XA_RETRY, XAER_RMFAIL));
...
private static RuntimeException wrapException(String action, Optional<Xid>
xid, Exception ex) {
if (ex instanceof XAException) {
XAException xa = (XAException) ex;
if (TRANSIENT_ERR_CODES.contains(xa.errorCode)) {
return new TransientXaException(xa);
} else {
return new JdbcConnectorException(...);
}
} else {
return new JdbcConnectorException(...);
}
}
```
and in `XaGroupOpsImpl.commit()`:
```java
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);
```
**Key findings, re-verified against the current head:**
- **The `throw`-inside-`wrapException` -> `return`-from-`wrapException`
change is the actual bug fix, not a refactor.** Every call site (`execute()`'s
`catch (XAException e)` branch via `.orElseThrow(() -> wrapException(...))`,
and `Command.fromRunnable`'s explicit `throw wrapException(...)`) previously
received a `JdbcConnectorException` thrown from inside `wrapException` itself,
with `TransientXaException` only reachable as `getCause()`. So `catch
(XaFacade.TransientXaException e)` in `XaGroupOpsImpl.java` and
`JdbcSinkAggregatedCommitter.java` could never match the old code. Now
`wrapException` returns the concrete exception object and the call sites throw
it directly, so `TransientXaException` reaches those `catch` blocks for real.
This one change is what makes "retry transiently" vs. "fail the checkpoint" a
real distinction again.
- **`XA_RBTRANSIENT` -> `XA_RETRY` is a correct XA-semantics fix, not a
cosmetic swap.** `XA_RBTRANSIENT` means the branch has already been rolled back
by the RM (terminal for that xid — retrying commit would just get `XAER_NOTA`).
`XA_RETRY` means the call had no effect and may be reissued — the actual
definition of retryable. Classifying `XA_RBTRANSIENT` as retryable was a latent
correctness bug independent of the swallowed-exception bug; this PR fixes both.
- **Ordering invariant holds as coded.** `XaGroupOpsImpl.commit()`'s loop
guard is `i.hasNext() && (result.hasNoFailures() || allowOutOfOrderCommits)`;
both production call sites (`JdbcSinkAggregatedCommitter.commitXidInfos` and
`replayRecoveredCheckpoint`) pass `allowOutOfOrderCommits=false`. On any
failure the loop stops advancing, and everything not yet attempted is appended
to `forRetry` via the same iterator state, preserving relative commit order
among not-yet-attempted XIDs — the invariant `replayRecoveredCheckpoint`'s
"first still-prepared index" logic depends on.
- **`replayRecoveredCheckpoint` fails closed on gaps, not just on total
absence.** `findFirstRecoveredIndex` returns the first checkpoint XID still
present in a fresh `recover()` scan; everything before that index is logged and
skipped as already-resolved, everything from that index onward is checked in a
loop where the first absent member throws before
`commitXidInfos(stillPrepared)` runs — a gap after a still-prepared transaction
cannot be silently skipped. An all-absent batch returns early without any
commit call.
- **Restore-time failure propagation reaches the master.**
`JdbcSinkAggregatedCommitter.restoreCommit()` throws directly out of the loop
on the first permanent failure. `SinkAggregatedCommitterTask.restoreState()`
has no surrounding try/catch and is declared `throws Exception`; its only
production caller, `NotifyTaskRestoreOperation.runInternal()`, runs
`task.restoreState(...)` inside a `try { ... } catch (Throwable e) {
task.getExecutionContext().sendToMaster(new
CheckpointErrorReportOperation(taskLocation, e)); }` block — a restore-time XA
failure is genuinely reported to the master, not lost inside an async executor.
**Runtime-path diagram** (checkpoint-triggered path, traced end to end):
```
Checkpoint N completes (live path)
CheckpointCoordinator.notifyCompleted()
-> SinkAggregatedCommitterTask.receivedWriterCommitInfo(...) already
buffered writer commit infos
-> JdbcSinkAggregatedCommitter.commit(aggregatedCommitInfos)
-> commitPreparedTransactions -> commitXidInfos(xidInfos)
-> XaGroupOpsImpl.commit(pending,
allowOutOfOrderCommits=false, maxCommitAttempts)
-> xaFacade.commit(xid, ignoreUnknown=false)
-> XaFacadeImplAutoLoad.execute(Command "commit")
XAException (permanent code, e.g. XAER_RMERR)
-> wrapException(...) returns
JdbcConnectorException
-> thrown out of commit(xid,...)
-> caught as `catch (Exception e)` -> result.failed(x, e)
-> loop stops (hasNoFailures()==false), remaining xids
appended to forRetry
-> result.throwIfAnyFailed("commit") ***fix: this used
to be commented out***
-> throws JdbcConnectorException
-> propagates out of commitXidInfos /
commitPreparedTransactions / commit(...)
-> propagates out of JdbcSinkAggregatedCommitter.commit(...) to the
engine's commit-notification path
-> job impact: commit notification fails instead of being reported
successful while data is
actually stuck prepared -> checkpoint is NOT silently marked complete
over lost data.
Restart / restore path (checkpoint state replay)
SinkAggregatedCommitterTask.restoreState(actionStateList) [no try/catch,
declared throws Exception]
-> JdbcSinkAggregatedCommitter.restoreCommit(aggregatedCommitInfos)
-> per batch: recoverCheckpointTransactions() [xaFacade.recover(),
bounded retry]
-> replayRecoveredCheckpoint(checkpointXids, recoveredXids)
-> findFirstRecoveredIndex(...) ; gap in suffix -> throw
JdbcConnectorException
-> else commitXidInfos(stillPrepared) [same commit path as
above, same fix applies]
-> exception propagates out of restoreState(...)
-> NotifyTaskRestoreOperation.runInternal()'s asyncExecuteFunction lambda:
try { task.restoreState(restoredState); }
catch (Throwable e) { task.getExecutionContext().sendToMaster(new
CheckpointErrorReportOperation(taskLocation, e)); }
-> job impact: master is notified of the restore failure instead of the
failure vanishing inside
the async executor; the task group does not silently continue as if
restore succeeded.
```
**On the two prior CHANGES_REQUESTED reviews:** unchanged from my last
round's re-verification — both target states this PR has since replaced
(nzw921rx's unconditional `ignoreUnknown=true` concern, and davidzollo's
mixed-outcome-restart concern, both structurally addressed by
`replayRecoveredCheckpoint`/`findFirstRecoveredIndex`, as I re-derived
line-by-line last round). Both reviews remain **formally** un-updated on GitHub
as of this comment; I cannot clear either as the PR author.
**New finding on the RocketMQ rider commit (`360a1c81aae`):**
`RocketMqSourceReader.close()` now calls `consumer.shutdown()` (via
`RocketMqConsumerThread.close()`) from the closing thread, *before*
`executorService.shutdownNow()` interrupts the consumer's own thread — whose
`run()` method already calls `this.consumer.shutdown()` in its `finally` block
on exit. That means `consumer.shutdown()` can now be invoked concurrently from
two different threads on the same `DefaultLitePullConsumer`. I checked this
against the actual RocketMQ client source (`rocketmq-client-4.9.4-sources.jar`,
not assumed): `DefaultLitePullConsumerImpl.shutdown()` is `synchronized` and
gated on `serviceState` (`RUNNING` -> does the real work and sets
`SHUTDOWN_ALREADY`; any other state, including a concurrent second call, is a
no-op `default: break`). So the double-shutdown pattern this commit introduces
is safe — it is idempotent and thread-safe by the library's own contract, not
by luck. What I can't ve
rify from a static read is whether shutting the client down from a second
thread while `run()`'s loop may still be mid-`task.accept(consumer)` (an
in-flight pull) causes that in-flight call to throw inside the consumer thread
— worth a quick look, but not a correctness risk to the shutdown call itself,
and not part of this PR's actual subject matter.
## 1.2 Compatibility Impact
Partially incompatible, and disclosed, for the JDBC XA change (unchanged
from my last review): the prefix-skip, all-absent-skip, and gap-fails-closed
cases are described accurately in
`docs/en(zh)/introduction/concepts/incompatible-changes.md`, and the doc
explicitly states that XA recovery cannot distinguish a SeaTunnel-committed XID
from one an external actor rolled back.
- No checkpoint/savepoint state-schema change: `XidInfo`'s serialized fields
(`xid`, `attempts`) are untouched.
- `TransientXaException`'s constructor widened from package-private to
`public` — additive, source-compatible.
- **Real behavior change for existing jobs**: a job that was previously
relying (even unknowingly) on a permanent XA commit failure being silently
swallowed will now see that checkpoint/restore fail loudly instead. This is the
entire point of the fix and is explicitly called out in the PR description and
the incompatible-changes doc.
- The RocketMQ rider commit is a separate, undisclosed behavior change to a
different connector (source `close()` now proactively shuts the consumer down
instead of relying solely on `executorService.shutdownNow()`'s interrupt to
trigger it through `run()`'s `finally`). It is not mentioned in this PR's
description or in the incompatible-changes doc, and arguably should not need to
be, since it does not change RocketMQ's documented user-facing contract — but
it means a RocketMQ behavior change is riding through review under a JDBC XA PR
title, which is exactly the hygiene problem Issue 3 (below) already flagged and
which just got materially worse.
## 1.3 Performance / Side-Effect Analysis
Unchanged from my last review for the JDBC XA path — restore does one
bounded `xaFacade.recover()` scan per restored batch, backoff is a bounded
`Thread.sleep` on the commit/restore control path only (not the per-record hot
path), and `XidInfo.attempts` is not persisted across restarts (a disclosed,
correct trade-off given Zeta doesn't checkpoint partial-commit progress today).
For the RocketMQ rider: `close()` now shuts the consumer client down
proactively rather than only through the interrupt-driven `finally` block,
which should make in-flight polls release faster on close/stop — a
resource-cleanup improvement, not a regression, and I found no
leaked-connection or double-close problem in it (see 1.1).
## 1.4 Error Handling and Logging
No swallowed exceptions on any path I traced, including the
restore-to-master propagation path in 1.1. Log levels are appropriate: WARN for
skip/already-resolved conditions and transient-retry attempts, thrown
exceptions (not just logs) for permanent failures. No credentials, connection
strings, or other sensitive data appear in any touched log statement or
exception message.
### Issue 1 (carried, not new): `XaFacade.commit(xid, ignoreUnknown=true)`
has no production caller
- **Location**: `XaGroupOpsImpl.java:61` (call site, hardcoded `false`);
`XaFacadeImplAutoLoad.java` (the `ignoreUnknown` parameter and
`buildCommitErrorDesc`)
- **Problem**: The 4-arg `commit(..., ignoreUnknown)` capability is
exercised only by unit tests; both production call sites pass `false`, so
`XAER_NOTA` always surfaces as a permanent failure in production rather than
being tolerated.
- **Potential risk**: Low. If an external actor resolves a checkpoint XID
between `recover()` and the following `commit()` call for the same batch, the
live commit or restore call fails hard (correctly), the job restarts, and the
next restore's fresh recovery scan finds that XID absent and skips it — an
extra, self-healing restart cycle, not silent data loss.
- **Best improvement**: Either document the unused `true` path as an
intentionally-retained escape hatch, or remove it so the API surface matches
what production actually reaches.
- **Severity**: Low
- **Raised by another reviewer**: Yes — li3zhi4 (2026-08-21), discussed
further with dybyte (2026-08-27); still open.
### Issue 2 (carried, not new): PR description text vs. implemented behavior
for the all-absent-batch case
- **Location**: PR description only (not code)
- **Problem**: The PR description characterizes the all-absent case as
"fails closed," while the code and docs correctly treat it as already-resolved
and skip it without any commit attempt.
- **Best improvement**: Edit the PR description to match the implemented and
documented behavior.
- **Severity**: Low
- **Raised by another reviewer**: Yes — li3zhi4 (2026-08-21); still open.
### Issue 3 (carried and escalating — severity raised this round): this
branch keeps accumulating unrelated commits outside its stated JDBC XA scope,
and this round crossed from test-only to production code
- **Location**: `AbstractAzureCosmosDBIT.java`, `PostgresCDCIT.java`,
`DorisErrorIT.java`, `JdbcHanaIT.java`, `MilvusIT.java`, `CouchbaseIT.java`,
`RocketMqIT.java`, `CoordinatorServiceTest.java`,
`seatunnel-engine-client/src/test/resources/hazelcast{,-client}.yaml` (all
test/CI infra, prior rounds) plus, new this round,
`RocketMqConsumerThread.java` and `RocketMqSourceReader.java` — production
source in an unrelated connector.
- **Problem**: An eighth unrelated commit landed since my last review:
`360a1c81aae` adds a `close()` method to `RocketMqConsumerThread` and calls it
from `RocketMqSourceReader.close()`. Every prior rider on this branch was
confined to test code or CI/test-infra config, which I treated as a hygiene
issue rather than a risk to certify individually. This one is different in
kind: it is production behavior in a connector this PR's title, description,
and incompatible-changes doc say nothing about. I verified it in full (see 1.1)
and did not find a correctness bug in it, but "the reviewer had to
independently verify a stranger's production fix in an unrelated connector to
be able to sign off on a JDBC XA PR" is exactly the cost this issue has been
warning about, now realized rather than hypothetical.
- **Potential risk**: Low correctness risk in this specific commit (verified
safe against the actual RocketMQ client's shutdown() contract), but a real and
now-materially-worse review/maintainability/audit cost: a `git blame` or `git
revert` on the JDBC XA fix would also have to reason about an unrelated
RocketMQ source-close change bundled into the same commit history, and any
future targeted revert of just the XA fix is more entangled than it was last
round.
- **Best improvement**: Split the RocketMQ fix and the remaining
CI-stabilization commits into their own PR(s) before merge, or at minimum have
the PR description explicitly acknowledge and justify the expanded scope,
especially now that it includes production code.
- **Severity**: Medium (raised from Low-Medium last round — process/hygiene,
not a correctness blocker in what I could verify, but the kind of change riding
along has materially changed).
- **Raised by another reviewer**: No.
No new correctness issue was found in the JDBC XA core this round; those
files are unchanged since my last review.
# 2. Code Quality Assessment
## 2.1 Coding Standards
Re-confirmed: `restoreCommit`, `replayRecoveredCheckpoint`,
`recoverCheckpointTransactions`, `findFirstRecoveredIndex`,
`containsEquivalentXid`, `backoffBeforeRetry`, `wrapException`, and `XidKey`
all carry Javadoc explaining intent. ASF license headers present on all touched
files. No wildcard imports. The new `RocketMqConsumerThread.close()` also
carries a Javadoc explaining its ordering rationale relative to interrupt.
## 2.2 Test Coverage and Test Stability: Stable
Unchanged from my last review for the JDBC XA path:
`JdbcSinkAggregatedCommitterTest` (9 tests), `XaGroupOpsImplTest` (3),
`XaFacadeImplAutoLoadTest` (5) cover the classification and reconciliation
logic with deterministic, mocked `XAException`s. `XaGroupOpsImplIT` (real MySQL
testcontainer) stops the container to force a genuine commit failure and
asserts propagation — real evidence against an actual resource manager, not
just mocks.
For the new RocketMQ commit: no unit test was added for
`RocketMqConsumerThread.close()` itself (e.g., asserting `consumer.shutdown()`
is called, or that a concurrent `run()`-thread shutdown and an external
`close()` don't conflict). The only related test change is the earlier
`waitForTopicRoute()` stabilization in `RocketMqIT.java` (already reviewed last
round, addresses topic-route flakiness, not this close-ordering fix). This is
consistent with the rest of this connector's existing test coverage style, so I
am not calling it a blocker, but it means the fix's own correctness rests on my
static read of the RocketMQ client source rather than on a test asserting the
ordering it changes.
## 2.3 Documentation Updates
`docs/en(zh)/connectors/sink/Jdbc.md` and
`docs/en(zh)/introduction/concepts/incompatible-changes.md` remain accurate
against the current JDBC XA code. No documentation was added or needed for the
RocketMQ commit, since it does not change RocketMQ's user-facing contract.
# 3. Architectural Soundness
## 3.1 Elegance of the Solution
Unchanged: a precise, minimal fix scoped to the XA classes themselves,
correcting a specific error-code misclassification, restoring
previously-disabled failure propagation, and adding order-preserving,
evidence-based restore reconciliation rather than inventing a new
persisted-state mechanism.
## 3.2 Maintainability
Good for the XA fix itself. Worse for the branch as a whole this round: see
Issue 3. The `XidKey`, `commit`/`restoreCommit` split, and retry/backoff
helpers remain small, single-purpose, and well documented.
## 3.3 Extensibility
Unchanged: the reconciliation logic operates on canonical XID values through
the existing `XaFacade`/`XaGroupOps` abstractions, extending to other
XA-capable JDBC dialects without connector-specific branching.
## 3.4 Historical-Version Compatibility
Unchanged: no checkpoint/savepoint state-schema change; `XidInfo`'s
serialized shape is untouched, so a job with in-flight XA-prepared transactions
from before this upgrade restores through the same deserialization path as one
prepared after it.
# 4. Issue Summary
| Number | Issue | Location | Severity |
| --- | --- | --- | --- |
| 1 | `XaFacade.commit(xid, ignoreUnknown=true)` has no production caller |
`XaGroupOpsImpl.java:61`, `XaFacadeImplAutoLoad.java` | Low |
| 2 | PR description overstates the all-absent case as "fails closed" | PR
description text | Low |
| 3 | Branch keeps accumulating unrelated commits outside its stated JDBC XA
scope; this round's addition is production code (RocketMQ consumer
close-ordering), not just test/CI infra | `RocketMqConsumerThread.java`,
`RocketMqSourceReader.java`, plus prior test/CI files across 8 commits | Medium
|
No High or Medium correctness issues found in the JDBC XA core on the
current head.
# 5. Merge Recommendation
### Conclusion: Ready to merge after fixes
The JDBC XA correctness fix is unchanged and remains sound on this round;
nothing in the new commit touches it. On that basis:
- **nzw921rx's 2026-07-27 CHANGES_REQUESTED** and **davidzollo's 2026-08-06
CHANGES_REQUESTED** target states this PR has since structurally replaced, as
re-derived in my prior rounds. Both remain **formally** un-updated on GitHub as
of this comment (`reviewDecision` still reflects nzw921rx's stale review;
davidzollo's latest state is `COMMENTED`, not a cleared `CHANGES_REQUESTED`).
Neither reviewer has revisited the current head. I cannot clear either as the
PR author.
- **CI on the current head (`360a1c81aae`)** is a fresh run (fork run
`33607981524`, queued/in-progress as of this comment) and has not finished. I
am not certifying it green from a prior head's results — the last head's run
(`e59b7276c92`) had completed all JDBC-specific buckets successfully with one
unrelated, out-of-scope Hazelcast test-infra flake (`edge-agent-it`), but that
result does not carry over to a new commit and must be re-confirmed on this
exact head once it finishes.
1. Blockers:
- Process, not code: the two stale review-decision states from nzw921rx
and davidzollo need a maintainer with write access (or the reviewers
themselves) to revisit the current head and either re-affirm a concrete
objection or clear the gate.
- CI on the current head needs to finish and be confirmed clean,
including the JDBC buckets, before this can be certified.
2. Recommended fixes (non-blocking for the correctness of the XA fix itself):
- Issue 3: split the RocketMQ fix and the remaining CI-stabilization
commits out of this branch, or at minimum have the PR description explicitly
justify the expanded scope now that it includes production code in an unrelated
connector. I'm raising this to Medium this round because the kind of rider
changed, not just its count.
- Issues 1 and 2: low-severity carryovers from li3zhi4, cheap to fold in
whenever this PR is next touched.
Overall assessment: the JDBC XA fix — error classification, the ordering
invariant, the restore reconciliation, and the async catch-and-report-to-master
path — is correct, and I found nothing new to correct in it this round. What
changed is procedural: CI restarted on a new head and needs to finish, and the
branch's scope-creep issue is now materially worse (production code, not just
test infra) and should be resolved by splitting the branch before merge rather
than accumulating further.
--
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]