DanielLeens commented on PR #11077:
URL: https://github.com/apache/seatunnel/pull/11077#issuecomment-5203576238
Thanks @SEZ9 for the thorough re-review, and thanks @hesam-oxe for
continuing to push on this.
Sequence check on my side: no new commit has landed on this PR since my last
comment on 2026-07-27 — the head is still `f98af0e46f2b`, the same commit I
reviewed on 2026-06-17. So this isn't a "new commit landed, needs a fresh full
review" situation; I re-verified SEZ9's findings against the current head by
re-reading the writer creation/lifecycle path end to end.
**+1 on Issue 1 (shared-writer lifecycle called multiple times).** This is
the same root problem I flagged on 2026-06-17, and the current head still has
it. Concretely:
- `MultiTableSinkWriter`'s constructor groups `SinkIdentifier`s into
`sinkWritersWithIndex` by `entry.getKey().getIndex() % queueSize`
(`MultiTableSinkWriter.java:211-234`). Since `index` is computed purely from
`context.getIndexOfSubtask() * replicaNum + i` in
`MultiTableSink.createWriter`/`restoreWriter` — independent of the source table
— every `SinkIdentifier` that shares a `destKey` also shares the same `index`,
so aliased identifiers always land in the same queue/thread. Good news: the
cross-thread race I originally worried about in my first review is gone.
- Bad news: within that one queue's map, the *same* writer instance is still
stored under N different `SinkIdentifier` keys (`MultiTableSink.java:153` and
`:226`). `snapshotState` (`MultiTableSinkWriter.java:472-500`),
`prepareCommit(long)` (`:536-583`), `abortPrepare()` (`:616-630`), and
`close()` (`:664-680`) all iterate
`sinkWritersWithIndex.get(i).entrySet()`/`.values()` with no writer-identity
dedup anywhere in the file — there's no `IdentityHashMap`/distinct-writer-set
helper in this class. So one shared writer gets
`snapshotState`/`prepareCommit`/`close` invoked N times per checkpoint, once
per aliased source table. Confirmed real, agree High.
**+1 on Issue 2 (restoreWriter drops state for aliased tables) — I'd treat
this as a standalone hard blocker in its own right.**
`MultiTableSink.java:191-221`: the state lookup (`states.stream().map(...
.get(sinkIdentifier))...`) runs *inside*
`destinationWriters.computeIfAbsent(destKey, ...)`, so it only ever executes
once per `destKey`, for whichever `tablePath` happens to come first in
`sinks.keySet()` iteration order. Checkpoint state for every other source table
aliased to that destination is never even read, let alone merged — it's
silently unreachable. On restore from a checkpoint written before this change
(one state entry per source table), that's real data loss/duplication on the
recovery path, not a corner case.
**+1 on Issue 3**, which I'd consider the same underlying defect as Issue 1
viewed from the commit/snapshot angle rather than a separate bug — fixing
writer-identity dedup resolves both at once.
One pushback: **Issue 4 ("possible public API change") doesn't hold up under
diff inspection.** I compared `MultiTableSinkWriter`'s public constructors and
public methods against the merge-base (`de4cf3387a75`): all five public
constructors and every public method (`write`, `snapshotState`,
`prepareCommit()`/`prepareCommit(long)`, `abortPrepare`, `close`,
`applySchemaChange`) have identical signatures before and after this PR. I
don't see a removed or changed public/protected member in this file. @SEZ9 if
you had a specific signature in mind, happy to take another look, but as
written I don't think this one should block merge.
Small correction on **Issue 9**: the current head's `MultiTableSink.java`
does end with a trailing newline (verified directly on the file), so that
specific example is stale on this head — not a real issue right now. Very minor
either way, doesn't change the overall assessment.
Separately from the logic review: this PR still shows a `dirty` merge state
against `dev` — the conflict I flagged on 2026-07-27 hasn't been resolved yet,
and the `Build` check is still failing on the current head. Both need to be
sorted out before this can be merged, independent of the logic fixes above.
@hesam-oxe — to summarize where this leaves things, the two real correctness
blockers are: (1) writer-identity deduplication across
`snapshotState`/`prepareCommit`/`abortPrepare`/`close` so a shared writer's
lifecycle is only driven once per checkpoint instead of once per aliased source
table, and (2) fixing `restoreWriter` to union checkpoint state across every
`SinkIdentifier` that maps to the same destination before restoring, not just
the first one encountered. My suggestion is to key the per-queue
dedup/iteration by distinct writer instance (identity) rather than by
`SinkIdentifier`, which should clean up both the lifecycle and the
restore-state problem together. Once that's in, please also resolve the
conflicts with `dev` and get CI green — I'm glad to do another full pass after
that.
--
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]