{"record":{"id":"463f875be648757a","repo":"apache/seatunnel","slug":"collect-worker-resource-snapshot-failed","errorCode":null,"errorMessage":"Collect worker resource snapshot failed: {}","messagePattern":"Collect worker resource snapshot 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":296,"sourceCode":"        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            }\n            List<WorkerResourceDiagnostic> workers =\n                    registerWorker.entrySet().stream()\n                            .map(entry -> convertWorker(entry.getKey(), entry.getValue()))\n                            .filter(worker -> worker != null)\n                            .sorted(Comparator.comparing(WorkerResourceDiagnostic::getAddress))\n                            .collect(Collectors.toList());\n            snapshot.setAvailable(true);\n            snapshot.setWorkers(workers);\n        } catch (Exception e) {\n            log.warn(\"Collect worker resource snapshot failed: {}\", ExceptionUtils.getMessage(e));\n        }\n        return snapshot;\n    }\n\n    private static WorkerResourceDiagnostic convertWorker(\n            Address registeredAddress, WorkerProfile workerProfile) {\n        if (workerProfile == null) {\n            return null;\n        }\n        WorkerResourceDiagnostic diagnostic = new WorkerResourceDiagnostic();\n        Address address = workerProfile.getAddress();\n        if (address == null) {\n            address = registeredAddress;\n        }\n        diagnostic.setAddress(address == null ? \"UNKNOWN\" : address.toString());\n        if (workerProfile.getAttributes() != null) {\n            diagnostic.setTags(new HashMap<>(workerProfile.getAttributes()));\n        } else {","sourceCodeStart":278,"sourceCodeEnd":314,"githubUrl":"https://github.com/apache/seatunnel/blob/cf67b549a7a6c35fa0beb12d83c62892427ea919/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/diagnostic/PendingDiagnosticsCollector.java#L278-L314","documentation":"PendingDiagnosticsCollector.collectWorkerResourceSnapshot builds a per-worker resource diagnostic projection from the master's ResourceManager. On exception, it logs 'Collect worker resource snapshot failed', marks the snapshot unavailable, and returns it so the collector can still emit a partial cluster snapshot.","triggerScenarios":"Called from collectClusterSnapshot; fails when iterating worker profiles throws — a worker disconnected mid-iteration, its WorkerProfile is null/incomplete, or an RPC to fetch worker resource info fails.","commonSituations":"Flapping workers during scale-up/down, network partitions isolating a worker from the master, or running diagnostics on an already-degraded cluster where some workers cannot answer resource queries.","solutions":["Check the exception message to identify which worker failed to report.","Confirm the worker process is running and reachable; restart it if dead.","Re-run diagnostics after the cluster stabilizes to capture the full worker snapshot.","If only some workers appear in snapshots, check master-worker RPC timeouts and network stability."],"exampleFix":"// before: diagnose with missing worker data\n// after: restore worker then re-collect\nbin/seatunnel.sh --shutdown-worker <address> ; bin/seatunnel.sh --start-worker  // then rerun pending diagnostics","handlingStrategy":"try-catch","validationCode":"// Filter null/incomplete worker profiles before building the snapshot\nList<WorkerProfile> profiles = workers.stream().filter(Objects::nonNull).collect(Collectors.toList());\nif (profiles.isEmpty()) { log.warn(\"no worker profiles available; snapshot will be marked unavailable\"); }","typeGuard":null,"tryCatchPattern":"try {\n    snapshot = collector.collectWorkerResourceSnapshot(resourceManager);\n    if (!snapshot.isAvailable()) {\n        // re-run later or fall back to last good snapshot\n    }\n} catch (Exception e) {\n    log.warn(\"diagnostic collection degraded: {}\", e.getMessage());\n}","preventionTips":["Monitor worker liveness so resource queries don't hit dead nodes.","Avoid running diagnostics during worker scale events.","Treat snapshot.available=false as a signal to re-collect later.","Increase master-worker RPC timeout if profiles are large or networks are slow."],"tags":["diagnostics","workers","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"}