davidzollo opened a new pull request, #12105:
URL: https://github.com/apache/seatunnel/pull/12105
### Purpose
`CheckpointCoordinator#restoreTaskState` has a modulo-based per-task state
remap specifically for handling **parallelism changes** across a
checkpoint/savepoint restore
(`seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/checkpoint/CheckpointCoordinator.java`):
```java
for (int i = tuple.f1(); i < actionState.getParallelism(); i +=
currentParallelism) {
ActionSubtaskState subtaskState = actionState.getSubtaskStates().get(i);
...
}
```
This branch is only meaningfully exercised when `currentParallelism !=
actionState.getParallelism()` (the restored parallelism differs from the
parallelism the checkpoint was taken at). Every existing restore IT in this
package (`CheckpointRestoreWithStopIT`, `SavepointRestoreIT`) restores at the
**same** parallelism the checkpoint was taken at, so this remap path has never
been exercised end-to-end by any test.
### What this PR adds
Two new E2E tests under
`seatunnel-e2e/seatunnel-engine-e2e/connector-seatunnel-e2e-base`:
- `SavepointRestoreScaleUpIT` — runs a job at parallelism 2, takes a
savepoint, stops it, and restores it at parallelism 4 (scale-up).
- `SavepointRestoreScaleDownIT` — runs a job at parallelism 4, takes a
savepoint, stops it, and restores it at parallelism 2 (scale-down).
Both sub-scenarios (scale-up **and** scale-down) are covered, since the
remap formula is asymmetric between the two directions (scale-down merges
multiple old subtask indices' state onto each new instance; scale-up leaves new
instances whose index is `>=` the old parallelism with no inherited state at
all).
Each test:
1. Submits the job via the normal `--config <conf> --set-job-id <id>` path,
waits for a completed checkpoint, then takes a savepoint (`-s`) and confirms
the source job exits cleanly.
2. Restores the **same job id** with a **different** conf file (`--config
<restoreConf> -r <id>`) whose only difference is `env.parallelism`. This is the
same CLI restore path (`ClientCommandArgs.getRestoreJobId()` →
`SeaTunnelClient#restoreExecutionContext` → `ClientJobExecutionEnvironment` →
`MultipleTableJobConfigParser`) already used by `SavepointRestoreIT`; it always
parses a fresh `LogicalDag`/parallelism from the config passed on that
invocation, restore or not, so this is a real, already-supported mechanism —
not an assumed one.
3. Verifies exact offset reconciliation (no loss, no duplication) across the
savepoint boundary, matching the rigor of `SavepointRestoreIT`
(`assertRestoreContinuesAfterBoundary` / `assertNoOffsetDuplicates`,
byte-for-byte copied assertions).
4. Additionally verifies the **restored job's actual physical task count**
via the `/trace/task-mapping` REST endpoint
(`RestConstant.REST_URL_TRACE_TASK_MAPPING`, served by
`RestHttpGetCommandProcessor` → `TraceTaskMappingService` →
`TaskMappingBuilder.build()`, which reads the live
`JobMaster#getPhysicalPlan()` on the active master) — a genuine white-box check
of the deployed plan, not just trusting that the configured parallelism was
applied. The expected task-count delta (`2 * (new - old)` parallelism) was
derived by tracing `ExecutionPlanGenerator`/`PhysicalPlanGenerator`: with no
transform stage, source and sink each contribute one task per parallel
instance, and coordinator-type vertices (split enumerator, aggregated
committer) are singletons that cancel out of the delta regardless of their
exact count.
### Files changed
-
`seatunnel-e2e/seatunnel-engine-e2e/connector-seatunnel-e2e-base/src/test/java/org/apache/seatunnel/engine/e2e/SavepointRestoreScaleUpIT.java`
(new)
-
`seatunnel-e2e/seatunnel-engine-e2e/connector-seatunnel-e2e-base/src/test/java/org/apache/seatunnel/engine/e2e/SavepointRestoreScaleDownIT.java`
(new)
-
`seatunnel-e2e/seatunnel-engine-e2e/connector-seatunnel-e2e-base/src/test/resources/savepoint-restore-rescale/stream_p2_to_localfile_scaleup.conf`
(new)
-
`seatunnel-e2e/seatunnel-engine-e2e/connector-seatunnel-e2e-base/src/test/resources/savepoint-restore-rescale/stream_p4_to_localfile_scaleup.conf`
(new)
-
`seatunnel-e2e/seatunnel-engine-e2e/connector-seatunnel-e2e-base/src/test/resources/savepoint-restore-rescale/stream_p4_to_localfile_scaledown.conf`
(new)
-
`seatunnel-e2e/seatunnel-engine-e2e/connector-seatunnel-e2e-base/src/test/resources/savepoint-restore-rescale/stream_p2_to_localfile_scaledown.conf`
(new)
No `src/main/**` changes — this is test-only.
### Test plan
- `./mvnw spotless:apply -pl
seatunnel-e2e/seatunnel-engine-e2e/connector-seatunnel-e2e-base -am -nsu` —
passed, no formatting changes needed for the new files.
- `./mvnw install -pl
'seatunnel-e2e/seatunnel-engine-e2e/connector-seatunnel-e2e-base,!seatunnel-engine/seatunnel-engine-ui'
-am -nsu -Dmaven.gitcommitid.skip=true -DskipTests -Dspotless.check.skip=true
-T 3C` — passed; the full ~64-module reactor built and installed successfully
(fail-fast reactor reached and installed the last module,
`connector-seatunnel-e2e-base`), and both new test classes were confirmed
compiled on disk:
`target/test-classes/org/apache/seatunnel/engine/e2e/SavepointRestoreScaleUpIT.class`
and `SavepointRestoreScaleDownIT.class` (plus their anonymous `$1` inner
classes from the `TypeRef` usage).
- Docker/Testcontainers-based execution of the new ITs themselves was not
run locally for this PR (consistent with this repository's standard practice of
validating E2E/Testcontainers runs through CI); CI is the authoritative signal
for the actual container-based run.
--
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]