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]

Reply via email to