apache/seatunnel · warning

Collect worker count failed

Error message

Collect worker count failed: {}

What it means

PendingDiagnosticsCollector.collectClusterSnapshot queries the ResourceManager for the cluster's worker count to fill the diagnostics snapshot. Any exception (e.g. a dead worker or an RPC failure to a node) is swallowed and logged as 'Collect worker count failed', leaving the snapshot's worker count unset/absent rather than failing the whole diagnostic collection.

Solutions

  1. This log is a diagnostic degradation, not a job failure; check the included exception message for the failing worker address.
  2. Verify all workers are alive (`worker list` / engine REST diagnostics) and restart any dead nodes.
  3. Re-run the pending-diagnostics collection once the cluster is stable to get a complete snapshot.
  4. Check master-worker network connectivity/RPC timeouts if failures recur.

Example fix

// before: nothing to fix in config usually; restart dead workers
// after: confirm cluster health, then re-collect
bin/seatunnel.sh --list-worker   // ensure all expected nodes respond, then rerun diagnostics
Defensive patterns

Strategy: try-catch

Validate before calling

// Confirm all workers responsive before collecting diagnostics
workers.forEach(w -> requireAlive(w.getAddress(), rpcTimeoutMs));

Try / catch

try {
    int workerCount = resourceManager.workerCount(tags);
} catch (Exception e) {
    // treat as degraded diagnostics, not fatal; log and continue
    log.warn("worker count unavailable: {}", e.getMessage());
}

Prevention

When it happens

Trigger: collectClusterSnapshot -> resourceManager.workerCount(tags) throws, typically because a worker node is unreachable, its profile is stale, or the underlying RPC fails while gathering resource tags.

Common situations: A worker crashed or was decommissioned between slot queries and the worker-count query, network flakiness between master and workers, or diagnosing a cluster that is already unhealthy (which is exactly why diagnostics are running).

Related errors


AI-assisted analysis of apache/seatunnel@cf67b549a7 (2026-09-10). Data as JSON: /api/errors/57bf997a60790f3f. Report an issue: GitHub.

Appendix: source

Thrown at seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/diagnostic/PendingDiagnosticsCollector.java:268

            return snapshot;
        }
        Map<String, String> tags =
                tagFilter == null ? Collections.emptyMap() : new HashMap<>(tagFilter);
        List<SlotProfile> assignedSlots = Collections.emptyList();
        List<SlotProfile> unassignedSlots = Collections.emptyList();
        try {
            assignedSlots = resourceManager.getAssignedSlots(tags);
            unassignedSlots = resourceManager.getUnassignedSlots(tags);
        } catch (Exception e) {
            log.warn("Collect slots info failed: {}", ExceptionUtils.getMessage(e));
        }
        snapshot.setAssignedSlots(assignedSlots.size());
        snapshot.setFreeSlots(unassignedSlots.size());
        snapshot.setTotalSlots(assignedSlots.size() + unassignedSlots.size());
        try {
            snapshot.setWorkerCount(resourceManager.workerCount(tags));
        } catch (Exception e) {
            log.warn("Collect worker count failed: {}", ExceptionUtils.getMessage(e));
        }
        snapshot.setWorkers(collectWorkerResourceSnapshot(resourceManager).getWorkers());
        return snapshot;
    }

    /** Collects a point-in-time projection of the resources registered by the master. */
    public static WorkerResourceSnapshot collectWorkerResourceSnapshot(
            ResourceManager resourceManager) {
        WorkerResourceSnapshot snapshot = new WorkerResourceSnapshot();
        snapshot.setCollectedAt(System.currentTimeMillis());
        if (resourceManager == null) {
            return snapshot;
        }
        try {
            Map<Address, WorkerProfile> registerWorker = resourceManager.getRegisterWorker();
            if (registerWorker == null) {
                return snapshot;
            }

View on GitHub (pinned to cf67b549a7)