DanielLeens opened a new pull request, #12195:
URL: https://github.com/apache/seatunnel/pull/12195

   ## Purpose of this pull request
   
   Fixes a state-duplication bug in `CheckpointCoordinator#restoreTaskState` 
that affects every checkpoint/savepoint restore of a pipeline whose actions run 
with parallelism > 1.
   
   ### Root cause
   
   - `TaskLocation#getTaskVertexId()` 
(`seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/execution/TaskLocation.java`)
 returns the per-subtask `taskID`, whose low-order digits are the subtask's own 
parallelism index, so the value is unique per subtask.
   - `CheckpointCoordinator#getPipelineTasks` groups `TaskLocation`s by that 
id, so every group has size 1.
   - `restoreTaskState` used that group size as the step of its modular remap 
loop (`for (i = tuple.f1(); i < actionState.getParallelism(); i += 
currentParallelism)`). With a step of 1 the loop degenerates: restored subtask 
`j` receives every checkpointed subtask state with index >= `j`.
   
   Consequences: on a rescale (for example 2 -> 4 or 4 -> 2) the same split 
state is delivered to several readers and each re-emits the same records; on a 
same-parallelism failover subtask 0 receives every subtask's state. Most 
bundled connectors mask this because their enumerators dedup splits by id in 
`addSplitsBack`; enumerators that keep a plain list/deque expose it as 
duplicate output. This is the root cause of the `Found duplicate offsets 
(exactly-once violated)` failures in #12105 (`SavepointRestoreScaleUpIT` / 
`SavepointRestoreScaleDownIT`); the CI log of that PR shows each restored 
reader registering with a single-element `subtaskStats`, i.e. `pipelineTasks` 
values of 1 at runtime.
   
   ### Fix
   
   Derive the remap step from the checkpoint plan's subtask-to-action mapping 
instead: `getActionParallelism` counts, per `ActionStateKey`, how many 
`(ActionStateKey, index)` tuples carry a real (non-`COORDINATOR_INDEX`) index, 
which is exactly the number of subtasks currently running that action. The 
modular loop, the `COORDINATOR_INDEX` handling, `pipelineTasks` and its only 
other caller (`getTaskStatistics`) are unchanged; the new map is additive and 
used only on the restore path. The checkpoint format is untouched (restore side 
only), so previously written checkpoints and savepoints restore unchanged.
   
   ### Verification
   
   - New 
`CheckpointCoordinatorTest#testRestoreTaskStateDeliversEveryCheckpointedSubtaskStateExactlyOnce`
 drives the real `restoreTaskState` path for every old/new parallelism pair in 
1..4 x 1..4 and asserts that each checkpointed subtask state reaches exactly 
one restored subtask (no duplicate, no drop), that a subtask only receives 
indexes congruent to its own modulo the new parallelism, and that the 
coordinator task receives exactly the coordinator state. It fails on `dev` 
before this change and passes with it.
   - The full `CheckpointCoordinatorTest` class passes (15 tests).
   - The rescale ITs in #12105 and the existing same-parallelism 
restore/failover ITs (`SavepointRestoreIT`, `CheckpointRestoreWithStopIT`, 
`CheckpointCoordinatorFailoverIT`) are left to GitHub CI as the verification of 
record for this PR head.
   
   ### Related, deliberately not changed
   
   `SourceSplitEnumeratorTask#addSplitsBack` synchronizes on `this` while 
`receivedReader`/`triggerBarrier` synchronize on `enumeratorContext`, and 
`requestSplit`/`handleSourceEvent` are unsynchronized, contrary to the 
`SourceSplitEnumerator` contract. That is a separate contract-hardening change 
on the steady-state hot path of every dynamic-split connector and needs its own 
review; it is not the cause of the duplication fixed here.
   
   ## Does this PR introduce any user-facing change?
   
   No. Restores with parallelism > 1 no longer deliver duplicated subtask 
state; the checkpoint format and all configuration are unchanged.
   
   ## How was this patch tested?
   
   Unit test matrix described above; engine E2E via GitHub CI on this PR head.
   
   ## Check list
   
   * [x] Code changed are covered with tests, or it does not need tests for 
reason:
   * [ ] If any new Jar binary package adding in your PR, please add License 
Notice according
     [New License 
Guide](https://github.com/apache/seatunnel/blob/dev/docs/en/contribution/new-license.md)
   * [ ] If necessary, please update the documentation to describe the new 
feature. https://github.com/apache/seatunnel/tree/dev/docs
   * [ ] If you are contributing the connector code, please check that the 
following files are updated:
     1. Update change log that in connector document. For more details you can 
refer to 
[connector-v2](https://github.com/apache/seatunnel/tree/dev/docs/en/connector-v2)
     2. Update 
[plugin-mapping.properties](https://github.com/apache/seatunnel/blob/dev/plugin-mapping.properties)
 and add new connector information in it
     3. Update the pom file of 
[seatunnel-dist](https://github.com/apache/seatunnel/blob/dev/seatunnel-dist/pom.xml)
   * [ ] Update the 
[`release-note`](https://github.com/apache/seatunnel/blob/dev/release-note.md).
   
   🤖 Generated with [Claude Code](https://claude.com/claude-code)
   


-- 
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]

Reply via email to