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

   ### Purpose of this pull request
   
   Tier-3 (scale/stress) E2E for the Zeta engine master's shared 
`CoordinatorService` executor — scenario **L4** from the ongoing Zeta 
engine-core extreme-case E2E initiative (Tier 1 shipped 11 PRs: #12027-#12035, 
#12095, #12098). This scenario had zero existing coverage. Test-only change; no 
`src/main` or `pom.xml` changes.
   
   `CoordinatorService#createCoordinatorExecutor()` builds a single, node-wide 
`ThreadPoolExecutor` shared across every job's state-transition callbacks, 
checkpoint I/O, and the pending-job scheduler on that node — one pool, no 
per-job or per-pipeline bound. 
`CoordinatorService#restoreAllRunningJobFromMasterNodeSwitch()` is itself 
dispatched onto this same executor when a node becomes the new active master, 
and fans every job still needing restore out onto it in one unthrottled pass. 
With N jobs alive on the old master at failover time, this submits N restore 
tasks to the shared pool in one tight loop, and because the backing queue is 
synchronous with an unbounded max pool size, this is a genuine 
unbounded-thread-creation risk under a mass-failover event. This test builds a 
small-but-real version of that event and asserts what the source actually does 
about it, without proposing or requiring a fix.
   
   **Source verification (traced against this branch's HEAD):**
   
   - `CoordinatorService.java:287-298` (`createCoordinatorExecutor()`): `new 
ThreadPoolExecutor(coreThreadNum, maxThreadNum, 60L, TimeUnit.SECONDS, new 
SynchronousQueue<>(), ..., new ThreadPoolStatus.RejectionCountingHandler())` — 
zero queue capacity, so any task beyond an already-idle pool thread spawns a 
brand-new thread immediately.
   - `ServerConfigOptions.java:456-466`: `core-thread-num` defaults to `10`, 
`max-thread-num` defaults to `Integer.MAX_VALUE`, wired through 
`CoordinatorServiceConfig.java:29-33`.
   - `CoordinatorService.java:1241-1271` (`checkNewActiveMaster()`, 
synchronized): the executor is only recreated if already 
`isShutdown()`/`isTerminated()` (line 1246-1248); otherwise the same 
pre-existing, idle `executorService` instance is reused across the promotion. 
`initCoordinatorService()` runs first (line 1249), then `isActive = true` is 
set (line 1251) — so the instance this test samples via `isCoordinatorActive()` 
is exactly the one that receives the restore burst, not a freshly swapped-in 
one.
   - `CoordinatorService.java:611-662` (`initCoordinatorService()`), line 
658-661: submits exactly one 
`CompletableFuture.runAsync(this::restoreAllRunningJobFromMasterNodeSwitch, 
executorService)` — the top-level restore task that, once running, does the 
actual fan-out.
   - `CoordinatorService.java:989-1093` 
(`restoreAllRunningJobFromMasterNodeSwitch()`): retry-wrapped IMap read 
(992-1005), a terminal-state ("zombie") pre-filter loop (1016-1042) that does 
not apply to genuinely running jobs, a worker-registration wait gate 
(1047-1055, a no-op here since the worker is already registered before 
failover), the N-way fan-out itself (1056-1083: 
`needRestoreFromMasterNodeSwitchJobs.stream().map(entry -> 
CompletableFuture.runAsync(..., executorService))`), and a blocking join on all 
N futures (1085-1092).
   - `Constant.java:40,42`: `OPERATION_RETRY_TIME = 30`, `OPERATION_RETRY_SLEEP 
= 2000` (ms) — the retry budget backing `RetryUtils.retryWithException`, used 
both by the IMap read above and by 20+ other state-transition call sites this 
same executor also carries.
   - `CoordinatorService.java:2253-2270` (`getThreadPoolStatusMetrics()`): 
reads `ThreadPoolExecutor`'s own getters directly (`getActiveCount`, 
`getCorePoolSize`, `getMaximumPoolSize`, `getPoolSize`, 
`getCompletedTaskCount`, `getTaskCount`, `getQueue().size()`) plus the 
rejection counter — O(1), safe to sample in a tight loop. Already public and 
already used in production for telemetry export (`SeaTunnelServer.java:421-422` 
-> `JobThreadPoolStatusExports.java:44`) and covered by 
`TelemetryCollectorCoordinatorGuardTest`; this PR adds no `src/main` code, it 
only calls an existing public API from a new test.
   
   **What `CoordinatorExecutorMassFailoverStormIT` does** (new file, 
`seatunnel-e2e/seatunnel-engine-e2e/connector-seatunnel-e2e-base`):
   
   - Boots a split-role, 3-node embedded cluster (2 master-only nodes + 1 
worker-only node), mirroring the pattern already established by the merged 
`SplitClusterPendingJobLifecycleFailoverIT` (#12034).
   - Submits **30 concurrent long-running streaming jobs** (reusing the 
existing `pending_jobs_streaming_lifecycle.conf` resource — a `FakeSource` at 
parallelism 2 against an `InMemory` sink with a synthetic writer sleep, so each 
job is CPU-light and never completes on its own) and waits until all 30 report 
RUNNING, so all 30 are genuinely replicated into `runningJobInfoIMap` before 
the failover below.
   - Kills the active master (the same master-kill technique already used 
throughout this test family) to trigger a genuine mass-simultaneous-restore 
condition on the promoted standby.
   - Immediately after the standby's `CoordinatorService` reports active, 
samples its executor via `getThreadPoolStatusMetrics()` in a 20ms-interval loop 
for 20 seconds, tracking peak pool size, peak active count, and peak rejection 
count.
   - Asserts what is actually observed: the pool grows past its configured 
`core-thread-num` (10), no task is ever rejected (consistent with 
`max-thread-num` defaulting to unbounded), and a generous CI-safety ceiling 
(500 threads) is never crossed — a backstop against a future pathological 
regression, not a claim that today's design bounds growth in general.
   - Cleans up: waits for the restore to fully finish, cancels all 30 jobs, and 
shuts down all nodes.
   
   **Scale rationale:** 30 concurrent streaming jobs on a single JVM-embedded, 
3-node Hazelcast cluster (no Docker, no separate processes) — "dozens," per 
this initiative's brief for scale/stress scenarios: large enough to reliably 
push the shared pool well past its default core size of 10, small enough to 
stay safe on a shared CI runner. The default test config (`seatunnel.yaml` in 
this module) sets `slot-service.dynamic-slot: true` (resource-accounting rather 
than a hard slot count), so 30 lightweight streaming jobs schedule without 
artificially blocking in PENDING — cross-referenced against 
`DefaultSlotService#selectBestMatchSlot`.
   
   ### Does this PR introduce _any_ user-facing change?
   
   No. Test-only addition; no production code, configuration, or documentation 
changes.
   
   ### How was this patch tested?
   
   This adds a new E2E test (`CoordinatorExecutorMassFailoverStormIT`); no 
existing test behavior is changed.
   
   Locally, per this fork's convention for the Apache SeaTunnel repo, only 
formatting and compile verification were run (not the IT itself):
   
   - `./mvnw spotless:apply -pl 
seatunnel-e2e/seatunnel-engine-e2e/connector-seatunnel-e2e-base -am -nsu` -> 
`BUILD SUCCESS`; spotless only reflowed a few Javadoc lines to the column 
limit, no semantic changes.
   - `./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` 
-> `BUILD SUCCESS`; confirms the new IT compiles. Verified the compiled 
`.class` file actually exists under `target/test-classes` (not just `BUILD 
SUCCESS` -- `-Dmaven.test.skip=true` would fake that by skipping test 
compilation entirely, so this build intentionally omits that flag).
   
   Actual execution of this IT (the concrete peak pool size/active 
count/rejection count observed under load) is left to this PR's GitHub CI run, 
which is the authoritative test-execution signal for this repo's contribution 
workflow.
   
   ### Check list
   
   * [ ] If any new Jar binary package adding in your PR, please add License 
Notice according [New License 
Guide](https://github.com/apache/seatunnel/blob/dev/docs/en/developer/new-license.md)
 -- N/A, no new dependency.
   * [ ] If necessary, please update the documentation to describe the new 
feature. -- N/A, test-only change, no user-facing feature.
   * [ ] If necessary, please update `incompatible-changes.md` to describe the 
incompatibility caused by this PR. -- N/A, no incompatible change.
   * [ ] If you are contributing the connector code, please check connector 
registration files are updated. -- N/A, not a connector change.
   


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