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]