davidzollo opened a new pull request, #12031:
URL: https://github.com/apache/seatunnel/pull/12031
## What this guards
Guards the fix in #10836 ("[Fix][Zeta] Job stuck permanently after master
failover, unable to complete (affects BATCH / bounded source / job shutdown
phase)").
Before that fix, `CheckpointCoordinator#readyToCloseStartingTask` — the set
of bounded-source starting subtasks that have already finished emitting data
and are waiting for the final `COMPLETED_POINT_TYPE` checkpoint to formally
close the pipeline — lived only in the pre-failover master's JVM heap. A
subtask that had already reported "ready to close" never reports again (there
is nothing left for it to signal once it has finished). So a fresh master
starting from an empty set could never again observe
`readyToCloseStartingTask.size() == plan.getStartingSubtasks().size()`, the
completing checkpoint never fired, and an otherwise-finished BATCH job stayed
`RUNNING` forever.
The fix persists that bookkeeping into `runningJobStateIMap` (keyed by
`CheckpointCoordinator#getReadyToCloseImapKey()`) as each subtask reports in,
and restores it in `CheckpointCoordinator#restoreCoordinator` after a master
failover — including a special case for when *every* subtask had already
reported ready right before the crash, which immediately re-triggers the
completing checkpoint on restore.
## Why the existing coverage doesn't reach this window
`CheckpointCoordinatorFailoverIT#testBatchJobCompletesAfterMasterFailover`
(added by #10836 itself) already covers master failover for a BATCH job, but it
deliberately triggers the kill once `observedRows > expectedTotalRows / 4` — I
confirmed by reading it at current HEAD that this lands squarely in the middle
of active source production (25% of expected output), nowhere near the close
handshake this fix protects. It cannot exercise the `readyToCloseStartingTask`
recovery path at all.
## What the new test does
Adds `testBatchJobCompletesAfterMasterFailoverDuringCloseHandshake` to that
same class (reusing its `createTestResources` helper and template-conf
pattern), targeting the close-handshake window directly instead of guessing
from timing:
- A new conf template,
`batch_fake_to_localfile_close_handshake_failover_template.conf`, defines two
independent FakeSource operators feeding one shared LocalFile sink, so the job
compiles into exactly one pipeline / one `CheckpointCoordinator` (asserted in
the test via `PhysicalPlan#getPipelineList()`), with 4 starting subtasks total
(parallelism 2 per source):
- `table_fast`: `row.num = 5`, a single split — finishes in well under a
second.
- `table_slow`: `row.num = 300` spread across 30 splits with a 300ms read
interval between them — at least ~8.7 seconds to drain, so it is still actively
producing rows long after `table_fast` is done.
- `checkpoint.interval` is set to 600000ms (far beyond the test's real
runtime) so the completing checkpoint is the only checkpoint ever attempted,
keeping the scenario isolated to the exact mechanism under test.
- The test polls `runningJobStateIMap` directly — via
`CheckpointCoordinator#getReadyToCloseImapKey()`, the same persisted key
`restoreCoordinator` reads back from — until its value's size is strictly
between 0 and 4 (i.e. `table_fast`'s two subtasks have reported ready to close,
`table_slow`'s two have not), then immediately kills the active master.
- The cluster uses two dedicated master-only nodes
(`createMasterHazelcastInstance`) plus a separate worker node
(`createWorkerHazelcastInstance`) that is never killed — unlike this class's
other two tests, which use combined master+worker nodes. I confirmed in
`SubPlan#restorePipelineState`
(`jobMaster.getCheckpointManager().reportedPipelineRunning(pipelineId,
allTaskRunning.get())`) that the "no redeploy" recovery path is gated on every
task vertex still being observed as `RUNNING` after failover, which
`PhysicalVertex#restoreExecutionState` verifies by checking the owning worker
is still a live cluster member. Since the worker here is never killed, every
task should stay `RUNNING` through the failover, so recovery is pure
coordinator-side bookkeeping with nothing to redeploy or redo — which is what
lets the final assertion check the post-recovery row count **exactly** (`610`)
instead of with a `">="` tolerance like the existing mid-production test.
- Final assertion: the job reaches `FINISHED` within a bounded timeout with
exactly the expected row count — not stuck `RUNNING` forever (the pre-fix bug)
and not `FAILED`.
## Why this trigger construction is reliable
This reads the fix's own persisted bookkeeping directly (the same IMap entry
the fix's `restoreCoordinator` restores from) rather than inferring readiness
from a row-count threshold or timing guess. The `table_fast`/`table_slow`
production gap is a deterministic ~8+ second window by construction (driven by
`split.num` and `split.read-interval`, not scheduler jitter), so the polling
loop (20ms interval) reliably lands the kill inside the intended
partial-close-handshake window rather than racing a narrow, timing-dependent
state transition.
## Test plan
- `./mvnw spotless:apply -pl
seatunnel-e2e/seatunnel-engine-e2e/connector-seatunnel-e2e-base -am -nsu` —
passed, no additional formatting changes beyond what's in this diff.
- `./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` — `BUILD SUCCESS` across the full reactor (`seatunnel-engine-ui`
excluded since it is untouched and its `npm install` step is unrelated to this
change).
- `./mvnw install -pl
seatunnel-e2e/seatunnel-engine-e2e/connector-seatunnel-e2e-base -nsu
-Dmaven.gitcommitid.skip=true -DskipTests -Dspotless.check.skip=true` — `BUILD
SUCCESS`, with test-compile genuinely enabled this time (`Compiling 52 source
files to .../target/test-classes`), confirming the new/modified test class
compiles cleanly; verified `CheckpointCoordinatorFailoverIT.class` exists under
`target/test-classes` with a fresh timestamp.
- Per this series' local-verification policy for Apache SeaTunnel, actual
test execution (including this new test) was intentionally not run locally; it
will run under GitHub CI on this PR's head commit.
--
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]