DanielLeens commented on PR #11597: URL: https://github.com/apache/seatunnel/pull/11597#issuecomment-5379365882
# What Problem Does This PR Solve? - **User pain point**: Zeta's `/overview` endpoint only reports cluster-wide slot totals; operators diagnosing a stuck or over-provisioned cluster have no way to ask "how many slots is running job X holding right now, and on which workers/pipelines?" - **Fix approach**: a read-only REST endpoint, `/running-jobs/slot-usage` (v2, Jetty) and `/hazelcast/rest/maps/running-jobs/slot-usage` (v1), backed by a stateless aggregation helper (`RunningJobSlotUsageBuilder`) that intersects the coordinator's owned-slot map with the resource manager's assigned-slot snapshot, grouped by job/pipeline/worker. - **One-sentence summary**: a purely additive, read-only per-running-job slot-usage diagnostics endpoint for the Zeta REST API. This is a genuinely fresh, from-scratch verification pass rather than a rehash of the earlier rounds on this PR. I did not trust the prior write-ups; I pulled every changed file straight from the GitHub Contents API at the exact head commit (`20538ae1ebe78cc2d35871fbd02e1d79b960e154`) and re-derived every conclusion below directly from that source, plus fresh live checks of CI and branch drift. Simple example: ```bash curl http://127.0.0.1:8080/running-jobs/slot-usage ``` ```json [ { "jobId": "733584788375093248", "slotCount": 3, "pipelineSlotCounts": { "1": 2, "2": 1 }, "workerSlotCounts": { "10.0.0.8:5801": 2, "10.0.0.9:5801": 1 } } ] ``` # 1. Code Change Review ## 1.1 Core Logic Analysis Runtime path, traced end to end against the current head: ```text HTTP GET /running-jobs/slot-usage (v2) HTTP GET /hazelcast/rest/maps/running-jobs/slot-usage (v1) -> JettyService: ServletHolder(RunningJobSlotUsageServlet) bound at REST_URL_RUNNING_JOBS_SLOT_USAGE -> RunningJobSlotUsageServlet.doGet() -> RestHttpGetCommandProcessor.handleRunningJobsSlotUsage() -> RunningJobSlotUsageService.getRunningJobSlotUsageJson() <-------------------------------------------/ | +-- this node IS master -> RunningJobSlotUsageBuilder.build(server) [in-process, same thread] +-- this node NOT master -> NodeEngineUtil.sendOperationToMasterNode(new GetRunningJobSlotUsageOperation()).join() -> master: GetRunningJobSlotUsageOperation.run() -> RunningJobSlotUsageBuilder.build(server) [inline, on whatever thread invokes run()] RunningJobSlotUsageBuilder.build(server): IMAP_RUNNING_JOB_INFO.keySet() .forEach(jobId -> if (coordinatorService.shouldShowAsRunningJob(jobId)) runningJobIds.add(jobId)) <- calls getCoordinatorService() per id assignedSlots = resourceManager == null ? [] : resourceManager.getAssignedSlots({}) build(IMAP_OWNED_SLOT_PROFILES, runningJobIds, assignedSlots) -> intersect owned-slot entries with assignedSlots by (worker, slotId, ownerJobId, sequence) -> group into per-job SlotUsage(slotCount, pipelineSlotCounts, workerSlotCounts) -> toResponse() ``` I independently re-verified the full registration matrix (this PR touches connector-adjacent REST wiring, so I checked it like a plugin registration even though it isn't one): - `JettyService.java`: servlet holder registered and mapped to `REST_URL_RUNNING_JOBS_SLOT_USAGE` — present. - `RestConstant.java`: new constant `/running-jobs/slot-usage` — present, no collision with `/running-jobs` or `/running-jobs/summary`. - `RestHttpGetCommandProcessor.java` (v1 Hazelcast path): field, constructor wiring (both constructors), dispatch case, and handler method — all present and route to the same `RunningJobSlotUsageService`. - `ClientToServerOperationDataSerializerHook.java`: new class id `17` — I compared this file at the PR head against the current `dev` head (48 commits ahead) and confirmed `dev`'s highest existing id is still `16`. **Id 17 is still collision-free even after the branch's drift**, so this specific rebase concern is a non-issue in practice, though a rebase is still good hygiene (see Issue 8). Key findings (all independently confirmed by reading the actual files, not by trusting a prior summary): - **Issue 1 confirmed real**: `resourceManager == null ? Collections.emptyList() : resourceManager.getAssignedSlots(...)` (`RunningJobSlotUsageBuilder.java:68-73`) means that when the resource manager is not yet initialized (a real state during master failover/cold-start — `getInitializedResourceManager()` returns `null` in that window), `assignedSlots` is empty, so *every* slot fails the `isAssignedToJob` filter and every running job reports `slotCount: 0`. That is byte-for-byte identical to what a genuinely idle job would report. The API gives no signal to distinguish "the source I need isn't ready yet" from "this job truly holds zero slots" — which is exactly the failover scenario an operator would reach for this endpoint to diagnose. - **Issue 2 confirmed real, and I independently re-derived the exact number**: `server.getCoordinatorService()` is called once per running job id inside the `forEach` filter lambda (`RunningJobSlotUsageBuilder.java:58-66`). I opened `SeaTunnelServer.java` myself and read `getCoordinatorService()` (`SeaTunnelServer.java:286-315`): on a master node where the coordinator isn't active yet, it loops `maxRetry = 3` times doing `Thread.sleep(retryPause)` with `retryPause = 500`, i.e. **up to 1.5 seconds of blocking sleep per running job**, not just once per request. With N running jobs during a coordinator-warmup window, that's up to `N × 1.5s` of blocking — see Issue 3 for where that blocking actually happens. - **Issue 3 confirmed real, with a precise sibling comparison**: `GetRunningJobSlotUsageOperation.run()` (`GetRunningJobSlotUsageOperation.java:48-52`) calls `RunningJobSlotUsageBuilder.build(service)` directly and synchronously — there is no executor hand-off. I pulled the directly comparable sibling, `GetRunningJobMetricsOperation.java`, and its `run()` wraps the equivalent work in `CompletableFuture.supplyAsync(() -> ..., getNodeEngine().getExecutionService().getExecutor("get_running_job_metrics_operation"))` — i.e. it deliberately moves the heavy work off whatever thread invoked `run()` and onto a dedicated named executor. `GetRunningJobSlotUsageOperation` does not follow that established pattern in this codebase, and it is the one operation in this pair that can now block for up to `N × 1.5s` per Issue 2. Since `Operation.run()` for a non-partition-specific operation executes on Hazelcast's shared generic-operation thread pool by default, this can tie up a thread shared with other cluster-internal operations during exactly the failover window this endpoint is meant to help diagnose. - **Issue 6 confirmed real** (and independently corroborated by `SEZ9`): `RunningJobSlotUsageBuilderTest.java` (108 lines, one `@Test`) only exercises the package-private 3-argument `build(ownedSlotProfiles, runningJobIds, assignedSlots)` overload with a real numeric slot-reuse case. It never calls the public `build(SeaTunnelServer)` entry point, never exercises the `resourceManager == null` degraded path, never touches the operation/serializer, and there is no REST-layer test (no `RunningJobSlotUsageRestIT` following the existing `PendingJobsRestIT`/`RealtimeMetricsRestIT` precedent in this module). - **Issue 7 plausible and worth fixing, kept Low**: `slotCount` is incremented once per `TaskGroupLocation -> SlotProfile` entry (`RunningJobSlotUsageBuilder.java:114-118`, `addSlot` at `189-195`), not deduplicated by `SlotKey`. If SeaTunnel's slot-sharing model ever colocates more than one task group in the same physical slot for a job (multiple `TaskGroupLocation`s mapping to slot profiles that compare equal under `SlotKey`), this would double-count that slot. I did not find a counter-example in the diff that rules this out, and the field name strongly implies "distinct slots," so this remains a real, if narrow, naming/semantics mismatch. **Correction to the location citations in this PR's own review history.** While fetching the actual files to verify each issue, I found that several previously-reported secondary line ranges do not exist in the real file: `RunningJobSlotUsageBuilder.java` is only 206 lines total (`wc -l` and the GitHub Contents API agree), yet earlier rounds cited `446-448` (Issue 1), `423-426` (Issue 4), and `503-509` (Issue 7) in addition to the correct primary ranges. Likewise `GetRunningJobSlotUsageOperation.java` is only 58 lines, not the `92-96` cited for Issue 3. The primary location for each issue (the first range cited each time) is accurate and I've re-confirmed it myself above; the extra trailing ranges appear to be stale artifacts from an earlier draft and should simply be dropped. This doesn't change any technical conclusion, but I want to flag it plainly rather than silently propagate an inaccurate citation into another round. ## 1.2 Compatibility Impact **Fully compatible.** Purely additive REST surface (`/running-jobs/slot-usage` v2, `/hazelcast/rest/maps/running-jobs/slot-usage` v1); new serializer class id `17`, independently re-verified as still non-colliding against the current `dev` HEAD (48 commits ahead, top id still `16`); no existing endpoint, config `Option`, default value, checkpoint/savepoint format, or job-state serialization is touched. `SlotUsage.toResponse()` (`RunningJobSlotUsageBuilder.java:199`) does `String.valueOf(jobId)` — I confirmed both v1 and v2 converge on this single code path, so the string-typed `jobId` is guaranteed identical by construction, and it's the right call since Zeta job ids are Snowflake-style longs that can exceed `2^53` and would silently lose precision as a raw JSON number in JS consumers. One rolling-upgrade nuance worth naming explicitly: a new-version node forwarding this new operation to an old-version master that doesn't know class id `17` would surface `IllegalArgumentExceptio n("Unknown type id 17")` on that node — but that's inherent to adding any new cross-version operation and doesn't regress anything that works today. ## 1.3 Performance / Side-Effect Analysis See Issues 2 and 3 above for the concurrency-relevant findings (per-job blocking sleep, on a shared operation thread). Beyond that: the full `IMAP_OWNED_SLOT_PROFILES` map is iterated on every request (`RunningJobSlotUsageBuilder.java:88-93`) with no cache or rate limit (Issue 4) — for a cluster with many pipelines this is a real but bounded cost, not a leak. No new locking, no retries, no idempotency concerns (the endpoint is read-only), no resource-release concerns (no resources are acquired). ## 1.4 Error Handling and Logging `RunningJobSlotUsageService.java` (44 lines, read in full) has zero `LOGGER` usage anywhere in the new code path. The non-master forwarding branch does `NodeEngineUtil.sendOperationToMasterNode(...).join()` with no `try/catch` — an exception during mastership loss on that join will propagate unhandled and surface as a generic 500 with no diagnostic log line explaining why (Issue 5, Low — this matches the pre-existing pattern of sibling servlets in this module, so it isn't a regression, just a missed opportunity to do better). Formal issues, sorted by severity: **Issue 1 (Medium): `slotCount: 0` is indistinguishable from "the assigned-slot source wasn't ready"** - Location: `RunningJobSlotUsageBuilder.java:68-73` (build resource-manager-null fallback), `132-134` (`isAssignedToJob` filter). - Best improvement: add a `degraded`/`slotSourceReady` boolean to the response when `resourceManager == null`, or fall back to counting directly from `IMAP_OWNED_SLOT_PROFILES` in that case with a flag noting it's an owned-not-assigned count. Add a unit test for this path. - Raised by another reviewer: No. **Issue 2 (Medium): `getCoordinatorService()` invoked per running job inside the collection lambda, can sleep up to 1.5s per call** - Location: `RunningJobSlotUsageBuilder.java:58-66` (call site); `SeaTunnelServer.java:286-315` (retry/sleep logic I traced independently: `maxRetry=3`, `retryPause=500`ms). - Best improvement: hoist `server.getCoordinatorService()` out of the lambda into a single local variable before the loop. - Raised by another reviewer: No. **Issue 3 (Medium): aggregation runs inline on the invoking thread instead of a dedicated executor, unlike the directly comparable sibling operation** - Location: `GetRunningJobSlotUsageOperation.java:48-52`; compare `GetRunningJobMetricsOperation.java`'s `CompletableFuture.supplyAsync(..., getExecutor("get_running_job_metrics_operation"))`. - Best improvement: follow the same pattern — offload `RunningJobSlotUsageBuilder.build(service)` onto a named executor and complete the operation's response asynchronously. - Raised by another reviewer: No. **Issue 4 (Low): full cluster-wide owned-slot map read per request, no cache/TTL** - Location: `RunningJobSlotUsageBuilder.java:88-93`. - Raised by another reviewer: No. **Issue 5 (Low): no logging in the new path; forwarding-branch mastership loss surfaces as an unhandled 500** - Location: `RunningJobSlotUsageService.java:35-43`. - Raised by another reviewer: No. **Issue 6 (Low): no coverage of the degraded path, the public `build(SeaTunnelServer)` entry point, the serializer round-trip, or either REST layer** - Location: `RunningJobSlotUsageBuilderTest.java` (whole file, 108 lines). - Raised by another reviewer: Yes (@SEZ9, 2026-08-21) — I independently confirmed the exact same gap by reading the test file myself. **Issue 7 (Low): `slotCount` counts task-group entries, not deduplicated distinct physical slots** - Location: `RunningJobSlotUsageBuilder.java:114-118` (aggregation loop), `189-195` (`addSlot`). - Raised by another reviewer: No. **Issue 8 (Low, informational): branch drift** - Live re-check at review time: `dev...20538ae1eb` is now `ahead_by: 3, behind_by: 48` (grown from the `46` in the prior round, as expected with time). The 3 ahead-commits still touch only JDBC-connector/doc files with zero overlap with the engine files here. Not a blocker, but please rebase before the final merge so CI runs against a current base. - Raised by another reviewer: Yes (@SEZ9, 2026-08-21). # 2. Code Quality Assessment ## 2.1 Coding Standards ASF headers present on all new files; no wildcard imports; every new class carries a class-level Javadoc; non-trivial methods have explanatory comments. One real gap: `build(SeaTunnelServer)`'s Javadoc doesn't state the freshness/consistency contract between its two joined data sources (`IMAP_OWNED_SLOT_PROFILES` vs. the resource manager's assigned-slot snapshot), which is exactly the subtlety behind Issue 1 — worth a one-line addition. ## 2.2 Test Coverage and Test Stability **Rating: Stable.** I read the one existing test end to end: no `Thread.sleep`, no Awaitility/polling, no shared static state, no order-sensitive or floating-point assertions — nothing in what exists can contribute to CI flakiness. The gap is a coverage-completeness problem (Issue 6), not a stability problem. ## 2.3 Documentation Updates All four REST reference docs (`docs/en|zh/engines/zeta/rest-api-v1.md`, `rest-api-v2.md`) gained 36 lines each — I confirmed this matches the PR's file list exactly, and the example payload matches the code's actual output shape, including the `jobId`-as-string choice. Missing doc follow-ups that would help operators: the `slotCount: 0` ambiguity (Issue 1), the counting unit (Issue 7), and the per-request cost characteristic (Issue 4). # 3. Architectural Soundness ## 3.1 Elegance of the Solution **Precise fix with one debatable design choice.** The servlet -> service -> operation -> builder decomposition matches every other Zeta REST endpoint I cross-checked in this module. The two-source join (cluster-replicated `IMap` for ownership vs. a master-local, heartbeat-refreshed snapshot for assignment) is architecturally honest about what each source represents, but its weaker source (the assigned-slot snapshot) goes empty in exactly the failover scenario this endpoint exists to help diagnose (Issue 1). ## 3.2 Maintainability Small, focused classes; clear names; no shared mutable state; the aggregation logic is already unit-testable in isolation (proven by the existing test using the package-private overload) — closing Issue 6 is straightforward from here. ## 3.3 Extensibility Solid foundation for a future Web UI panel or slot-pressure alerting on top of this data. ## 3.4 Historical-Version Compatibility No existing REST contract, config option, serialized job/checkpoint/savepoint state is touched. New serializer id `17` re-verified collision-free against current `dev`. No `docs/en/introduction/concepts/incompatible-changes.md` entry is required. # 4. Issue Summary | # | Issue | Location | Severity | Raised by another reviewer | | --- | --- | --- | --- | --- | | 1 | `slotCount: 0` indistinguishable from "assigned-slot source not ready" during failover/cold-start | `RunningJobSlotUsageBuilder.java:68-73, 132-134` | Medium | No | | 2 | `getCoordinatorService()` called per running job in the collection lambda, can sleep up to 1.5s/call | `RunningJobSlotUsageBuilder.java:58-66`; `SeaTunnelServer.java:286-315` | Medium | No | | 3 | Aggregation runs inline on the invoking operation thread instead of a dedicated executor, unlike the comparable sibling operation | `GetRunningJobSlotUsageOperation.java:48-52` | Medium | No | | 4 | Full cluster-wide owned-slot map read per request, no cache/rate limit | `RunningJobSlotUsageBuilder.java:88-93` | Low | No | | 5 | No logging in new path; forwarding-branch mastership loss is an unhandled 500 | `RunningJobSlotUsageService.java:35-43` | Low | No | | 6 | No coverage of degraded path, public entry point, serializer round-trip, or REST layer | `RunningJobSlotUsageBuilderTest.java` | Low | Yes (@SEZ9) | | 7 | `slotCount` counts task-group entries, not deduplicated distinct slots | `RunningJobSlotUsageBuilder.java:114-118, 189-195` | Low | No | | 8 | Branch now 48 commits behind `dev`; rebase before merge | branch state | Low | Yes (@SEZ9) | # 5. Merge Recommendation ### Conclusion: Ready to merge after fixes **1. Blockers — must be fixed before merge:** - Issue 1 — make "the slot source wasn't ready" distinguishable from "genuinely zero slots," with a regression test for the degraded path. - Issue 2 — hoist `getCoordinatorService()` out of the per-job lambda; a trivial two-line fix. - Issue 3 — offload the aggregation off the shared operation thread, following the existing `GetRunningJobMetricsOperation` precedent in the same package. **2. Recommended fixes — non-blocking:** - Issue 4 — narrow the `IMap` read or add a short-TTL cache. - Issue 5 — add a WARN log on the degraded/forwarding-failure path. - Issue 6 — add the REST-layer IT and the degraded-path unit cases; I'd treat this as a practical precondition alongside the three Medium fixes even though its own severity is Low, since a Medium-severity behavior fix (Issue 1) without a regression test for it isn't really closed. - Issue 7 — de-duplicate by `SlotKey`, or rename the field so the counting unit is unambiguous. - Issue 8 — rebase onto current `dev` before the final merge. **Overall assessment.** I re-derived every technical conclusion above directly from source rather than trusting the PR's own review history, and I reach the same bottom line the last round did: the feature is genuinely useful and the plumbing (routing, serializer id, compatibility) is sound, but the endpoint currently behaves worst exactly when it would be needed most — during a master/coordinator transition it can silently report all-zero slot counts while blocking a thread shared with other cluster-internal work for up to 1.5s per running job, with nothing logged to explain why. Live CI at review time is green (`Build`: SUCCESS on the fork run backing this head), and `mergeStateStatus` is `BLOCKED` only because self-approval is blocked on this account, not because of any check failure. Since this is my own PR, I'm posting this as a plain comment rather than a formal review/approval, as GitHub does not permit self-approval here — a maintainer with approval rights is the remai ning gate once Issues 1-3 are addressed. -- 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]
