{"record":{"id":"57bf997a60790f3f","repo":"apache/seatunnel","slug":"collect-worker-count-failed","errorCode":null,"errorMessage":"Collect worker count failed: {}","messagePattern":"Collect worker count failed: (.+?)","errorType":"console","errorClass":null,"httpStatus":null,"severity":"warning","filePath":"seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/diagnostic/PendingDiagnosticsCollector.java","lineNumber":268,"sourceCode":"            return snapshot;\n        }\n        Map<String, String> tags =\n                tagFilter == null ? Collections.emptyMap() : new HashMap<>(tagFilter);\n        List<SlotProfile> assignedSlots = Collections.emptyList();\n        List<SlotProfile> unassignedSlots = Collections.emptyList();\n        try {\n            assignedSlots = resourceManager.getAssignedSlots(tags);\n            unassignedSlots = resourceManager.getUnassignedSlots(tags);\n        } catch (Exception e) {\n            log.warn(\"Collect slots info failed: {}\", ExceptionUtils.getMessage(e));\n        }\n        snapshot.setAssignedSlots(assignedSlots.size());\n        snapshot.setFreeSlots(unassignedSlots.size());\n        snapshot.setTotalSlots(assignedSlots.size() + unassignedSlots.size());\n        try {\n            snapshot.setWorkerCount(resourceManager.workerCount(tags));\n        } catch (Exception e) {\n            log.warn(\"Collect worker count failed: {}\", ExceptionUtils.getMessage(e));\n        }\n        snapshot.setWorkers(collectWorkerResourceSnapshot(resourceManager).getWorkers());\n        return snapshot;\n    }\n\n    /** Collects a point-in-time projection of the resources registered by the master. */\n    public static WorkerResourceSnapshot collectWorkerResourceSnapshot(\n            ResourceManager resourceManager) {\n        WorkerResourceSnapshot snapshot = new WorkerResourceSnapshot();\n        snapshot.setCollectedAt(System.currentTimeMillis());\n        if (resourceManager == null) {\n            return snapshot;\n        }\n        try {\n            Map<Address, WorkerProfile> registerWorker = resourceManager.getRegisterWorker();\n            if (registerWorker == null) {\n                return snapshot;\n            }","sourceCodeStart":250,"sourceCodeEnd":286,"githubUrl":"https://github.com/apache/seatunnel/blob/cf67b549a7a6c35fa0beb12d83c62892427ea919/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/diagnostic/PendingDiagnosticsCollector.java#L250-L286","documentation":"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.","triggerScenarios":"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.","commonSituations":"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).","solutions":["This log is a diagnostic degradation, not a job failure; check the included exception message for the failing worker address.","Verify all workers are alive (`worker list` / engine REST diagnostics) and restart any dead nodes.","Re-run the pending-diagnostics collection once the cluster is stable to get a complete snapshot.","Check master-worker network connectivity/RPC timeouts if failures recur."],"exampleFix":"// before: nothing to fix in config usually; restart dead workers\n// after: confirm cluster health, then re-collect\nbin/seatunnel.sh --list-worker   // ensure all expected nodes respond, then rerun diagnostics","handlingStrategy":"try-catch","validationCode":"// Confirm all workers responsive before collecting diagnostics\nworkers.forEach(w -> requireAlive(w.getAddress(), rpcTimeoutMs));","typeGuard":null,"tryCatchPattern":"try {\n    int workerCount = resourceManager.workerCount(tags);\n} catch (Exception e) {\n    // treat as degraded diagnostics, not fatal; log and continue\n    log.warn(\"worker count unavailable: {}\", e.getMessage());\n}","preventionTips":["Keep workers healthy before running pending-job diagnostics.","Check master-worker RPC connectivity when snapshots repeatedly lack worker counts.","Re-run diagnostics after cluster stabilization to fill gaps.","Alert on dead workers, which is the usual root cause."],"tags":["diagnostics","cluster","resource-manager","rpc"],"backgroundTag":"upstream-api-error","analyzedSha":"cf67b549a7a6c35fa0beb12d83c62892427ea919","analyzedAt":"2026-09-10T21:44:55.265Z","contentChangedAt":"2026-09-10T21:44:55.265Z","schemaVersion":2},"datasetVersion":"2026-09-23T08:17:48.524Z"}