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]
