hutiefang76 commented on issue #12509:
URL: https://github.com/apache/seatunnel/issues/12509#issuecomment-5863304332
Thanks for the detailed checklist. The fix and regression are already in
#12510. I also ran the requested simpler-sink comparison today: the same real
Master + two Worker JVM setup, with the existing Console Sink and
checkpointable bounded sequence fixture, without any DuckDB connection,
DuckLake extension or S3 service.
Both runs killed Worker `59602` after checkpoint 1 (saved offsets `[209,
10]`) and started replacement Worker `59603`. The original engine reproduced
`WrongTargetSlotException: Unknown slot` and did not finish within the
90-second observation window, despite the replacement joining. With the current
#12510 fix (`f135f7ceb`), the job restored three times, executed on `59603`,
and finished. Console logs contain all 400 distinct IDs, 0–399 (401 output
lines, including one replay).
The harness builds a LogicalDag directly rather than submitting a REST job
file. Its configuration is:
```text
Master: CLUSTER coordinator-only
Workers: 2 initially, 2 fixed slots each, dynamicSlot=false, SLOT_RATIO
Source: bounded CheckpointableSequenceSource fixture
start_offset=0, end_offset=400, split_num=2
records_per_poll=1, emit_interval_ms=100, parallelism=2
Sink: Console, log.print.data=true, parallelism=2
Job: BATCH, checkpoint enabled, env checkpoint.interval=2000ms
Engine: checkpoint timeout=15000ms, local temporary checkpoint storage
Retry settings: unchanged defaults, 3 retries / 3-second interval
Fault: SIGKILL a SourceTask Worker after completed checkpoint 1,
then start a new Worker at a different address
```
Here is the original Master allocation/release sequence:
```text
[1790569084713] 2026-09-28 12:18:10,673 INFO [o.a.s.e.s.m.JobMaster
] [seatunnel-coordinator-service-7] - release the pipeline Job
console-zeta-recovery-1790569084713 (1790569084713), Pipeline: [(1/1)] resource
[1790569084713] 2026-09-28 12:18:10,680 WARN [o.a.s.e.s.m.JobMaster
] [seatunnel-coordinator-service-7] - Pre resource application failed for job:
1790569084713, success: 0, failed: 1/3, error: can't apply resource request:
ResourceProfile{cpu=CPU{core=0}, heapMemory=Memory{bytes=0}}
[] 2026-09-28 12:18:10,680 WARN [s.e.s.r.ResourceRequestHandler]
[SeaTunnel-CompletableFuture-Thread-6] - Apply resource not success for job:
1790569084713, required: 1 slots, applied: 0 slots, releasing slots: [],
remaining: 1 slots not assigned
[] 2026-09-28 12:18:10,680 WARN [s.e.s.r.ResourceRequestHandler]
[SeaTunnel-CompletableFuture-Thread-5] - request slot with retry error:
org.apache.seatunnel.engine.server.resourcemanager.NoEnoughResourceException:
can't apply resource request: ResourceProfile{cpu=CPU{core=0},
heapMemory=Memory{bytes=0}}
[1790569084713] 2026-09-28 12:18:10,680 WARN [o.a.s.e.s.m.JobMaster
] [seatunnel-coordinator-service-7] - Filtering failed resource for job
1790569084713 during release: can't apply resource request:
ResourceProfile{cpu=CPU{core=0}, heapMemory=Memory{bytes=0}}
[] 2026-09-28 12:18:10,680 WARN [s.e.s.r.ResourceRequestHandler]
[SeaTunnel-CompletableFuture-Thread-5] - Apply resource not success for job:
1790569084713, required: 1 slots, obtained: 0 slots, first unassigned resource
at index 0: ResourceProfile{cpu=CPU{core=0}, heapMemory=Memory{bytes=0}}
[1790569084713] 2026-09-28 12:18:10,683 INFO [o.a.s.e.s.d.p.SubPlan
] [seatunnel-coordinator-service-7] - Job console-zeta-recovery-1790569084713
(1790569084713), Pipeline: [(1/1)] state process is start
[1790569084713] 2026-09-28 12:18:10,683 INFO [o.a.s.e.s.d.p.SubPlan
] [seatunnel-coordinator-service-7] - Job console-zeta-recovery-1790569084713
(1790569084713), Pipeline: [(1/1)] turned from state CREATED to SCHEDULED.
```
And the receiving Worker's stack:
```text
org.apache.seatunnel.engine.server.service.slot.WrongTargetSlotException:
Unknown slot in slot service, slot profile:
SlotProfile{worker=[127.0.0.1]:59601, slotID=1, ownerJobID=1790569084713,
assigned=true, resourceProfile=ResourceProfile{cpu=CPU{core=0},
heapMemory=Memory{bytes=402653184}},
sequence='cf8439ac-2276-4146-9454-902658f72fa4'}
at
org.apache.seatunnel.engine.server.service.slot.DefaultSlotService.getSlotContext(DefaultSlotService.java:198)
~[classes/:?]
at
org.apache.seatunnel.engine.server.task.operation.DeployTaskOperation.runInternal(DeployTaskOperation.java:53)
~[classes/:?]
at
org.apache.seatunnel.engine.server.task.operation.TracingOperation.run(TracingOperation.java:42)
~[classes/:?]
at
com.hazelcast.spi.impl.operationservice.Operation.call(Operation.java:189)
[seatunnel-shade-hazelcast-5.1-3.0.0.jar:5.1-3.0.0]
at
com.hazelcast.spi.impl.operationservice.impl.OperationRunnerImpl.call(OperationRunnerImpl.java:273)
[seatunnel-shade-hazelcast-5.1-3.0.0.jar:5.1-3.0.0]
at
com.hazelcast.spi.impl.operationservice.impl.OperationRunnerImpl.run(OperationRunnerImpl.java:248)
[seatunnel-shade-hazelcast-5.1-3.0.0.jar:5.1-3.0.0]
at
com.hazelcast.spi.impl.operationservice.impl.OperationRunnerImpl.run(OperationRunnerImpl.java:471)
[seatunnel-shade-hazelcast-5.1-3.0.0.jar:5.1-3.0.0]
at
com.hazelcast.spi.impl.operationexecutor.impl.OperationThread.process(OperationThread.java:197)
[seatunnel-shade-hazelcast-5.1-3.0.0.jar:5.1-3.0.0]
at
com.hazelcast.spi.impl.operationexecutor.impl.OperationThread.process(OperationThread.java:137)
[seatunnel-shade-hazelcast-5.1-3.0.0.jar:5.1-3.0.0]
at
com.hazelcast.spi.impl.operationexecutor.impl.OperationThread.executeRun(OperationThread.java:123)
[seatunnel-shade-hazelcast-5.1-3.0.0.jar:5.1-3.0.0]
at
com.hazelcast.internal.util.executor.HazelcastManagedThread.run(HazelcastManagedThread.java:102)
[seatunnel-shade-hazelcast-5.1-3.0.0.jar:5.1-3.0.0]
```
On failed allocation, `preApplyResources` releases the partial allocation
without replacing the previous futures. The patch therefore enters `FAILING`
locally and returns before `restorePipeline()` can deploy from those stale
futures. A later successful allocation installs fresh futures. Handling the
failure inside the state machine also covers the direct Master-takeover entry;
the public-entry regression reproduces the earlier escaping exception and
passes with the current fix. The Console run verifies Worker loss/replacement,
not a live Master-loss scenario.
--
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]