{"record":{"id":"acd4c599d86d1861","repo":"apache/seatunnel","slug":"source-close-failed","errorCode":null,"errorMessage":"source close failed","messagePattern":"source close failed","errorType":"console","errorClass":null,"httpStatus":null,"severity":"error","filePath":"seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/task/flow/SourceFlowLifeCycle.java","lineNumber":340,"sourceCode":"     * <p>Sets the {@code prepareClose} flag to {@code true} and sends a {@link\n     * SourceNoMoreElementOperation} to the remote enumerator, deregistering this reader from\n     * further split assignment.\n     *\n     * @throws RuntimeException if the deregistration message fails to send\n     */\n    public void signalNoMoreElement() {\n        // ready close this reader\n        try {\n            this.prepareClose = true;\n            runningTask\n                    .getExecutionContext()\n                    .sendToMember(\n                            new SourceNoMoreElementOperation(\n                                    currentTaskLocation, enumeratorTaskLocation),\n                            enumeratorTaskAddress)\n                    .get();\n        } catch (Exception e) {\n            log.warn(\"source close failed\", e);\n            throw new RuntimeException(e);\n        }\n    }\n\n    /**\n     * Registers this reader with the remote split enumerator.\n     *\n     * <p>Sends a {@link SourceRegisterOperation} to the enumerator at the previously resolved\n     * address, informing it that this reader subtask is ready to receive splits.\n     *\n     * @throws RuntimeException if registration fails due to communication errors\n     */\n    private void register() {\n        try {\n            runningTask\n                    .getExecutionContext()\n                    .sendToMember(\n                            new SourceRegisterOperation(","sourceCodeStart":322,"sourceCodeEnd":358,"githubUrl":"https://github.com/apache/seatunnel/blob/cf67b549a7a6c35fa0beb12d83c62892427ea919/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/task/flow/SourceFlowLifeCycle.java#L322-L358","documentation":"SourceFlowLifeCycle.signalNoMoreElement sends a SourceNoMoreElementOperation to the member hosting the split enumerator, signalling the source reader has finished. If the remote call fails, this message is logged at WARN and the exception is wrapped in a RuntimeException, which fails the source task. The enumerator may never learn the reader finished, stalling split completion for that source.","triggerScenarios":"The blocking .sendToMember(...).get() to enumeratorTaskAddress throws — e.g. the enumerator task's member died, network partition, or the future timed out — during source close/finish handling.","commonSituations":"Worker node crash or restart while a source reader finishes; network instability between Zeta nodes; source job teardown racing with enumerator task completion.","solutions":["Check whether the enumerator task/node is still alive in engine logs; if the node died, the job will be restarted by fault tolerance — verify job state.","Inspect connectivity between cluster members (ports, firewall) if this recurs on healthy nodes.","Retry/re-run the job; checkpoint restore should resume from the last completed checkpoint.","If caused by a task-cancel race, prefer graceful stop and upgrade if a fixed teardown race exists in a newer SeaTunnel version."],"exampleFix":"// before: node crashed, remote call hangs then fails\nsendToMember(new SourceNoMoreElementOperation(...), deadMemberAddress).get();\n// after: restore/reattach after node recovery, job resumes from checkpoint\nseaTunnel.sh -r <jobId> // resubmit with restore from last checkpoint","handlingStrategy":"retry","validationCode":null,"typeGuard":null,"tryCatchPattern":"try {\n    sourceJob.execute();\n} catch (RuntimeException e) {\n    if (e.getCause() != null && e.getCause().getMessage() != null\n            && e.getCause().getMessage().contains(\"source close failed\")) {\n        retryJobWithCheckpointRestore();\n    } else {\n        throw e;\n    }\n}","preventionTips":["Run the source and its enumerator on stable, connected nodes; avoid mixing with spot/preemptible workers.","Enable checkpointing so a failed finish signal is recoverable by restore.","Monitor inter-node connectivity; network partitions cause this.","Prefer graceful stop over node kill during source finish."],"tags":["source","split-enumerator","cluster-communication","task-failure"],"backgroundTag":"cluster-node-unreachable","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"}