DanielLeens commented on PR #11597: URL: https://github.com/apache/seatunnel/pull/11597#issuecomment-5340775378
Re-checking my own PR today. Since GitHub does not allow self-approval, I'm leaving this as a plain comment with my technical conclusion instead of a formal review. I independently re-traced the whole diff against the current head `20538ae1ebe78cc2d35871fbd02e1d79b960e154` rather than trusting my previous write-up. The diff stat is byte-for-byte identical to what I reviewed on 2026-08-16 (`git diff --stat dev...20538ae1e` = 13 files, 626 insertions, 0 deletions, `dev...head` compare `status=diverged, ahead_by=3, behind_by=28`, and all 3 ahead-commits touch only JDBC-connector/docs paths with zero overlap with the engine files below), so there is no new commit to re-review here. What has changed is the CI signal, and my last comment's CI note is now stale, so I want to correct it explicitly rather than leave outdated guidance on the page. # What Problem Does This PR Solve? - User pain point: Zeta exposes cluster-wide slot totals through `/overview`, but there was no endpoint that answers "how many slots is each running job holding right now, and on which workers". Operators diagnosing a stuck or over-provisioned cluster had to infer it from logs. - Fix approach: add a read-only REST endpoint `/running-jobs/slot-usage` (Jetty v2) and `/hazelcast/rest/maps/running-jobs/slot-usage` (Hazelcast v1), backed by a master-side operation and a stateless aggregation helper that intersects the coordinator's owned-slot map with the resource manager's assigned-slot snapshot, grouped by job, pipeline and worker. - One-sentence summary: this PR adds a purely additive, read-only per-job slot usage diagnostics endpoint to the Zeta REST API. 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, re-traced against the current head with `path:line`: ```text HTTP GET /running-jobs/slot-usage (Jetty v2) -> JettyService.java:190-231, registered under the same BasicAuthFilter as every other REST endpoint (JettyService.java:172-175, filter is mapped to "/*", so this endpoint is not an auth gap) -> RunningJobSlotUsageServlet.doGet() -> RunningJobSlotUsageService.getRunningJobSlotUsageJson() -- this node IS master -> RunningJobSlotUsageBuilder.build(server) [in-process, Jetty thread] -- this node is NOT master -> forward GetRunningJobSlotUsageOperation to master -> GetRunningJobSlotUsageOperation.run() [runs inline on a Hazelcast generic operation thread] -> RunningJobSlotUsageBuilder.build(server) RunningJobSlotUsageBuilder.build(SeaTunnelServer) [RunningJobSlotUsageBuilder.java:48-75] step 1 IMAP_RUNNING_JOB_INFO.keySet().forEach(jobId -> server.getCoordinatorService().shouldShowAsRunningJob(jobId)) [:58-66] step 2 ResourceManager resourceManager = server.getCoordinatorService().getInitializedResourceManager(); [:68-69] assignedSlots = resourceManager == null ? emptyList() : resourceManager.getAssignedSlots(emptyMap()) [:70-73] step 3 IMAP_OWNED_SLOT_PROFILES.forEach(...) -> aggregatePipelineSlots(...) -> keep only slots whose SlotKey is in assignedSlotSet [:88-119] step 4 SlotUsage.addSlot() -> slotCount / pipelineSlotCounts / workerSlotCounts -> toResponse() -> JsonUtils.toJsonString ``` What "slot usage" actually means here: the numerator is the set of entries in `IMAP_OWNED_SLOT_PROFILES` (coordinator-side, cluster-replicated) that also appear in `ResourceManager.getAssignedSlots(...)` (master-local, in-memory `registerWorker` snapshot). That is the intersection of two sources with different freshness characteristics, and it is only as good as the weaker one. The docs (`docs/en/engines/zeta/rest-api-v2.md:336-339`) describe it as "the number of assigned task group slots owned by the running job" and "a running job that has not received any slot yet is returned with `slotCount` set to `0`" — that text matches what the code does, so docs and implementation agree; the gap is between that behavior and what a diagnostics API user needs, covered in Issue 1/7 below. Read-only verification (re-checked, not assumed): every call in the new path — `IMap.keySet()`, `IMap.forEach`, `CoordinatorService.shouldShowAsRunningJob`, `getInitializedResourceManager()` (a volatile field read, `CoordinatorService.java:138`), `AbstractResourceManager.getAssignedSlots` — is a read. Nothing calls `requestResource`, `releaseResource`, `slotActiveCheck`, `heartbeat`, or `ResetResourceOperation`. There is no slot allocation/deallocation side effect from querying this endpoint. Iteration safety under concurrent job submission/completion, re-verified: no `ConcurrentModificationException` is possible. `IMap.forEach` walks a materialized `entrySet()` snapshot, not a live view; `AbstractResourceManager.registerWorker` is a `ConcurrentHashMap` with a weakly-consistent iterator. Torn reads are possible (the running-job id set, the assigned-slot snapshot, and the owned-slot map are read at three different instants with no consistency barrier), which is acceptable for a diagnostics endpoint but is not distinguishable from Issue 1's degradation, which is the actual correctness concern. ## 1.2 Compatibility Impact **Fully compatible.** Purely additive: new constant, new servlet, new service, new operation, new serializer class id (`17`, no collision — current `dev` tops out at `16`). No existing constant, servlet mapping, service method, or response shape is modified (13 files, 626 insertions, 0 deletions). No config option added or changed. No checkpoint/savepoint/state-serialization surface touched. Rolling-upgrade note: class id `17` is only ever sent worker -> master; a mixed-version cluster with an old-version master would reject it with `IllegalArgumentException("Unknown type id 17")`, but that failure is confined to this new endpoint and breaks no pre-existing path. ## 1.3 Performance/Side-Effect Analysis - CPU/memory: linear aggregation, nothing quadratic, small maps/sets — fine. - The real cost is data movement and thread blocking, not CPU, and it scales with total cluster state rather than with the number of running jobs, on every request, with no caching or rate limit (Issue 4). - `getCoordinatorService()` runs once per running job inside the `forEach` filter (`RunningJobSlotUsageBuilder.java:59-66`); it is not a plain getter — `SeaTunnelServer.getCoordinatorService()` can retry up to 3 times with `Thread.sleep(500)` while the coordinator is not yet active. With N running jobs during a coordinator hiccup, one request turns into up to N × 1.5s of blocking (Issue 2). - The whole aggregation runs inline on a Hazelcast generic operation thread in `GetRunningJobSlotUsageOperation.run()` (`GetRunningJobSlotUsageOperation.java:47-52`) — a small, shared, cluster-wide pool — unlike the directly comparable `GetRunningJobMetricsOperation.run()`, which deliberately offloads to `getNodeEngine().getExecutionService().getExecutor("get_running_job_metrics_operation")`. The new operation also implements `AllowedDuringPassiveState`, so it can be invoked exactly when the coordinator is most likely inactive and the sleep path most likely to trigger (Issue 3). - Retry/idempotency: read-only, idempotent. Resource release: nothing to release. ## 1.4 Error Handling and Logging - The entire new path (servlet, service, operation, builder) contains zero log statements, so the degraded-zero condition in Issue 1 leaves no trace anywhere. - `RunningJobSlotUsageService.getRunningJobSlotUsageJson()` calls `getSeaTunnelServer(true)` and `getSeaTunnelServer(false)` as two independent calls and `.join()`s the forwarded invocation with no try/catch; a mastership change between the two calls surfaces as an unhandled `CompletionException`/`SeaTunnelEngineException` -> Jetty 500. This mirrors the existing pattern in `RunningJobsServlet`/`JobInfoService`, so it is a consistency-level nit rather than a regression introduced by this PR — flagged at Low. Formal issues (verified against the current head, all still present in unmodified code, same numbering as my 2026-08-16 comment): **Issue 1: The endpoint silently reports `slotCount: 0` when the assigned-slot source is unavailable, and `0` is also a documented legitimate value** - Location: `seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/trace/RunningJobSlotUsageBuilder.java:68-73, 115, 132-134` - Problem: `getInitializedResourceManager()` deliberately never lazy-inits and returns `null` whenever the resource manager has not been touched yet on this master (cold master, just after failover); the builder then forces `assignedSlots = Collections.emptyList()`, and every running job renders as `slotCount: 0` even though `IMAP_OWNED_SLOT_PROFILES` already holds full ownership data. In steady state `registerWorker` is refreshed synchronously on assign/release, so this is a failover/cold-start window rather than a permanent staleness problem, but the docs (`docs/en/engines/zeta/rest-api-v2.md:339`) make `0` a legitimate business value, so a caller can't tell "job genuinely holds no slots" from "we couldn't read the slot source at all". - Risk: an operator investigating a master failover — exactly when this API is most useful — sees `slotCount: 0` for every still-running job and may conclude the jobs lost their resources, with no log line to contradict it. - Best improvement: make unavailability explicit (a `degraded`/`slotSourceAvailable` field), or fall back to counting directly from `IMAP_OWNED_SLOT_PROFILES` when `resourceManager == null`, and log once at WARN. Add a unit test for the degraded case. - Severity: Medium - Raised by another reviewer: No (carryover, first raised in my 2026-08-03 comment, still unresolved) **Issue 2: `getCoordinatorService()` is invoked once per running job inside the collection lambda, and it can sleep** - Location: `RunningJobSlotUsageBuilder.java:59-66` - Problem: the filter lambda calls `server.getCoordinatorService().shouldShowAsRunningJob(jobId)` once per entry in `IMAP_RUNNING_JOB_INFO`; `getCoordinatorService()` can retry up to 3× with `Thread.sleep(500)` while the coordinator is not active, logging a WARNING each time. The reference is loop-invariant. - Risk: on a large cluster during master transition, a single REST call blocks its thread for tens of seconds and floods the master log; combined with Issue 3 this thread is a Hazelcast operation thread. - Best improvement: `CoordinatorService coordinatorService = server.getCoordinatorService();` once before the loop, reuse for both the filter and the `getInitializedResourceManager()` call at line 69. - Severity: Medium - Raised by another reviewer: No (carryover) **Issue 3: The master operation performs the whole aggregation inline on a Hazelcast operation thread, unlike its sibling operation** - Location: `GetRunningJobSlotUsageOperation.java:47-52` - Problem: `run()` calls `RunningJobSlotUsageBuilder.build(service)` directly on the invoking generic operation thread. `GetRunningJobMetricsOperation.run()` in the same package deliberately offloads its body to a named executor for exactly this reason, and the new operation also carries `AllowedDuringPassiveState`, so it is most likely to run this blocking work exactly when the coordinator is inactive. - Risk: blocking generic operation threads degrades unrelated cluster operations, not just this endpoint — under a polling UI plus a coordinator transition this is a plausible route to cluster-wide operation-thread starvation. - Best improvement: offload `build(...)` to `getNodeEngine().getExecutionService().getExecutor(...)`, mirroring `GetRunningJobMetricsOperation`. - Severity: Medium - Raised by another reviewer: No (carryover) **Issue 4: Every request materializes the entire owned-slot map for the whole cluster, including jobs it will discard** - Location: `RunningJobSlotUsageBuilder.java:50-56, 88-92` - Problem: `ownedSlotProfiles.forEach(...)` resolves to `Map.forEach`, which walks `entrySet()` — a distributed operation fetching and deserializing every pipeline entry in the cluster, even entries for jobs later discarded because they're not running. No caching, no TTL, no rate limit. - Risk: on a cluster with many pipelines, a tight polling loop turns into repeated full-map deserialization on the master, amplified by Issue 3. - Best improvement: restrict the fetch to running-job keys (`ownedSlotProfiles.getAll(...)` or a predicate read), or add a short-TTL cache in front of the aggregation. - Severity: Low - Raised by another reviewer: No (carryover) **Issue 5: No logging anywhere in the new path, and the forwarded call surfaces mastership loss as a 500** - Location: `RunningJobSlotUsageService.java:35-43`; `RunningJobSlotUsageBuilder.java:48-75` - Problem: zero log statements in the whole feature, and `getSeaTunnelServer` is called twice without capturing a single reference, so a mastership change between the calls surfaces as an unhandled 500 rather than a clean transient error. - Risk: a monitoring scrape sees an opaque 500 during a routine master change, and the silent-zero path (Issue 1) is undiagnosable after the fact. - Best improvement: capture the server reference once, wrap the `.join()`/build in a try/catch that logs at WARN, and log once when `getInitializedResourceManager()` returns `null` while running jobs exist. - Severity: Low - Raised by another reviewer: No (carryover) **Issue 6: Test coverage does not reach the failure modes, the public entry point, or the REST layer** - Location: `RunningJobSlotUsageBuilderTest.java:40-82` - Problem: the single test exercises the package-private 3-arg overload with real value assertions (good), but does not cover `assignedSlots` empty/null while `ownedSlotProfiles` is populated (Issue 1's exact case), the public `build(SeaTunnelServer)` entry point, a serializer round-trip for class id `17`, or any REST-level assertion. The repo has clear precedent for REST-level IT (`RestApiIT`, `PendingJobsRestIT`, `RealtimeMetricsRestIT`, `JobInfoDagStabilityRestIT` in `seatunnel-e2e/seatunnel-engine-e2e/connector-seatunnel-e2e-base`). - Risk: a routing regression, a serialization regression on id `17`, or a JSON shape change would all pass CI silently. - Best improvement: add the degraded-path unit cases plus a `RunningJobSlotUsageRestIT` asserting `slotCount` against a real job's parallelism. - Severity: Low - Raised by another reviewer: No (carryover) **Issue 7: `slotCount` counts task-group entries, not distinct slots, which the endpoint name does not convey** - Location: `RunningJobSlotUsageBuilder.java:113-118, 189-195` - Problem: `addSlot` increments once per `TaskGroupLocation` entry that passes the filter, not per distinct `SlotProfile`. The docs word this carefully ("assigned task group slots"), so docs and code agree; the mismatch is between the code and the endpoint's/field's own name, which both read as "distinct slots occupied". - Risk: capacity-planning conclusions drawn from `slotCount` could overstate real slot occupancy, invisibly, because the two numbers agree in the common case. - Best improvement: de-duplicate by `SlotKey` before counting, or rename the field to `taskGroupSlotCount` so the unit is unambiguous. - Severity: Low - Raised by another reviewer: No (carryover) # 2. Code Quality Assessment ## 2.1 Coding Standards ASF license headers present on all 9 new source files. Formatting matches AOSP/Spotless, no wildcard imports. Comment discipline is good — every new class has a class-level Javadoc and the non-trivial methods carry explanatory comments; the one gap is that `build(SeaTunnelServer)` doesn't document the freshness contract of the two data sources it joins, which Issue 1 makes worth calling out explicitly in the comment once fixed. ## 2.2 Test Coverage and Test Stability Unit tests: one, on the package-private overload, with real value assertions and deterministic `TreeMap` ordering. Integration/REST tests: none (Issue 6). Flaky-pattern check: clean — no `Thread.sleep`, no `Awaitility`, no network/container dependency, no shared static state, no ordering dependency beyond the deterministic map. **Test stability rating: Stable.** Nothing in the existing test contributes flakiness risk; the gap is coverage of the failure modes, not stability of what's there. ## 2.3 Documentation Updates All four REST reference files updated (`docs/en` and `docs/zh`, both v1 and v2). The zh docs are a real translation, not a copy. The example payload matches the code exactly. Missing pieces are downstream of the code fixes: the `slotCount: 0` ambiguity (Issue 1), the counting unit (Issue 7), and the per-request cost characteristic (Issue 4) are not yet spelled out in the docs, but that's a doc follow-up rather than an independent doc bug. # 3. Architectural Soundness ## 3.1 Elegance of the Solution **Precise fix within a temporary-workaround aggregation choice.** The servlet -> service -> operation -> builder decomposition matches every other Zeta REST endpoint and is a clean, minimal addition. The one architectural choice I'm still not fully comfortable with is intersecting a cluster-replicated IMap with a master-local in-memory snapshot as the source of truth — the answer is only as good as the weaker source, and the weaker source is empty in exactly the scenario (failover) the API exists to diagnose. ## 3.2 Maintainability Good — small classes, clear names, no shared mutable state, easy to unit test. The main risk is reputational: if operators learn this endpoint sometimes reports unexplained zeros, the API loses its diagnostic value. Fixing Issues 1 and 5 removes that risk. ## 3.3 Extensibility A per-job, per-pipeline, per-worker breakdown is a good foundation for a future UI panel or slot-pressure alerting. Once something polls it, Issues 3 and 4 stop being theoretical. ## 3.4 Historical-Version Compatibility No existing REST contract, config option, default value, serialized state, checkpoint, or savepoint format is touched; old clients are unaffected. New serializer id `17` doesn't collide with anything on current `dev`. No upgrade note or `incompatible-changes.md` entry is required. # 4. Issue Summary | # | Issue | Location | Severity | | --- | --- | --- | --- | | 1 | Silent `slotCount: 0` when the assigned-slot source is unavailable, indistinguishable from a documented legitimate zero | `RunningJobSlotUsageBuilder.java:68-73, 115, 132-134` | Medium | | 2 | `getCoordinatorService()` called once per running job inside the loop; can sleep up to 1.5s per call | `RunningJobSlotUsageBuilder.java:59-66` | Medium | | 3 | Aggregation runs inline on a Hazelcast operation thread; sibling metrics operation deliberately offloads | `GetRunningJobSlotUsageOperation.java:47-52` | Medium | | 4 | Full cluster-wide owned-slot map materialized per request, including discarded jobs; no cache or rate limit | `RunningJobSlotUsageBuilder.java:50-56, 88-92` | Low | | 5 | No logging in the entire new path; mastership loss surfaces as an unhandled 500 | `RunningJobSlotUsageService.java:35-43` | Low | | 6 | No coverage of the degraded path, the public entry point, serializer round-trip, or the REST layer | `RunningJobSlotUsageBuilderTest.java:40-82` | Low | | 7 | `slotCount` counts task-group entries rather than distinct slots; name implies otherwise | `RunningJobSlotUsageBuilder.java:113-118, 189-195` | Low | # 5. Merge Recommendation ### Conclusion: Ready to merge after fixes **1. Blockers — must be fixed** - Issue 1 (Medium) — a diagnostics endpoint must not encode "I couldn't read the data" as the same value it uses for "the answer is genuinely zero". Unresolved since my first pass on 2026-08-03. - Issue 2 (Medium) — hoist `getCoordinatorService()` out of the per-job lambda; a two-line fix that removes an O(N) × 1.5s blocking amplification. - Issue 3 (Medium) — offload the aggregation off the Hazelcast operation thread, following the `GetRunningJobMetricsOperation` precedent in the same package. Blocking shared operation threads is an engine-wide hazard, not an endpoint-local one. **2. Recommended fixes — non-blocking** - Issue 4 — narrow the IMap read to running-job keys or add a short-TTL cache. - Issue 5 — single server lookup, try/catch around the forwarded `join()`, WARN log when the slot source is unavailable. - Issue 6 — add the degraded-path unit cases and one REST-level IT asserting real slot counts. - Issue 7 — de-duplicate by `SlotKey` or rename the field so the counting unit is unambiguous. **CI correction (supersedes my 2026-08-16 note):** I previously wrote that the apache-side `Build` check was red because of a `paimon-connector-it` infra flake (`429 Too Many Requests` on the Maven wrapper download) and that a rerun was needed. Rechecking the same head right now, `Build` is green (`gh pr checks 11597` and the commit's check-runs both show `conclusion: success`). So that specific CI note is stale and I'm withdrawing it — please disregard "needs a rerun" from my last comment. That said, CI passing does not clear the review: all three Medium blockers above are unrelated to CI and are still present in this exact, unmodified code, so my merge conclusion is unchanged from 2026-08-16. Since GitHub blocks self-approval on my own PR, `mergeStateStatus` will keep showing `BLOCKED`/`REVIEW_REQUIRED` regardless of my own conclusion here — that gate needs an independent maintainer regardless of the fixes above. **Overall assessment:** the feature is worth having and the plumbing is genuinely solid — routing is correct on both API versions, the serializer id is free, `SlotProfile.ownerJobID` really is on the wire so the intersection works cross-node, the resource-manager field is properly `volatile`, the endpoint is verifiably free of any allocation/deallocation side effect, iteration is CME-safe, compatibility is cleanly additive, and all four doc files are accurate in both languages. What holds it back is that the three blockers cluster around the same theme: the endpoint behaves worst exactly when it's needed most — during a master transition it can return all zeros with no log line while blocking a shared cluster thread for up to 1.5s per running job. Two of the three are a handful of lines each; Issue 1 needs a deliberate decision between exposing a degraded flag or falling back to the coordinator-only view. Alternative implementation: - Option A: aggregate directly from `IMAP_OWNED_SLOT_PROFILES` as the source of truth and add a separate, explicit "assignment confirmed" indicator per slot instead of using it as a hard filter. - Option B: keep the current intersection semantics but expose a top-level `degraded: true/false` flag driven by `resourceManager == null`, so callers can programmatically distinguish "no data" from "confirmed zero" without changing the existing fields. -- 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]
