davidzollo opened a new pull request, #12109:
URL: https://github.com/apache/seatunnel/pull/12109

   ## What
   
   Adds 
`SplitClusterFaultToleranceIT#testManyPipelinesRestoreContentionInWorkerDown`, 
an E2E probe for a verified-fragile part of the Zeta engine's pipeline-restore 
path: **synchronous, backoff-free slot allocation on restore**.
   
   This is part of the Zeta-engine-core E2E test initiative (Tier 2: 
verified-fragile current architecture, no existing coverage). It documents 
production behavior; it does not change any production code.
   
   ## The verified architecture gap
   
   `ResourceUtils.applyResourceForPipeline()` requests every slot a pipeline 
needs, joins the futures, and if even one is missing, throws 
`NoEnoughResourceException` immediately — there is no wait/retry at this layer:
   
   
https://github.com/apache/seatunnel/blob/dev/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/dag/physical/ResourceUtils.java#L49-L67
   
   ```java
   // TODO If there is no enough resources for tasks, we need add some wait 
profile
   allocateResources(subPlan, futures, preApplyResourceFutures);
   ...
   if (futures.size() != slotProfiles.size()) {
       throw new NoEnoughResourceException();
   }
   ```
   
   `SubPlan.stateProcess()`'s `SCHEDULED` case catches *any* exception from 
that call — including this one — and calls `makePipelineFailing(e)`:
   
   
https://github.com/apache/seatunnel/blob/dev/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/dag/physical/SubPlan.java#L669-L701
   
   That consumes one unit of the pipeline's own fixed `pipelineMaxRestoreNum` 
budget (default 3, from `job.retry.times`; each SubPlan tracks this 
independently — see `canRestorePipeline()` at L341-343), with a 
`pipelineRestoreIntervalSeconds` sleep (default 3, from 
`job.retry.interval.seconds`) before each attempt (`prepareRestorePipeline()`, 
L489-515).
   
   **Net effect**: if several pipelines restore at the same time (e.g. a worker 
dies while hosting slots for more than one of them) and briefly contend for a 
shrunken slot pool, a pipeline can burn its *entire* retry budget purely on 
allocation timing and permanently fail — even if the contention would have 
resolved moments later.
   
   ## What the test does
   
   Rather than hoping placement gets unlucky, the test forces genuine 
contention by construction:
   
   - The job config 
(`cluster_batch_fake_to_localfile_slot_contention_template.conf`) declares 4 
independent FakeSource -> LocalFile chains sharing no table/transform, so the 
planner's connected-components pipeline split (`PipelineGenerator`) produces 4 
separate `SubPlan`s, each with its own independent restore budget.
   - Each simple 1-parallelism pipeline needs exactly 3 slots: 1 
source-enumerator coordinator, 1 sink-committer coordinator (LocalFile's sink 
always registers an aggregated file-commit coordinator, independent of 
`is_enable_transaction`), and 1 fused reader/writer task group. 4 pipelines x 3 
slots = 12 total slot demand.
   - Both workers are pinned to a **fixed** slot pool (`dynamic-slot=false`, 
`slotNum=8`) — confirmed as a hard ceiling in 
`ResourceRequestHandler.preCheckWorkerResource`, which only falls back to 
on-demand slot creation for `isDynamicSlot()` workers. 16 total capacity 
comfortably admits the initial 12-slot demand, but no single 8-slot worker can 
host more than 2 complete pipelines (6 of 8 slots).
   - This makes the contention capacity-forced, not lucky: whichever worker is 
killed must be carrying at least `12 - 8 = 4` of the 12 slots (since the other 
worker caps at 8), and 4 slots can only belong to 2+ distinct pipelines (each 
pipeline is at most 3 slots) — so killing either worker always strands part of 
more than one pipeline at once, forcing them to restore-race the survivor's 
remaining fixed capacity.
   - `is_enable_transaction=true` on every LocalFile sink ensures a 
canceled/aborted attempt's in-progress writes stay in an uncommitted temp 
location and never surface in the counted output directory, so the 
exact-row-count assertion is sound regardless of how many restore attempts 
occur.
   - The test then asserts the desired/correct outcome (job reaches `FINISHED` 
with exactly `testRowNumber * pipelineNum` rows). Today's fixed 9-second 
guaranteed-sleep budget (3 attempts x 3s, before any real 
cancel/checkpoint-cancel/resource-release/RPC overhead) either is enough 
patience for this contention to resolve, or it isn't — this test surfaces the 
current, honest answer at HEAD via its pass/fail result rather than asserting a 
pre-decided outcome.
   
   ## Test plan
   
   - `./mvnw spotless:apply -pl 
seatunnel-e2e/seatunnel-engine-e2e/connector-seatunnel-e2e-base -am -nsu` — 
BUILD SUCCESS, no formatting diffs left uncommitted.
   - `./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 65-module reactor; confirmed the 
compiled `SplitClusterFaultToleranceIT.class` is freshly produced under 
`target/test-classes` (this local machine's shared-resource contention made the 
reactor build itself take ~3h wall clock; that is unrelated to this change).
   - Per this repository's local-verification policy for the main SeaTunnel 
repo, local verification was limited to formatting and compilation (no local 
E2E execution, no Docker). The actual runtime outcome for 
`testManyPipelinesRestoreContentionInWorkerDown` — whether the current 
architecture tolerates this specific contention level or a pipeline spuriously 
exhausts its restore budget — is exactly what this PR's CI run will show, and 
is the empirical finding this test exists to surface.
   
   ## Notes
   
   - No production code is modified. This is a test-only change (one 
new/modified test class, one new test resource template).
   - Not a regression test for a fixed bug — it is a probe for a currently-open 
architectural gap (the `TODO` in `ResourceUtils`), consistent with prior E2E 
hardening work in this class.
   


-- 
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