apache/flink · error · IllegalStateException
RecoverableWriter has been closed
Error message
RecoverableWriter has been closed
What it means
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.
Source
Thrown at flink-filesystems/flink-s3-fs-native/src/main/java/org/apache/flink/fs/s3native/writer/NativeS3RecoverableWriter.java:246
if (recoverable instanceof NativeS3Recoverable) {
return (NativeS3Recoverable) recoverable;
}
throw new IllegalArgumentException(
"Native S3 File System cannot recover recoverable for other file system: "
+ recoverable);
}
@Override
public void close() {
if (!closed.compareAndSet(false, true)) {
return;
}
LOG.debug("Closing S3 recoverable writer");
}
private void checkNotClosed() {
if (closed.get()) {
throw new IllegalStateException("RecoverableWriter has been closed");
}
}
public static NativeS3RecoverableWriter writer(
NativeS3ObjectOperations s3AccessHelper,
String localTmpDir,
long userDefinedMinPartSize,
int maxConcurrentUploadsPerStream) {
return new NativeS3RecoverableWriter(
s3AccessHelper, localTmpDir, userDefinedMinPartSize, maxConcurrentUploadsPerStream);
}
}
View on GitHub (pinned to 2f3c205e92)
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.
Example fix
// before
static RecoverableWriter WRITER = fs.createRecoverableWriter();
// one task closes WRITER; others still call WRITER.open()
// after
// per-operator instance, closed only in that operator's close()
private transient RecoverableWriter writer;
public void open(...) { writer = fs.createRecoverableWriter(); }
public void close() { if (writer != null) writer.close(); writer = null; } Defensive patterns
Strategy: validation
Validate before calling
// track writer lifecycle alongside the operator that owns it private transient RecoverableWriter writer; // never static/shared // open(): writer = fs.createRecoverableWriter(); // close(): writer.close(); writer = null;
Try / catch
try {
writer.recover(rec);
} catch (IllegalStateException e) {
if (e.getMessage().equals("RecoverableWriter has been closed")) {
// recreate writer and retry once
}
} Prevention
- One writer per operator lifecycle; recreate after close.
- Make close() the terminal action in operator lifecycle.
- Never share writers across tasks.
When it happens
Trigger: 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.
Common situations: 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.
Related errors
- Stream is closed
- Stream is already closed
- Corrupt data, magic number mismatch. Expected %8x, found %8x
- Corrupt data, magic number mismatch. Expected %8x, found %8x
- Unrecognized version: {}
AI-assisted analysis of apache/flink@2f3c205e92 (2026-08-14).
Data as JSON: /api/errors/b4d88349def875e6.
Report an issue: GitHub.