apache/seatunnel · warning

Error sink writer failed during shutdown

Error message

Error sink writer failed during shutdown

What it means

The DefaultErrorSinkWriter drain loop caught a Throwable from sink writing or lifecycle. Because the writer is already in the closed state, the failure happened during shutdown and is logged at WARN instead of ERROR. The Throwable is recorded as workerFailure and rethrown if it is an Error; data still in the queue at this point may not be written.

Solutions

  1. Inspect the logged underlying exception 'e' for the real sink-side cause and fix that (connectivity, schema, auth).
  2. Retry the job or re-send the error records that were pending at shutdown, since they may be lost.
  3. Gracefully cancel jobs so the drain loop empties the queue before close.
  4. Harden the error sink connector (timeouts, retries) so shutdown writes succeed.

Example fix

// before: hard-cancel drops drain during shutdown
seaTunnel.sh -can <jobId>  // immediate cancel while sink blocked
// after: stop gracefully / with timeout so drain completes
seaTunnel.sh -s <jobId>   // graceful stop, drain loop finishes first
Defensive patterns

Strategy: try-catch

Try / catch

try {
    errorSinkWriter.write(record);
} catch (Exception e) {
    log.error("error record dropped", e);
    deadLetterQueue.add(record); // preserve data
}

Prevention

When it happens

Trigger: drainLoop processes records or closes the underlying sinkWriter after close() has set closed=true, and an exception occurs (e.g. sink write fails, sinkWriter.close() fails, I/O error while draining remaining records).

Common situations: Job cancellation/stop races with an in-flight error-sink write; error sink connection dropped right at shutdown; sink close() throwing; disk/network failure during final drain.

Related errors


AI-assisted analysis of apache/seatunnel@cf67b549a7 (2026-09-10). Data as JSON: /api/errors/c4447721afe3e6c6. Report an issue: GitHub.

Appendix: source

Thrown at seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/task/error/DefaultErrorSinkWriter.java:533

                try {
                    synchronized (writerLock) {
                        sinkWriter.write(row);
                    }
                } catch (Throwable writeEx) {
                    workerFailure = writeEx;
                    throw writeEx;
                } finally {
                    pendingRows.decrementAndGet();
                }
            }
        } catch (Throwable e) {
            if (e instanceof Error) {
                workerFailure = e;
                throw (Error) e;
            }
            workerFailure = e;
            if (closed) {
                log.warn("Error sink writer failed during shutdown", e);
            } else {
                log.error("Error sink writer failed", e);
            }
        } finally {
            try {
                if (sinkWriter != null) {
                    sinkWriter.close();
                }
            } catch (Throwable closeEx) {
                if (closeEx instanceof Error) {
                    throw (Error) closeEx;
                }
                if (workerFailure != null) {
                    workerFailure.addSuppressed(closeEx);
                } else {
                    workerFailure = closeEx;
                }
                log.warn("Failed to close error sink writer", closeEx);

View on GitHub (pinned to cf67b549a7)