{"record":{"id":"f1a6264152350236","repo":"apache/seatunnel","slug":"errorsinkrowwriter-is-already-closed","errorCode":null,"errorMessage":"ErrorSinkRowWriter is already closed","messagePattern":"ErrorSinkRowWriter is already closed","errorType":"exception","errorClass":"IllegalStateException","httpStatus":null,"severity":"error","filePath":"seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/task/error/SynchronizedErrorSinkRowWriter.java","lineNumber":39,"sourceCode":"\n/** Thread-safe wrapper for ErrorSinkRowWriter to serialize write/close calls. */\npublic final class SynchronizedErrorSinkRowWriter<T> implements ErrorSinkRowWriter<T> {\n\n    private static final long serialVersionUID = 1L;\n\n    private final ErrorSinkRowWriter<T> delegate;\n    private final Object lock = new Object();\n    private volatile boolean closed;\n\n    public SynchronizedErrorSinkRowWriter(ErrorSinkRowWriter<T> delegate) {\n        this.delegate = Objects.requireNonNull(delegate, \"delegate must not be null\");\n    }\n\n    @Override\n    public void write(RowErrorContext ctx, T row, Throwable t) throws Exception {\n        synchronized (lock) {\n            if (closed) {\n                throw new IllegalStateException(\"ErrorSinkRowWriter is already closed\");\n            }\n            delegate.write(ctx, row, t);\n        }\n    }\n\n    @Override\n    public boolean writeAndCheckAccepted(RowErrorContext ctx, T row, Throwable t) throws Exception {\n        synchronized (lock) {\n            if (closed) {\n                throw new IllegalStateException(\"ErrorSinkRowWriter is already closed\");\n            }\n            return delegate.writeAndCheckAccepted(ctx, row, t);\n        }\n    }\n\n    @Override\n    public void flush() throws Exception {\n        synchronized (lock) {","sourceCodeStart":21,"sourceCodeEnd":57,"githubUrl":"https://github.com/apache/seatunnel/blob/cf67b549a7a6c35fa0beb12d83c62892427ea919/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/task/error/SynchronizedErrorSinkRowWriter.java#L21-L57","documentation":"SynchronizedErrorSinkRowWriter.write delegates error-row writes to an underlying ErrorSinkRowWriter under a lock. If the writer has already been closed (typically after checkpoint teardown or sink close), any further write throws IllegalStateException 'ErrorSinkRowWriter is already closed'. This protects the closed resource from late writes.","triggerScenarios":"Calling write(ctx, row, t) after close() was invoked on the writer, e.g. an upstream error path still emitting rows while the task is shutting down.","commonSituations":"Late error events arriving during task cancellation or checkpoint completion; asynchronous error handlers not stopped before the writer is closed.","solutions":["Stop feeding error rows before closing the writer; close only after all error paths have finished.","Check the closed state (or catch IllegalStateException) in late-emit paths and drop/queue the row instead.","Ensure task lifecycle ordering: close error consumers/writers as the last step of task teardown."],"exampleFix":"// before\nwriter.close();\nwriter.write(ctx, row, t); // throws\n// after\nwriter.write(ctx, row, t);\nwriter.close();","handlingStrategy":"try-catch","validationCode":"null","typeGuard":"null","tryCatchPattern":"try {\n    writer.write(ctx, row, t);\n} catch (IllegalStateException e) {\n    if (e.getMessage().contains(\"already closed\")) {\n        log.warn(\"Dropping error row after writer close\");\n    } else throw e;\n}","preventionTips":["Close the error writer only after all error paths are stopped","Join/stop async error handlers before close","Order teardown: stop emitting, flush, then close"],"tags":["lifecycle","state","error-handler"],"backgroundTag":"invalid-state-transition","analyzedSha":"cf67b549a7a6c35fa0beb12d83c62892427ea919","analyzedAt":"2026-09-10T21:44:55.265Z","contentChangedAt":"2026-09-10T21:44:55.265Z","schemaVersion":2},"datasetVersion":"2026-09-14T05:17:10.506Z"}