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

  1. Create a fresh writer (filesystem.createRecoverableWriter()) per lifecycle instead of reusing a closed one.
  2. Ensure close() is the terminal action: audit that recover/commit calls cannot be ordered after cleanup in your operator.
  3. Never share one RecoverableWriter instance across concurrent operators/tasks — each should own its own.
  4. 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

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


AI-assisted analysis of apache/flink@2f3c205e92 (2026-08-14). Data as JSON: /api/errors/b4d88349def875e6. Report an issue: GitHub.