{"record":{"id":"c52525d3a0f92ac9","repo":"apache/seatunnel","slug":"flush-errorhandler-for-transform-stage-failed","errorCode":null,"errorMessage":"Flush ErrorHandler for transform stage failed","messagePattern":"Flush ErrorHandler for transform stage failed","errorType":"exception","errorClass":"RuntimeException","httpStatus":null,"severity":"error","filePath":"seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/task/flow/TransformFlowLifeCycle.java","lineNumber":391,"sourceCode":"        super.close();\n    }\n\n    private void flushErrorHandler() {\n        flushErrorHandler(null);\n    }\n\n    private void flushErrorHandler(Long checkpointId) {\n        if (errorHandler == null) {\n            return;\n        }\n        try {\n            if (checkpointId == null) {\n                errorHandler.flush();\n            } else {\n                errorHandler.flush(checkpointId);\n            }\n        } catch (Exception e) {\n            throw new RuntimeException(\"Flush ErrorHandler for transform stage failed\", e);\n        }\n    }\n\n    private void snapshotErrorHandler(long checkpointId) {\n        if (errorHandler != null) {\n            errorHandler.snapshotState(checkpointId);\n        }\n    }\n\n    @Override\n    public void notifyCheckpointComplete(long checkpointId) {\n        if (errorHandler != null) {\n            errorHandler.notifyCheckpointComplete(checkpointId);\n        }\n    }\n\n    @Override\n    public void notifyCheckpointAborted(long checkpointId) {","sourceCodeStart":373,"sourceCodeEnd":409,"githubUrl":"https://github.com/apache/seatunnel/blob/cf67b549a7a6c35fa0beb12d83c62892427ea919/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/task/flow/TransformFlowLifeCycle.java#L373-L409","documentation":"flushErrorHandler() flushes buffered error records of the transform stage's ErrorHandler, either globally (no checkpoint) or up to a checkpointId. Any exception from flush is wrapped in a RuntimeException with this message, failing the calling task flow (received() or flush paths).","triggerScenarios":"flushErrorHandler() called from received() (per-record) or checkpoint flush; underlying errorHandler.flush()/flush(checkpointId) throws because the error sink is unavailable, serialization fails, or the writer is closed.","commonSituations":"Error-output sink outage while error records are being written; error sink storage full; serialization issue in error payload; sink closed prematurely during concurrent flush and close.","solutions":["Inspect the cause for the actual sink error and restore the error-output target (network, credentials, disk space).","Confirm the error sink config (type, URL/path, auth) is valid in the job config.","Avoid closing the transform while flush is in flight (ordering of close vs barrier flush).","If error volume is huge, throttle error output or sample to avoid sink overload."],"exampleFix":"// before: raw flush may fail the whole flow on transient sink errors\nerrorHandler.flush(checkpointId);\n// after: add bounded retry for transient sink errors\nfor (int i = 0; i < 3; i++) {\n    try { errorHandler.flush(checkpointId); return; }\n    catch (Exception retryEx) { sleep(1000L * (i + 1)); }\n}\nthrow new RuntimeException(\"Flush ErrorHandler for transform stage failed\", lastEx);","handlingStrategy":"retry","validationCode":null,"typeGuard":null,"tryCatchPattern":"try { flushErrorHandler(checkpointId); } catch (RuntimeException e) { if (e.getMessage().startsWith(\"Flush ErrorHandler\")) { retryFlushWithBackoff(checkpointId); } else { throw e; } }","preventionTips":["Health-check the error sink periodically during the job","Bound error-record volume to avoid sink overload","Ensure close ordering: flush before close, never concurrently"],"tags":["io","error-handling","flush","transform"],"backgroundTag":"file-write-failed","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"}