DanielLeens commented on PR #11512:
URL: https://github.com/apache/seatunnel/pull/11512#issuecomment-5617146109
@SEZ9 Here's the grounding you asked for, checked directly against the
current head (`b921264bb`, unchanged since my last review):
**F4 — coordinator collection path, with the actual code:**
`CoordinatorService.collectCdcEnumeratorProgress()`
(`seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/CoordinatorService.java:2190-2248`)
is the implementation, and its own Javadoc (lines 2182-2188) states the
contract directly: "Enumerator tasks follow normal coordinator-task placement
and may execute on any member. Their current locations are derived from the
coordinator-owned slot map... worker-local registration is not part of the
ownership contract."
Walking the method itself:
- Lines 2195-2233: it iterates every running `JobMaster`'s physical plan,
looks up each pipeline's `ownedSlotProfilesIMap` (the coordinator-owned slot
assignment map), and for every coordinator vertex whose task group contains a
`SourceSplitEnumeratorTask`, resolves that task group's current location from
the slot map — this is the "derives enumerator task group locations from
running job plans and coordinator-owned slot assignments" step.
- Lines 2235-2247: it then calls the overloaded
`collectCdcEnumeratorProgress(taskGroupsByWorker, collect, update, onFailure)`,
where `collect` sends a
`CollectCdcEnumeratorProgressOperation(taskGroupLocations)` to each resolved
worker address via `NodeEngineUtil.sendOperationToMemberNode` — this is the
"requests reports from the assigned members" step — and completions feed into
`seaTunnelServer.getCdcProgressService().updateReports(reports)`.
- The static overload at lines 2255-2274 fires all worker requests up front
before attaching completion callbacks (matching the Javadoc note on that
method: callbacks may run on Hazelcast completion threads and must stay
non-blocking).
This is a genuine request/poll model driven entirely by coordinator-side
state (the slot map), with no worker-initiated registration step anywhere in
the method. I also diffed this against `docs/en/developer/cdc-progress.md`'s
"Runtime collection" section ("The active coordinator derives enumerator task
group locations from running job plans and coordinator-owned slot assignments.
It requests reports from the assigned members...") — the doc text and the code
match phrase-for-phrase, not just in spirit. F4 is grounded and resolved.
**F8 — explicit conclusion, since you asked me to state it:**
Yes, the deep-immutability guarantee holds at every layer, verified against
the actual (unchanged) production code, not just the new test:
- `CdcProgressPosition`'s constructor does a real defensive copy —
`Collections.unmodifiableMap(new LinkedHashMap<>(values))` — so mutating the
caller's original map after construction cannot reach the stored copy, and the
returned view itself rejects writes.
- `CdcSnapshotSplitProgress` is `public final class` with four `final`
fields and no setters — no mutation path exists post-construction regardless of
what's passed in.
- `CdcEnumeratorProgressReport`'s active-splits list is wrapped the same way
(`Collections.unmodifiableList`), and the new
`testEnumeratorReportActiveSplitWatermarksAreDeeplyImmutable` now exercises the
full nested chain (`CdcEnumeratorProgressReport` → `CdcSnapshotSplitProgress` →
`CdcProgressValue<CdcProgressPosition>` → the map), not just the outer list. I
confirmed the test is discriminating (it would fail if the defensive copy were
removed) rather than vacuously true.
F8 is resolved with both source-level confirmation and direct regression
coverage.
**F7 pointer** (since @goutamadwant's ask was specifically for a method
name): `CdcProgressModelTest.testEnumeratorReportRejectsInvalidCounts`
(`seatunnel-api/src/test/java/org/apache/seatunnel/api/cdc/CdcProgressModelTest.java:155-180`)
— asserts `IllegalArgumentException` for a negative snapshot-total count and
for a completed-count that exceeds the total, both via the
`CdcEnumeratorProgressReport` constructor directly.
**Explicit record on all of F1-F8 against `b921264bb`:** all eight remain
resolved from my side. Nothing in this head's one-file test diff touches F1-F6
(those depend on production code that hasn't changed since my 2026-08-31
line-by-line pass), and F7/F8 are now grounded above. I don't have an open item
on this list.
--
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]