li3zhi4 commented on PR #11885:
URL: https://github.com/apache/seatunnel/pull/11885#issuecomment-5364453320
Thanks for the thorough round-2 review, @DanielLeens — both High blockers
confirmed addressed. Two follow-ups on your remaining items:
**Issue 2 (first-split staleness) — confirmation, with the timing
guarantee:**
The first incremental split cannot observe a stale stop offset because of a
hard ordering guarantee in the enumerator, independent of `resolvedStopOffset`:
1. `HybridSplitAssigner.getNext()` only enters incremental allocation after
`snapshotSplitAssigner.isCompleted()` is true
(`HybridSplitAssigner.java:99-117`) — i.e. after the snapshot assigner has
turned `assignerCompleted` (all snapshot splits executed, parallelism=1
short-circuit at `SnapshotSplitAssigner.java:206-210`).
2. `IncrementalSplitAssigner.getNext()` (first call, `splitAssigned=false`)
executes `createIncrementalSplits()` → `createIncrementalSplit()` →
`StopConfig.getStopOffset(offsetFactory)` → `offsetFactory.latest()`. Since
step 1 guarantees this happens strictly *after* the snapshot phase finished
consuming the table, the resolved offset already covers every change written
during the snapshot window — it is not stale.
3. `completedSnapshotPhase()` itself is gated by
`checkArgument(splitAssigned && noMoreSplits())`
(`IncrementalSplitAssigner.java:324`), so it necessarily runs *after* the split
was created; the offset it resolves (`resolvedStopOffset`) is therefore
at-or-after the split-creation offset, and `createIncrementalSplit()` prefers
it whenever non-null.
So the roles are: split-creation resolution is already correct by ordering,
and `resolvedStopOffset` adds (a) the checkpoint-stable value reused on restore
(no drift — the actual fix for Issue 1 in round 1), and (b) a consistent
post-snapshot value for any split created after `completedSnapshotPhase`. Both
point to the same post-snapshot position in practice.
**Issue 1 (readiness signal) — attempted the structural fix, found it
infeasible for `earliest`, accepting the residual risk:**
I tried replacing the `RUNNING` gate with a real "snapshot has started
reading" signal: wait until the first bulk row (id=1, 'bulk') is visible in the
sink, with `snapshot.split.size = 20` to keep the 200-row table multi-split.
This works for `initial`, but **fails for `earliest`**: `earliest` takes no
snapshot phase (the incremental split's `completedSnapshotSplitInfos=[]` and
its startup offset is the earliest binlog position), so the sink never sees the
bulk rows until the binlog phase replays them from the beginning of the binlog
— which does not happen within the 60s readiness window. Log evidence from the
run: `IncrementalSplit(tableIds=[...], startupOffset={... mysql-bin.000001,
pos=4 ...}, stopOffset={...}, completedSnapshotSplitInfos=[], ...)` for the
earliest job, vs 10 `CompletedSnapshotSplitInfo` entries for the initial job.
Since the structural signal cannot be made uniform across all five startup
modes without significant test-harness surgery (e.g. waiting on binlog replay
progress), I reverted to the wide-margin approach (200-row bulk insert +
immediate post-`RUNNING` UPDATE) that CI demonstrated green, and I **explicitly
accept the residual (much smaller) flake risk** on the `initial`/`earliest`
readiness timing — same trade-off you offered as acceptable.
Current head is unchanged from the round-2 review point (`c58fc10c3`):
enumerator-side single resolution + checkpointed stop offset + unit test +
CI-green E2E. Happy to take any further feedback.
--
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]