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]

Reply via email to