davidzollo opened a new pull request, #12032:
URL: https://github.com/apache/seatunnel/pull/12032
## Summary
- Adds `JobExecutionIT#testFailedPipelineMetricsAreCleanedUp`, an E2E
regression test proving that a pipeline which ends in `FAILED` has its entry
removed from the coordinator's metrics `IMap`.
- Guards against the regression fixed by #10757 ("[Fix][Zeta] Clean failed
pipeline metrics without full-map cleanup scan"), a follow-up to #10418
("[Fix][Zeta] Fix memory leak caused by failed pipeline cleanup after unstable
worker communication").
## Background
Two related fixes live in `CoordinatorService#processPendingPipelineCleanup`
/ `PipelineCleanupRecord`:
- #10418 made pipeline cleanup (TaskGroupContext + metrics IMap entries)
retry through a persisted `PipelineCleanupRecord` in a new
`IMAP_PENDING_PIPELINE_CLEANUP`, instead of a fire-and-forget best-effort call
that silently dropped cleanup forever if the worker RPC was flaky exactly when
a pipeline ended.
- #10757 discovered that even without any RPC flakiness, a `FAILED`
pipeline's metrics were **never** cleaned at all:
`CoordinatorService#shouldCleanup` and `JobMaster#removeMetricsContext` only
matched `CANCELED`/`FINISHED`, missing `FAILED` entirely.
I verified in current `dev` HEAD that both fixes are present
(`shouldCleanup(PipelineCleanupRecord)` matches `FAILED`/`CANCELED`/`FINISHED`,
and `JobMaster#removeMetricsContext` / `#enqueuePipelineCleanupIfNeeded` do the
same), and found no existing E2E coverage asserting on
pipeline-cleanup/metrics-IMap state specifically for a `FAILED` pipeline — the
fault-tolerance suite covers CANCELED/FINISHED job lifecycle, but nothing
inspects the metrics IMap after a FAILED pipeline specifically.
## What the test does
`testFailedPipelineMetricsAreCleanedUp` (added next to `testGetErrorInfo`,
which already uses the same failing config):
1. Submits `batch_fakesource_to_console_error.conf` (a SQL transform that
casts a random string to `int`) — the same deterministic
`NumberFormatException` fixture `testGetErrorInfo` already relies on to reach
`JobStatus.FAILED` reliably, with no restart-budget exhaustion or multi-node
orchestration needed.
2. Immediately after submission, seeds one metrics entry for the job's
(single) pipeline directly through `SeaTunnelServer#updateMetrics` — the same
entry point a worker's periodic metrics reporter (`ReportMetricsOperation`,
default 10s interval) itself calls. This sidesteps a real race: because this
pipeline fails almost immediately after task deployment, relying on the
worker's real ~10s report cycle to populate the coordinator's metrics IMap
before the pipeline ends (and its task group context is torn down) is not
reliable, and would risk making "no metrics after FAILED" trivially/vacuously
true. The seeded entry is indistinguishable to the coordinator's cleanup logic
from a worker-reported one — cleanup matches purely on `PipelineLocation` and
the pipeline's real terminal status, not on how the entry was populated.
3. Waits for the real job to reach `JobStatus.FAILED`.
4. Asserts, with a 120s bound (comfortably above
`CoordinatorService#PIPELINE_CLEANUP_INTERVAL_SECONDS` = 60s), that the metrics
entry for that pipeline is eventually removed from the coordinator's metrics
IMap
(`SeaTunnelServer#getEngineContext().getStateStores().metricsSnapshotStore().containsPipeline(...)`),
via whichever of the two cleanup paths fires first: the immediate best-effort
call in `SubPlan#subPlanDone`, or the 60s-scheduled
`CoordinatorService#cleanupPendingPipelines` safety net from #10418.
## Sub-scenario coverage
This fix area splits into two invariants; I covered the first and
deliberately did not attempt the second in this E2E suite:
1. **FAILED-pipeline metrics cleanup (covered)** — the straightforward,
reliably-constructible half, proven above via a real pipeline failure and a
real coordinator metrics-IMap assertion.
2. **Cleanup survives flaky worker RPC (not attempted here)** — I looked for
a clean, honest way to inject real RPC flakiness at exactly the cleanup moment
in this black-box, multi-node-in-process E2E harness, and concluded it isn't
reliably constructible here without either modifying production code to add a
fault-injection hook (out of scope for a test-only PR) or entangling the test
with unrelated node-failure/restore machinery (killing the only worker before
cleanup would also derail the pipeline's own path to a clean terminal status,
no longer isolating the invariant being tested). This half is already
thoroughly covered at the unit level in `CoordinatorServicePipelineCleanupTest`
(added/extended by #10418 and #10757), including
`testCleanupUpdatesRecordAndKeepsItWhenTaskGroupCannotBeCleaned` (an
unreachable worker address leaves the record un-cleaned and retryable, with
`attemptCount` incrementing) and
`testRestoreInvalidationPreventsCapturedCleanupFromDeletingNewRoundMetr
ics` (a real-concurrency test proving a locked-and-captured cleanup call from
a stale round cannot corrupt a newer round). I'm relying on that existing
coverage rather than shipping a weaker or flakier E2E approximation of the same
invariant.
## Test plan
- `./mvnw spotless:apply -pl
seatunnel-e2e/seatunnel-engine-e2e/connector-seatunnel-e2e-base -am -nsu` —
BUILD SUCCESS.
- `./mvnw install -pl
seatunnel-e2e/seatunnel-engine-e2e/connector-seatunnel-e2e-base -nsu
-Dmaven.gitcommitid.skip=true -DskipTests -Dspotless.check.skip=true -o`
(targeted build against the already-fully-built local reactor for this same
worktree's unmodified upstream modules, to work around severe multi-agent
CPU/memory contention on the build host during this session) — BUILD SUCCESS;
directly confirmed `JobExecutionIT.class` (including the new test method)
exists under `target/test-classes`.
- No production code (`src/main/**`) touched — test-only change.
🤖 Generated with [Claude Code](https://claude.com/claude-code)
--
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]