DanielLeens commented on PR #11077:
URL: https://github.com/apache/seatunnel/pull/11077#issuecomment-5633072441
Thanks @SEZ9. Going through this in order, including the ordering question,
which I traced through the actual code rather than reasoning abstractly about
it.
**F2 follow-up (canonical-identifier ordering stability) — verified safe,
not a live concern.** I checked `SinkIdentifier` first: it's
`@EqualsAndHashCode` over `tableIdentifier` + `index` only
(`SinkIdentifier.java`), not object identity, so its hash is a deterministic
function of those two business fields. But the more important point is that
ordering doesn't actually matter here, because the restore path never depends
on *which* identifier was chosen as primary at snapshot time.
`MultiTableSink.getRestoredState()` (`MultiTableSink.java:338-350`) is called
with the *full* set of aliased identifiers for a destination key, recomputed
fresh at restore time from the current `sinks.keySet()`
(`identifiersByDestinationKey`, built at `MultiTableSink.java:224-236`), and it
does `identifiers.stream().flatMap(identifier ->
states...map(...).get(identifier))...` — i.e. it probes the persisted state map
for *every* possible alias and keeps whichever one(s) have a non-null entry. So
even
if `groupByIdentity`'s iteration order over the `ConcurrentHashMap`
(`MultiTableSinkWriter.java:324-334`) picks a different "primary" identifier
after a restart than it did before, the restore lookup still finds the state,
because it isn't indexed by position — it's a full scan-and-filter over all
known aliases for that destination. A config reorder that keeps the same set of
`(tableIdentifier, index)` pairs mapped to the same destination key is
unaffected either way, since `identifiersByDestinationKey` is keyed by content,
not by loop order. No existing test specifically exercises a reordered alias
list, but I don't think one is needed given the lookup is inherently
order-independent — happy to be shown a scenario where that reasoning breaks
down.
**F3 (schema/config divergence) — agreed your option (a) is the stronger
fix.** The per-connector mitigation (`BaseMultipleTableFileSink` folding row
type into its identifier) only protects connectors that happen to implement it
carefully; it's not a guarantee the shared `MultiTableSink` layer itself
enforces. A defensive check comparing `CatalogTable` schema compatibility
across every alias resolving to the same `DestinationKey`, failing fast with a
clear error, is preferable to relying on each future connector author
remembering to encode every compatibility-relevant dimension into
`getPhysicalDestinationIdentifier()`. I'd keep this Medium/non-blocking as
before, but I'd like to see (a) rather than (b) if @davidzollo has bandwidth
for it.
**F4-F8 — reposting the per-item status from my 2026-09-10T10:24:00Z comment
(id 5617102061), since the part after F3 seems to hit the same client-side
rendering cutoff we've run into before on this thread:**
- **F4 (docs for `getPhysicalDestinationIdentifier()` SPI method):**
Resolved — documented in `docs/en/developer/sink-connector-development.md` and
the `zh` counterpart, confirmed present and accurate on `153355ce8c`.
- **F5 (`proxyContexts` only registering the first alias via
`containsValue`, O(n^2) startup):** Not reproducible on the current head. Both
`createWriter` and `restoreWriter` call `proxyContexts.put(sinkIdentifier,
proxy)` unconditionally inside the per-table loop, for every alias — no
`containsValue` gate exists in this code path anymore.
- **F6 (changed `restoreWriter` contract undocumented):** Resolved, covered
by the same doc update as F4.
- **F7 (`IOException` wrapped as unchecked `RuntimeException` inside
`computeIfAbsent`):** Not reproducible on the current head.
`sink.createWriter(proxy)`/`sink.restoreWriter(proxy, state)` sit in a plain
`if (writer == null) {...}` block inside the outer `try`, not inside a
`computeIfAbsent` lambda; the declared `IOException` propagates normally to the
single outer `catch (IOException error) { closeCreatedWriters(...); throw
error; }`.
- **F8 (missing Javadoc param/return tags on `getDestinationKey`):**
Resolved — `@param tablePath`, `@param replicaIndex`, and `@return` are present
on the current head.
If any of F5/F7 are reproducible against a different line range than what I
quoted, a `path:line` pointer would help me recheck, since we should both be
looking at `153355ce8c`.
Once F3's defensive check (or an explicit decision to keep it as a
documented per-connector responsibility) lands, I don't have anything else open
on this PR from a correctness standpoint.
--
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]