DanielLeens commented on PR #11077:
URL: https://github.com/apache/seatunnel/pull/11077#issuecomment-5617102061
@SEZ9 Sorry about the rendering issue on your side — I've re-fetched my
2026-09-09T11:57:10Z review via the API and the body is intact and complete
(20k+ chars, review id 5153898629), so nothing was actually cut off in what I
posted; it looks like a client-side truncation when viewing it. Pasting the
part right after "substantiv" here so you don't have to fight the renderer:
> ...(which is real and substantive, but I wanted my own independent
confirmation given how many rounds this has been through).
>
> **Destination-key collision safety** (`MultiTableSink.java`,
`DestinationKey` inner class): the dedup key folds in `sink.getClass()`
whenever a connector supplies a `physicalDestinationIdentifier`, and falls back
to raw object identity (`sink == that.sink`) when either side has none...
[continues into the CI investigation and issue summary]
On the F1-F8 list: I went back through the current head (`153355ce8c`)
source line by line against each item just now, since your numbering doesn't
map 1:1 onto what I carried forward, and most of these are already resolved on
this head with evidence, not just "should be fine":
- **F1 (destination-key collision, HIGH):** Resolved.
`DestinationKey.equals()`/`hashCode()` fold in `sink.getClass()` plus the
connector-supplied `physicalDestinationIdentifier`, and fall back to raw object
identity when either side has none (`MultiTableSink.java`, `DestinationKey`
inner class). I re-derived this myself in the 09-09 review rather than trusting
prior rounds, and there's a named regression test,
`testSamePhysicalIdentifierDoesNotShareAcrossConnectorClasses` in
`MultiTableSinkWriterTest.java`.
- **F2 (restore-time state duplication, HIGH):** Resolved. `groupByIdentity`
(an `IdentityHashMap<SinkWriter, List<SinkIdentifier>>`) drives each distinct
writer instance exactly once and records the result under a single canonical
identifier (`aliasedIdentifiers.get(0)`), not fanned out to every alias —
that's what prevents the restore-time union from loading N duplicate state
copies. Covered by the shared-writer snapshot/restore round-trip test in
`MultiTableSinkWriterTest.java`, which asserts exactly one restored state list.
- **F3 (schema/config divergence unguarded, MEDIUM):** Guarded for the one
connector that actually opts into sharing today.
`BaseMultipleTableFileSink#getPhysicalDestinationIdentifier()` folds the row
type into the identifier (`path + ";row-type=" +
catalogTable.getSeaTunnelRowType()`), so two aliases with divergent row
layouts, or with schema evolution enabled, never resolve to the same
`DestinationKey` and never share a writer. `MultiTableSink` itself doesn't need
a separate guard because sharing only ever happens through a connector's own
identifier contract.
- **F5 (proxyContexts only registers the first alias via containsValue):**
Not reproducible on current head. Both `createWriter` and `restoreWriter` call
`proxyContexts.put(sinkIdentifier, proxy)` unconditionally for every alias's
`SinkIdentifier`, not gated by any `containsValue` check — I checked both
methods directly (`MultiTableSink.java`, inside the per-table loop in each).
Every aliased identifier gets a context entry; there's no O(n^2) containsValue
scan in this code path.
- **F7 (IOException wrapped as unchecked RuntimeException inside
computeIfAbsent):** Not reproducible on current head. The actual
`sink.createWriter(proxy)` / `sink.restoreWriter(proxy, state)` call sits in a
plain `if (writer == null) { ... }` block inside the outer `try`, not inside a
`computeIfAbsent` lambda — only the proxy-context lookup uses
`computeIfAbsent`, and that doesn't call into the connector. The declared
`IOException` propagates normally and is caught once by the outer `catch
(IOException error) { closeCreatedWriters(...); throw error; }`.
- **F4/F6 (undocumented SPI method / restoreWriter contract change):**
Resolved. `docs/en/developer/sink-connector-development.md` and the `zh`
counterpart document `getPhysicalDestinationIdentifier()` and the
writer-sharing contract for connector implementers — I confirmed this is
present and accurate on the current head.
- **F8 (missing Javadoc param/return tags on getDestinationKey):** Resolved.
Current `getDestinationKey` Javadoc has `@param tablePath`, `@param
replicaIndex`, and `@return` tags.
So of your eight items, none are open on the current head as far as I can
verify from the actual source — if you're seeing something different, a
`path:line` pointer would help me recheck against the exact code you're looking
at, since we're both looking at `153355ce8c`.
What is genuinely still open and carried from my own review (non-blocking,
Low severity): `initResourceManager` sizing by alias count instead of
distinct-writer count, `closeCreatedWriters` only triggering cleanup on
`IOException` (not on unchecked exceptions) from writer creation, and the
queue-colocation thread-safety invariant not being documented as a named
invariant/comment. None of these are correctness blockers for the
shard-to-one-destination use case this PR targets.
On CI: the only real failure on this head (`transform-v2-it-part-1`,
`TestFilterRowKindIT.testFilterRowKindMultiTable`) traces to
`AssertSinkWriter`'s static JVM-wide row counters (apache/seatunnel#12116,
which explicitly cites this PR's own prior CI run as one of its reproduction
cases), not to anything in this PR's diff — the Assert connector doesn't
override `getPhysicalDestinationIdentifier()`, so the writer-sharing path this
PR adds doesn't even fire for that test. That's process, not a code blocker,
but the branch-protection Build gate still needs to go green before merge
regardless.
--
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]