{"record":{"id":"b4d88349def875e6","repo":"apache/flink","slug":"recoverablewriter-has-been-closed","errorCode":null,"errorMessage":"RecoverableWriter has been closed","messagePattern":"RecoverableWriter has been closed","errorType":"validation","errorClass":"IllegalStateException","httpStatus":null,"severity":"error","filePath":"flink-filesystems/flink-s3-fs-native/src/main/java/org/apache/flink/fs/s3native/writer/NativeS3RecoverableWriter.java","lineNumber":246,"sourceCode":"        if (recoverable instanceof NativeS3Recoverable) {\n            return (NativeS3Recoverable) recoverable;\n        }\n        throw new IllegalArgumentException(\n                \"Native S3 File System cannot recover recoverable for other file system: \"\n                        + recoverable);\n    }\n\n    @Override\n    public void close() {\n        if (!closed.compareAndSet(false, true)) {\n            return;\n        }\n        LOG.debug(\"Closing S3 recoverable writer\");\n    }\n\n    private void checkNotClosed() {\n        if (closed.get()) {\n            throw new IllegalStateException(\"RecoverableWriter has been closed\");\n        }\n    }\n\n    public static NativeS3RecoverableWriter writer(\n            NativeS3ObjectOperations s3AccessHelper,\n            String localTmpDir,\n            long userDefinedMinPartSize,\n            int maxConcurrentUploadsPerStream) {\n\n        return new NativeS3RecoverableWriter(\n                s3AccessHelper, localTmpDir, userDefinedMinPartSize, maxConcurrentUploadsPerStream);\n    }\n}\n","sourceCodeStart":228,"sourceCodeEnd":260,"githubUrl":"https://github.com/apache/flink/blob/2f3c205e9266cb30240eb7f4fdab15cad629a70f/flink-filesystems/flink-s3-fs-native/src/main/java/org/apache/flink/fs/s3native/writer/NativeS3RecoverableWriter.java#L228-L260","documentation":"NativeS3RecoverableWriter uses an AtomicBoolean 'closed' flag; after close() flips it, any subsequent operation guarded by checkNotClosed() — open(), recover(), recoverForCommit() — throws IllegalStateException(\"RecoverableWriter has been closed\"). This enforces a deterministic shutdown contract instead of letting closed writers leak new uploads.","triggerScenarios":"Calling any guarded method on the writer after close(); typical sources: close() in a finally block that runs before a later recover attempt, lifecycle ordering bugs where cleanup precedes a final commit-for-recover, or shared writer instances closed by one operator while another still uses it.","commonSituations":"Sink operators caching the RecoverableWriter statically/instance-wide and one task's close() affecting another; a retry loop that closes the writer in a finally on first failure, then retries recovery on the same instance; unit tests reusing a writer across cases.","solutions":["Create a fresh writer (filesystem.createRecoverableWriter()) per lifecycle instead of reusing a closed one.","Ensure close() is the terminal action: audit that recover/commit calls cannot be ordered after cleanup in your operator.","Never share one RecoverableWriter instance across concurrent operators/tasks — each should own its own.","Catch IllegalStateException around recovery in shutdown hooks to distinguish 'already closed' from genuine S3 failures."],"exampleFix":"// before\nstatic RecoverableWriter WRITER = fs.createRecoverableWriter();\n// one task closes WRITER; others still call WRITER.open()\n\n// after\n// per-operator instance, closed only in that operator's close()\nprivate transient RecoverableWriter writer;\npublic void open(...) { writer = fs.createRecoverableWriter(); }\npublic void close() { if (writer != null) writer.close(); writer = null; }","handlingStrategy":"validation","validationCode":"// track writer lifecycle alongside the operator that owns it\nprivate transient RecoverableWriter writer; // never static/shared\n// open(): writer = fs.createRecoverableWriter();\n// close(): writer.close(); writer = null;","typeGuard":null,"tryCatchPattern":"try {\n    writer.recover(rec);\n} catch (IllegalStateException e) {\n    if (e.getMessage().equals(\"RecoverableWriter has been closed\")) {\n        // recreate writer and retry once\n    }\n}","preventionTips":["One writer per operator lifecycle; recreate after close.","Make close() the terminal action in operator lifecycle.","Never share writers across tasks."],"tags":["s3","lifecycle","writer","state","io"],"backgroundTag":null,"analyzedSha":"2f3c205e9266cb30240eb7f4fdab15cad629a70f","analyzedAt":"2026-08-14T08:48:24.518Z","schemaVersion":2},"datasetVersion":"2026-08-14T10:17:34.591Z"}