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]