DanielLeens commented on PR #11597:
URL: https://github.com/apache/seatunnel/pull/11597#issuecomment-5342958114
Re-reviewed from scratch at head `20538ae1ebe78cc2d35871fbd02e1d79b960e154`
— this is the **same commit** I reviewed on 2026-08-16 (no new commits landed
on this branch in the last 3 days). Rather than re-post a copy of that review,
I re-traced the whole call chain independently again (own-code trust boundary
applies to me too), re-verified every claim in the 2026-08-16 comment against
the current source, and I'm folding the result into one review below so there's
a single up-to-date place to look. Where a prior finding still holds I say so
and cite what I re-checked; I did not find anything that changes the earlier
conclusions.
# What Problem Does This PR Solve?
- User pain point: Zeta's `/overview` endpoint only gives cluster-wide slot
totals; there is no way to ask "how many slots is running job X holding right
now, and on which workers/pipelines". Operators diagnosing a stuck or
over-provisioned cluster have to reconstruct this from logs.
- Fix approach: adds a read-only REST endpoint, `/running-jobs/slot-usage`
(v2, Jetty) and `/hazelcast/rest/maps/running-jobs/slot-usage` (v1), backed by
a master-side operation and 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.
# 1. Code Change Review
## 1.1 Core Logic Analysis
Runtime path, re-traced against the current head:
```text
HTTP GET /running-jobs/slot-usage (Jetty v2)
-> JettyService.java:190,230 servlet mapped via convertUrlToPath() ->
"/running-jobs/slot-usage/*"
-> RunningJobSlotUsageServlet.doGet()
[RunningJobSlotUsageServlet.java:41-45]
-> RunningJobSlotUsageService.getRunningJobSlotUsageJson()
[RunningJobSlotUsageService.java:35-43]
|
+-- this node IS master (getSeaTunnelServer(true) != null)
| -> RunningJobSlotUsageBuilder.build(server) [in-process,
on the Jetty thread]
|
+-- this node is NOT master
->
NodeEngineUtil.sendOperationToMasterNode(GetRunningJobSlotUsageOperation)
[NodeEngineUtil.java:33-41, service name =
SeaTunnelServer.SERVICE_NAME]
-> master: GetRunningJobSlotUsageOperation.run()
[GetRunningJobSlotUsageOperation.java:92-96]
runs on a Hazelcast generic operation thread (see Issue 3)
-> RunningJobSlotUsageBuilder.build(server)
HTTP GET /hazelcast/rest/maps/running-jobs/slot-usage (v1)
-> RestHttpGetCommandProcessor.handle() ordering:
RUNNING_JOBS_SUMMARY -> RUNNING_JOBS_SLOT_USAGE -> RUNNING_JOBS
[RestHttpGetCommandProcessor.java:161-166]
-> handleRunningJobsSlotUsage() -> same RunningJobSlotUsageService
Aggregation (RunningJobSlotUsageBuilder.build(server), both entry points
converge here)
step 1 IMAP_RUNNING_JOB_INFO.keySet().forEach(...)
[RunningJobSlotUsageBuilder.java:59-66]
filter:
server.getCoordinatorService().shouldShowAsRunningJob(jobId)
-> CoordinatorService.java:1502 -> getJobStatus(jobId) ->
!status.isEndState()
step 2 server.getCoordinatorService().getInitializedResourceManager()
[line 68-69]
-> plain volatile field read, no lazy init
(CoordinatorService.java:1247-1249)
resourceManager == null ? emptyList() :
getAssignedSlots(emptyMap()) [line 70-73]
-> AbstractResourceManager.java:328-332, streams
registerWorker.values()
-> WorkerProfile.getAssignedSlots() (master-local,
heartbeat-refreshed state)
step 3 ownedSlotProfiles.forEach((pipelineLocation, taskGroupSlots) ->
...) [line 88-92, 414-433]
per pipeline: keep only slots whose SlotKey is present in the
assigned-slot set
SlotKey = (worker, slotId, ownerJobId, sequence) [line 450-490]
step 4 SlotUsage.addSlot() -> slotCount / pipelineSlotCounts /
workerSlotCounts
[line 503-509] -> toResponse() [line 511-518] ->
JsonUtils.toJsonString
```
Things I checked myself rather than take on faith, because each one would
have been a real bug if wrong:
- `SlotProfile.ownerJobID` is written/read in `writeData`/`readData` —
confirmed via `grep` on the constructor/serialization methods. Since `SlotKey`
includes `ownerJobId`, a non-serialized owner id would make the intersection
miss on every cross-node call and `slotCount` would be permanently `0` on any
real (multi-node) cluster. It's on the wire, so this genuinely works.
- `PipelineLocation` carries `jobId` as one of its two fields
(`PipelineLocation.java`), so `usageByJob.get(pipelineLocation.getJobId())` at
line 423 can never attribute one job's owned-slot entries to a different job's
`SlotUsage`. No cross-job leakage from the primary grouping key.
- New serializer class id `17` (`GET_RUNNING_JOB_SLOT_USAGE_OPERATION`) does
not collide — `ClientToServerOperationDataSerializerHook` on `dev` tops out at
`16` (`GET_JOB_TASK_MAPPING_OPERATION`); this PR is insert-only there.
- Servlet routing is correct on both API versions: `convertUrlToPath()`
appends `/*`, so `/running-jobs/slot-usage/*` is a strictly longer path than
`/running-jobs/*` and wins under the servlet longest-prefix rule; and the v1
`startsWith` chain places `REST_URL_RUNNING_JOBS_SLOT_USAGE` before the more
general `REST_URL_RUNNING_JOBS`, so it's not swallowed by the existing branch —
same technique already used for `REST_URL_RUNNING_JOBS_SUMMARY`.
- `getService()` resolves without `GetRunningJobSlotUsageOperation`
overriding `getServiceName()`, because
`NodeEngineUtil.sendOperationToMasterNode` builds the invocation with
`SeaTunnelServer.SERVICE_NAME` explicitly — same pattern as the other no-arg
operations in this package (e.g. `GetNodeHttpPortOperation`).
- Read-only, confirmed by walking every call the new code makes:
`IMap.keySet()`, `IMap.forEach`, `shouldShowAsRunningJob` (reads only),
`getInitializedResourceManager()` (field read),
`AbstractResourceManager.getAssignedSlots` (stream over `registerWorker`).
Nothing calls `requestResource`, `releaseResource`, `slotActiveCheck` or
`heartbeat`. No slot allocation/deallocation side effect from querying this
endpoint.
- `getInitializedResourceManager()` (not the lazy `getResourceManager()`) is
the right call for a read path — the sibling `GetOverviewOperation.java:82`
still uses the lazy-init variant, so `/overview` can trigger cluster RPCs and
runtime-state creation just from being polled. This endpoint deliberately
avoids that. Good instinct; see Issue 1 for the cost of that choice.
What "usage" actually measures, verified against the code rather than the
field names: the numerator is the set of `IMAP_OWNED_SLOT_PROFILES` entries
(coordinator-side, cluster-replicated, written at deploy time) that also appear
in `ResourceManager.getAssignedSlots(...)` (master-local, heartbeat-refreshed
`registerWorker` snapshot). That's "coordinator-owned task-group slots the
current master also currently believes are assigned" — an intersection of two
views with different freshness characteristics, not a single authoritative
read. The docs phrase this accurately ("assigned task group slots owned by the
running job"), so docs and code agree; the residual gap is between the code and
the endpoint's own name/field name (see Issue 7).
## 1.2 Compatibility Impact — **Fully compatible**
- Purely additive: new constant, new servlet, new service, new operation,
new serializer class id. The diff on `JettyService`,
`RestHttpGetCommandProcessor`, `RestConstant`,
`ClientToServerOperationDataSerializerHook` is insert-only (13 files, 626
additions, 0 deletions — confirmed via `git diff dev...HEAD --stat`).
- No existing endpoint's response shape changed (`/running-jobs`,
`/running-jobs/summary`, `/overview`, `/job-info` untouched).
- No config option added, renamed, or re-defaulted.
- No checkpoint/savepoint/state-serialization/job-restore surface touched.
- Rolling-upgrade nuance worth one line in the PR description: operation
class id `17` only ever travels worker -> master. In a mixed-version cluster
where a new-version node serves this REST call but an old-version node is
master, the old master's `ClientToServerOperationDataSerializerHook` factory
throws `IllegalArgumentException("Unknown type id 17")`. That failure is
confined to this new endpoint (no old client calls it), so it doesn't break any
existing upgrade path, but it's worth documenting.
## 1.3 Performance / Side-Effect Analysis
- Iteration safety re-checked: `IMap.forEach` walks a materialized
`entrySet()` snapshot, not a live view;
`AbstractResourceManager.registerWorker` is a `ConcurrentHashMap` with a
weakly-consistent iterator. `SlotUsage` instances are confined to one call, no
shared mutable state, no locks. No `ConcurrentModificationException` is
possible.
- Torn reads are real and inherent to the design: the running-job id set,
the assigned-slot snapshot, and the owned-slot map are read at three different
instants with no consistency barrier between them. A job that completes
mid-call can show a stale count; a job that starts mid-call can show a partial
one. Acceptable for a diagnostics endpoint, but combined with Issue 1 below,
the caller can't tell a torn read from a real zero.
- Cost re-verified: `IMAP_OWNED_SLOT_PROFILES.forEach` (line 88-92) pulls
and deserializes **every** pipeline entry in the cluster to the calling member,
including entries for jobs immediately discarded at line 423-426 because
they're not in `runningJobIds`. `IMAP_RUNNING_JOB_INFO.keySet()` is likewise a
distributed call. So per-request wire cost scales with total cluster state, not
with the number of running jobs, on every single request — no cache, no TTL, no
rate limit. To be fair to the PR, this same "read everything, filter after"
shape already exists elsewhere in the engine for periodic/internal paths; the
difference here is this one is externally, arbitrarily triggerable.
- The bigger cost isn't CPU, it's thread blocking, covered as Issues 2 and 3
below.
## 1.4 Error Handling and Logging
- No per-job 404 path exists by design — the endpoint always returns the
full running-job list, and a job that ends mid-call just comes back with a
stale/zero count rather than throwing. That's a reasonable shape for a
bulk-listing diagnostic endpoint.
- Re-verified: the whole new path (servlet, service, operation, builder)
contains zero log statements. Combined with Issue 1, the silent-zero condition
leaves no trace anywhere.
- Re-verified: `RunningJobSlotUsageService.getRunningJobSlotUsageJson()`
calls `getSeaTunnelServer(true)` and, on the forwarding branch, `.join()`s the
master invocation with no try/catch (`RunningJobSlotUsageService.java:37-41`).
If mastership is lost between that check and the local `build()` call on the
other branch, `getCoordinatorService()` throws `SeaTunnelEngineException`, and
that reaches Jetty as an unhandled 500. I checked `BaseServlet.doGet`
implementations for a sibling (`OverviewServlet`) and confirmed this exact gap
already exists there too — this is not a regression introduced by this PR, it's
the established (imperfect) pattern in this REST layer. I'm flagging it as Low
for that reason, matching my own 2026-08-16 assessment.
---
**Issue 1: `slotCount: 0` is both the "no slots yet" answer and the
"couldn't read the data" answer**
- **Location:**
`seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/trace/RunningJobSlotUsageBuilder.java:68-73`,
`:446-448`
- **Problem description:** `getInitializedResourceManager()` never
lazy-inits by design, so it's `null` whenever the resource manager hasn't been
touched yet on this master (fresh master after failover, or before the first
slot request). When that happens, `assignedSlots` is forced to
`Collections.emptyList()`, the assigned-slot set is empty, `isAssignedToJob`
rejects every entry, and every running job renders with `slotCount: 0` — even
though `IMAP_OWNED_SLOT_PROFILES` may already hold full ownership data for
those jobs. The docs (`docs/en/engines/zeta/rest-api-v2.md`) explicitly make
`0` a legitimate value ("a running job that has not received any slot yet is
returned with `slotCount` set to `0`"), so the response has no way to
distinguish "genuinely zero" from "couldn't verify".
- **Potential risk:** the moment an operator is most likely to query this
endpoint — right after a master failover, while diagnosing whether jobs kept
their resources — is exactly the moment it's most likely to lie and say every
job has zero slots, with nothing logged to contradict it.
- **Best improvement:** make unavailability explicit instead of encoding it
as `0`. Option A: add a `slotSourceAvailable`/`degraded` boolean to each
response entry when `resourceManager == null`, and document it. Option B: when
the resource manager isn't initialized yet but `ownedSlotProfiles` is
non-empty, fall back to counting directly from `IMAP_OWNED_SLOT_PROFILES` and
mark that response as coordinator-only/unverified. Either way, log once at WARN
with the affected job count. Add a unit test for this exact case (resource
manager `null`, `ownedSlotProfiles` non-empty).
- **Severity:** Medium
**Issue 2: `getCoordinatorService()` runs once per running job inside the
filter lambda, and it can sleep**
- **Location:**
`seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/trace/RunningJobSlotUsageBuilder.java:59-66`
- **Problem description:** re-verified line by line —
`runningJobInfo.keySet().forEach(jobId -> { if
(server.getCoordinatorService().shouldShowAsRunningJob(jobId)) ... })` calls
`getCoordinatorService()` once per entry in `IMAP_RUNNING_JOB_INFO`. That's not
a cheap accessor: `SeaTunnelServer.getCoordinatorService()`
(`SeaTunnelServer.java:286-303`) runs up to 3 retries of `Thread.sleep(500)`
plus a WARNING log per iteration while the coordinator isn't active yet. With N
running jobs and a coordinator that's briefly inactive (e.g. mid-transition), a
single REST call can turn into up to N x 1.5s of sleeping plus N log lines. The
reference is loop-invariant — trivial to hoist.
- **Potential risk:** on a cluster with many running jobs during a master
transition, one REST call blocks its thread for a long time and floods the log;
combined with Issue 3, that thread is a shared Hazelcast operation thread, not
a request-scoped one.
- **Best improvement:** resolve `CoordinatorService coordinatorService =
server.getCoordinatorService();` once before the loop and reuse the local
reference for both the `shouldShowAsRunningJob` filter and the
`getInitializedResourceManager()` call. Two-line fix, removes the O(N)
amplification entirely.
- **Severity:** Medium
**Issue 3: The master operation aggregates inline on a shared Hazelcast
operation thread; the directly comparable sibling operation deliberately
doesn't**
- **Location:**
`seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/operation/GetRunningJobSlotUsageOperation.java:92-96`
- **Problem description:** `run()` calls
`RunningJobSlotUsageBuilder.build(service)` directly on the invoking thread.
The operation declares no partition id, so it runs on a small, shared,
cluster-wide generic operation thread pool. I compared this against
`GetRunningJobMetricsOperation` (same package) and confirmed it deliberately
offloads its body via `CompletableFuture.supplyAsync(...,
getNodeEngine().getExecutionService().getExecutor("get_running_job_metrics_operation"))`
— a dedicated named executor — precisely to keep this kind of work off the
shared generic pool (the calling thread still blocks on `future.get()`, but a
saturated dedicated executor only stalls this operation type, not every
unrelated cluster operation sharing the generic pool).
`GetRunningJobSlotUsageOperation` also declares `AllowedDuringPassiveState`,
i.e. it can run while the cluster is passive — exactly when the coordinator is
most likely to be inactive and the Issue 2 sleep path most likely to fir
e.
- **Potential risk:** blocking generic operation threads degrades unrelated
cluster operations, not just this endpoint. A polling UI/monitor plus a
coordinator transition is a plausible route to operation-thread pressure across
the cluster, which on a production Zeta cluster is a job-availability concern,
not just a monitoring inconvenience.
- **Best improvement:** mirror `GetRunningJobMetricsOperation` — offload
`build(...)` to a named executor from `getNodeEngine().getExecutionService()`.
Fixing Issue 2 reduces the blocking duration but doesn't remove the
distributed-map read cost from the shared operation thread, so both fixes are
needed together.
- **Severity:** Medium
**Issue 4: Every request materializes the entire cluster-wide owned-slot
map, including entries it immediately discards**
- **Location:** `RunningJobSlotUsageBuilder.java:88-92`, `:423-426`
- **Problem description:** `ownedSlotProfiles.forEach(...)` on a Hazelcast
`IMap` resolves to `Map.forEach` over `entrySet()` — a distributed read that
pulls and deserializes every pipeline's `Map<TaskGroupLocation, SlotProfile>`
in the cluster to the calling member, before line 423-426 throws away
everything not belonging to a running job. No caching, no TTL, no rate limit on
the endpoint itself.
- **Potential risk:** on a cluster with many pipelines, a tight polling loop
turns into repeated full-map deserialization on the master, made worse by Issue
3 since it happens on a shared operation thread.
- **Best improvement:** narrow the read to running-job keys first (e.g. a
predicate-based `IMap` read, or fetch by known `PipelineLocation` keys once
`runningJobIds` is known) instead of a full `forEach`; alternatively a
short-TTL cache in front of the aggregation, following the existing running-job
DAG cache pattern already used elsewhere in `BaseService`.
- **Severity:** Low
**Issue 5: No logging anywhere in the new path; mastership loss on the
forwarding branch surfaces as an unhandled 500**
- **Location:** `RunningJobSlotUsageService.java:35-43`;
`RunningJobSlotUsageBuilder.java:48-75`
- **Problem description:** covered under 1.4 above — zero log statements in
the whole feature, and an un-caught `.join()` on the forwarding path. Confirmed
this matches the pre-existing pattern in sibling servlets, so it's a
consistency nit rather than a regression, which is why I keep it Low rather
than blocking on it.
- **Potential risk:** a monitoring scrape sees an opaque 500 during a
routine master transition instead of a clean transient error, and the Issue 1
silent-zero condition is undiagnosable after the fact because nothing was
logged.
- **Best improvement:** capture the server reference once instead of two
separate `getSeaTunnelServer` calls; wrap the forwarded `.join()` in a
try/catch that logs at WARN; add one WARN when
`getInitializedResourceManager()` returns `null` while running jobs exist.
- **Severity:** Low
**Issue 6: Test coverage doesn't reach the degraded path, the public entry
point, or the REST layer**
- **Location:**
`seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/trace/RunningJobSlotUsageBuilderTest.java`
- **Problem description:** the one test that exists is genuinely good: it
exercises the package-private 3-arg overload, asserts real numeric values
rather than a bare "200 OK", and correctly covers the "slot reused by a
different job" race by giving the reused slot a different `ownerJobId` in the
`assignedSlots` list than the one recorded in `ownedSlotProfiles` — I confirmed
the assertion (`firstPipelineSlotCounts.containsKey(2)` is false) genuinely
exercises `SlotKey` rejecting a stale/reassigned entry, which is the single
most important correctness property for this feature (not attributing a slot to
the wrong job). What it doesn't cover: (a) the `assignedSlots` empty /
resource-manager-null degradation from Issue 1; (b) the public
`build(SeaTunnelServer)` entry point, where `getCoordinatorService()` and the
two IMap reads actually live; (c) a serializer round-trip for the new class id
`17`; (d) anything at the REST layer. I verified the precedent cited for (d) is
real: `Pendi
ngJobsRestIT.java`, `RealtimeMetricsRestIT.java`, and
`JobInfoDagStabilityRestIT.java` all exist under
`seatunnel-e2e/seatunnel-engine-e2e/connector-seatunnel-e2e-base/src/test/java/org/apache/seatunnel/engine/e2e/`,
each asserting real endpoint output for a comparable diagnostics endpoint — so
this codebase does have an established convention of adding a dedicated
`*RestIT` for a new REST diagnostics surface, this PR just didn't follow it.
(I'll also note, in fairness, that the older `/running-jobs/summary` endpoint
itself has no dedicated `RestIT`, so the convention isn't applied 100%
uniformly — but three separate precedents is enough that I'd still ask for one
here.)
- **Potential risk:** a routing regression, a serialization regression on
class id `17`, or an accidental JSON shape change would all pass CI silently —
nothing in the current test suite proves the endpoint returns anything over
HTTP at all.
- **Best improvement:** add the two missing unit cases (null resource
manager; `ownedSlotProfiles` populated but `assignedSlots` empty/stale)
asserting the response is either non-zero or explicitly marked degraded per
Issue 1's fix; add a `RunningJobSlotUsageRestIT` that submits a job and asserts
`slotCount` matches its real parallelism, following the `PendingJobsRestIT`
structure.
- **Severity:** Low
**Issue 7: `slotCount` counts task-group entries, not distinct physical
slots — the field/endpoint name implies otherwise**
- **Location:** `RunningJobSlotUsageBuilder.java:428-433`, `:503-509`
- **Problem description:** `addSlot` increments once per `TaskGroupLocation`
entry that passes the filter. If the same physical `SlotProfile` were ever
referenced by two task groups, or the same pipeline showed up under two
`PipelineLocation` keys, it would be counted twice, even though `SlotProfile`'s
own `equals`/`hashCode` (worker + slotID + sequence) would treat those as the
same slot. The docs word this carefully ("assigned task group slots"), so docs
and code agree; the mismatch is between the implementation and what a reader
expects from the endpoint's name and the bare field name `slotCount`.
- **Potential risk:** capacity-planning conclusions drawn from `slotCount`
could overstate real slot occupancy in an edge case, and the divergence would
be invisible because the two numbers agree in the common case.
- **Best improvement:** either de-duplicate by `SlotKey` before counting (a
`Set<SlotKey>` per `SlotUsage`, increment only on first insert) so the field
means what its name says, or keep current semantics and rename to
`taskGroupSlotCount` so the counting unit is unambiguous. Keep docs in sync
either way.
- **Severity:** Low
# 2. Code Quality Assessment
## 2.1 Coding Standards
- ASF license headers present on all new source files — verified.
- No wildcard imports.
- Every new class carries a class-level Javadoc, and the non-trivial methods
(`build`, `aggregatePipelineSlots`, `run`, `getRunningJobSlotUsageJson`) have
explanatory comments — this meets the documentation bar from the project guide.
Gap: none of the comments state the freshness contract between the two data
sources being joined; given Issue 1, `build(SeaTunnelServer)` should say in the
Javadoc that `getAssignedSlots` is a master-local, heartbeat-refreshed snapshot
that can be empty on a fresh master, and what the response means when that
happens.
- Package placement (`server.trace`) is consistent with the existing
`TaskMappingBuilder` in the same package — no concern there.
- `GetRunningJobSlotUsageOperation` omits the empty
`writeInternal`/`readInternal` overrides that `GetRunningJobMetricsOperation`
carries; functionally irrelevant since both just delegate to `super`, purely
cosmetic.
## 2.2 Test Coverage and Test Stability
- Flaky-pattern check on `RunningJobSlotUsageBuilderTest`: no
`Thread.sleep`, no `Awaitility`/polling, no network or container dependency, no
shared static state, no order-sensitive assertions beyond the deterministic
`TreeMap`/`TreeSet` ordering already used in the production code, no
floating-point comparisons. Nothing in this test can contribute to CI flakiness.
- **Test stability rating: Stable**, backed by the above — no High-severity
flaky-test issue applies here.
- Coverage gap: as detailed in Issue 6, the untested paths are exactly the
failure modes, which is the wrong side of the line for a diagnostics API whose
value proposition is being trustworthy during incidents.
- Minor, non-blocking observation I hadn't called out before: every
`SlotProfile` constructed in the test uses the same literal `"test"` for
`sequence`, so the test's differentiation between "the same slot" and "a
different slot" relies entirely on `ownerJobId`/`slotId`/`worker`, never on
`sequence` itself. That's fine for what it's testing, but it means `sequence`'s
contribution to `SlotKey` uniqueness is exercised only implicitly.
## 2.3 Documentation Updates
- All four REST reference files updated:
`docs/en/engines/zeta/rest-api-v1.md`, `docs/en/engines/zeta/rest-api-v2.md`,
`docs/zh/engines/zeta/rest-api-v1.md`, `docs/zh/engines/zeta/rest-api-v2.md` —
re-diffed the English pair myself; both API versions and both languages are
covered, and the zh content is a real translation, not a placeholder copy.
- Example payload matches the code: `jobId` is a string
(`String.valueOf(jobId)`), `slotCount` an int, `pipelineSlotCounts` keyed by
pipeline id, `workerSlotCounts` keyed by worker address string. The v1 doc
correctly uses the `/hazelcast/rest/maps/` prefix.
- The explanatory notes are accurate for the normal path. Missing: the
meaning of `slotCount: 0` is ambiguous per Issue 1, the counting unit isn't
spelled out per Issue 7, and the per-request cost characteristic isn't
mentioned per Issue 4 — these are doc follow-ups to the code fixes, not
independent doc bugs.
# 3. Architectural Soundness
## 3.1 Elegance — **Precise fix with one debatable design choice**
The servlet -> service -> operation -> builder decomposition matches every
other Zeta REST endpoint's shape, and reusing `BaseServlet`/`BaseService` keeps
the new code small. Splitting the pure aggregation into a static builder with a
package-private overload is exactly what makes a mock-free, deterministic unit
test possible. The one choice I'd push back on is the two-source join itself:
intersecting a cluster-replicated `IMap` with a master-local in-memory snapshot
means the answer is only as trustworthy as the weaker source, and the weaker
source is empty in exactly the failover scenario the endpoint exists to help
diagnose (Issue 1). A single authoritative source, or an explicit "which
sources did we manage to read" field, would be a more honest design.
## 3.2 Maintainability
Good — small classes, clear names, no inheritance tricks, no shared mutable
state, easy to unit test in isolation. The main maintainability risk is
reputational: if operators learn this endpoint sometimes reports unexplained
zeros, the API loses its diagnostic value. Fixing Issue 1 plus the Issue 5 WARN
removes that risk.
## 3.3 Extensibility
A per-job/per-pipeline/per-worker breakdown is a solid foundation for a
future Web UI panel or slot-pressure alerting, and the response shape leaves
room for additional fields (e.g. a degraded flag) without breaking existing
callers. Once something actually polls it in production, 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 current `dev` (max `16`). The only
rolling-upgrade exposure is the new-node-calls-old-master case noted in 1.2,
which is confined to this new endpoint and breaks nothing pre-existing. No
`incompatible-changes.md` entry is required.
# 4. Issue Summary
| # | Issue | Location | Severity |
| --- | --- | --- | --- |
| 1 | `slotCount: 0` is indistinguishable from "couldn't read the slot
source", exactly during failover/cold-start |
`RunningJobSlotUsageBuilder.java:68-73, 446-448` | Medium |
| 2 | `getCoordinatorService()` called once per running job in the filter
lambda; can sleep up to 1.5s per call | `RunningJobSlotUsageBuilder.java:59-66`
| Medium |
| 3 | Aggregation runs inline on a shared Hazelcast generic operation
thread; sibling operation deliberately offloads |
`GetRunningJobSlotUsageOperation.java:92-96` | Medium |
| 4 | Full cluster-wide owned-slot map materialized per request, including
discarded jobs; no cache or rate limit |
`RunningJobSlotUsageBuilder.java:88-92, 423-426` | Low |
| 5 | No logging anywhere in the new path; mastership loss on the forwarding
branch is an unhandled 500 | `RunningJobSlotUsageService.java:35-43` | Low |
| 6 | No coverage of the degraded path, the public `build(SeaTunnelServer)`
entry point, serializer round-trip, or REST layer |
`RunningJobSlotUsageBuilderTest.java` | Low |
| 7 | `slotCount` counts task-group entries, not distinct slots; name
implies otherwise | `RunningJobSlotUsageBuilder.java:428-433, 503-509` | Low |
# 5. Merge Recommendation
### Conclusion: Ready to merge after fixes
**1. Blockers (must fix before merge)**
- Issue 1 — a diagnostics endpoint must not encode "I couldn't read the
data" as the same value it uses for "the answer is genuinely zero". Still
unaddressed since my 2026-08-03/2026-08-16 review of this same code.
- Issue 2 — hoist `getCoordinatorService()` out of the per-job lambda;
two-line fix, removes an O(N) x 1.5s blocking amplification.
- Issue 3 — offload the aggregation off the shared Hazelcast operation
thread, following the `GetRunningJobMetricsOperation` precedent in the same
package. Blocking shared operation threads is an engine-wide hazard, not an
endpoint-local inconvenience.
**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 when the slot source is unavailable.
- Issue 6 — add the two degraded-path unit cases and one REST-level IT
asserting real slot counts, following `PendingJobsRestIT`.
- Issue 7 — de-duplicate by `SlotKey`, or rename the field so the counting
unit is unambiguous.
**Overall assessment**
The feature is worth having and the plumbing is well done: routing is
correct on both API versions, the serializer id is free,
`SlotProfile.ownerJobID` really is on the wire so the cross-node intersection
works, 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 updated
accurately in both languages. Choosing `getInitializedResourceManager()` over
the lazy-initializing variant was the right instinct for a read-only path.
What's holding 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 silently return all zeros while blocking a shared
cluster thread for up to 1.5s per running job, with nothing logged anywhere to
explain why. Two of the three blockers are a handful of lines each (Issues 2
and 3); Issue 1 needs one deliberate decision (degraded flag vs.
coordinator-only fallback) plus a test for it. Once those land, with a
REST-level test proving real slot counts end to end, this is a solid, useful
addition to Zeta diagnostics.
**CI status at the time of this review:** `Build` is green (`SUCCESS`) at
this head — an improvement since 2026-08-16, when the apache-side check was red
due to an infrastructure-only `paimon-connector-it` Maven-wrapper `429 Too Many
Requests` failure in the fork run, unrelated to this diff. No CI action needed
right now beyond the source fixes above. `mergeStateStatus` is `BLOCKED`, which
on this repo means "awaiting maintainer review", not a merge conflict —
`mergeable` reports `MERGEABLE`.
--
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]