DanielLeens commented on PR #11885: URL: https://github.com/apache/seatunnel/pull/11885#issuecomment-5385713634
Thanks for the very thorough round of review, @SEZ9 — this pushed me to re-verify my own round-2/round-3 conclusions against the actual source rather than trusting my earlier read, and the result is a mix: one of your findings turns out not to hold up, but another is a real gap I missed. Splitting it out: **Issue 1 / Issue 6 (checkpoint serializer) — I don't think this holds up; pushing back with evidence** Both issues are built on the premise that enumerator state is persisted through a hand-written, field-by-field `PendingSplitsStateSerializer` that this PR forgot to update. I searched the full `dev` tree (`git ls-tree -r --name-only dev` / GitHub tree API) for any file matching `*Serializer*` under `connector-cdc-base`'s `enumerator` package, and there is no such class anywhere in this module — and the PR's own changed-file list (7 files) doesn't touch one either, because none exists to touch. What actually persists this state is plain Java serialization: `PendingSplitsState extends Serializable`, and `IncrementalPhaseState` (`.../enumerator/state/IncrementalPhaseState.java:26,29`) implements it directly via `@Data` with a `private static final long serialVersionUID = -6809026812298443356L` that the PR explicitly preserves unchanged (see the comment on it: "Preserve compatibility with checkpoints written when this state had no fields."). That's the correct, standard mechanism for this exact scenario: with the `serialVersionUID` held stable, `ObjectInputStream` reading an old stream (written before the `stopOffset` field existed) will default-initialize the missing field to `null` via normal Java field-evolution rules — no custom read/write logic is needed for that to work. Given there's no serializer to bypass, `IncrementalSplitAssignerTest`'s restore test — constructing the `IncrementalPhaseState` from `snapshotState()` and passing it straight into the restore constructor — isn't skipping a persistence layer; that object graph *is* the checkpoint payload. So Issue 6's "the test can't detect data loss on a real round-trip" concern doesn't apply here either. **Issue 2 (stale offset for the pre-resolution split, restore boundary drift) — you're right, and I missed this in rounds 2/3. Conceding with thanks.** I re-traced the actual call order rather than relying on the author's explanation this time, and it confirms your read: - `IncrementalSplitAssigner.completedSnapshotPhase()` (`IncrementalSplitAssigner.java:323-324`) is gated by `checkArgument(splitAssigned && noMoreSplits())` — it cannot run until `splitAssigned` is already `true`. - `splitAssigned` only becomes `true` inside `getNext()` (`IncrementalSplitAssigner.java:139-146`), immediately after `createIncrementalSplits()` has already built the split(s) using whatever `resolvedStopOffset` held *at that moment* — which is still `null` the first time, so it falls through to the live `sourceConfig.getStopConfig().getStopOffset(offsetFactory)` (`IncrementalSplitAssigner.java:304-307`), a separate `latest()` call. - `HybridSplitAssigner.getNext()` unlocks incremental allocation purely off `snapshotSplitAssigner.isCompleted()` (`HybridSplitAssigner.java:104`), and `SnapshotSplitAssigner.isCompleted()` flips via its own internal split-tracking bookkeeping (`SnapshotSplitAssigner.java:208`, `254-265`) — independent of, and typically earlier than, the enumerator receiving the reader's explicit `CompletedSnapshotPhaseEvent`, which is the only thing that triggers `IncrementalSplitAssigner.completedSnapshotPhase()` (`IncrementalSourceEnumerator.java:124-134`). So the ordering isn't just possible, it's effectively the normal path with the default `incrementalParallelism`: `createIncrementalSplit()` resolves and hands the running reader one `latest()` value, then `completedSnapshotPhase()` resolves a second, later `latest()` value purely for the checkpoint. On restore, `createIncrementalSplit()` prefers the checkpointed `resolvedStopOffset` (`IncrementalSplitAssigner.java:304-307` again, this time non-null), so the post-restart split can stop at a materially different binlog position than the one the reader was actually running against pre-restart — which is exactly the class of restart-drift bug this PR exists to close, just triggered one step earlier than round 1's original report. I'd treat this as a blocking item alongside (really, in place of) Issue 1, not just a follow-up: e.g. resolving `resolvedStopOffset` once inside `createIncrementalSplit()` itself (guarded by `resolvedStopOffset == null`, same pattern already used in `comple tedSnapshotPhase()`) so there's a single authoritative resolution point for both the live split and the checkpoint, as your suggested fix proposes. I haven't independently re-verified Issues 3/4/5/7/8 in this pass — this reply is scoped to the two points above where the review's conclusion needed direct confirmation. @li3zhi4, given Issue 2 is confirmed real, I'd suggest prioritizing that fix; happy to take another full pass once it's addressed. -- 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]
