apache/seatunnel · warning

Error sink worker thread did not terminate within

Error message

Error sink worker thread did not terminate within {} ms, interrupting it

What it means

close() waits for the error-sink worker thread to finish. If the thread is still alive after the join timeout, this warning is logged and the thread is interrupt()ed, followed by a shorter (max 5s) grace join. It means the error sink writer could not drain/finish within the allotted time, usually because the error-sink backend is slow or blocked.

Solutions

  1. Check error-sink backend health/latency at the time of the warning.
  2. Increase the error-sink writer close/drain timeout so slow backends can finish.
  3. Reduce error row volume (fix the upstream row errors) or increase error-sink parallelism/batch size.
  4. If the worker is deadlocked in a connector, capture a thread dump and fix or upgrade that connector.

Example fix

null
Defensive patterns

Strategy: retry

Validate before calling

null

Type guard

null

Try / catch

// if you observe this warning, capture a thread dump of the worker before the interrupt escalates
if (workerThread.isAlive()) {
    threadDump(); // diagnose where write/close is blocked
}

Prevention

When it happens

Trigger: workerThread.join(timeoutMillis) elapses with the thread alive — the worker is stuck in sinkWriter.write() or close() (e.g. blocking network I/O to the error sink, lock contention on writerLock, huge pending queue).

Common situations: Error-sink database latency spike or deadlock; error sink configured with very large flush intervals; error row rate far exceeds sink throughput so the queue never drains before close.

Related errors


AI-assisted analysis of apache/seatunnel@cf67b549a7 (2026-09-10). Data as JSON: /api/errors/bbca5d336db0c479. 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:617

        }
    }

    private void waitForWorkerTermination(long timeoutMillis) {
        if (workerThread == null) {
            return;
        }
        if (!workerThread.isAlive()) {
            return;
        }
        try {
            workerThread.join(timeoutMillis);
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
            log.warn("Interrupted while waiting for error sink worker to close");
        }

        if (workerThread.isAlive()) {
            log.warn(
                    "Error sink worker thread did not terminate within {} ms, interrupting it",
                    timeoutMillis);
            workerThread.interrupt();
            try {
                workerThread.join(Math.min(5_000L, timeoutMillis));
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
                log.warn(
                        "Interrupted while waiting for error sink worker to close after interrupt");
            }
        }
    }

    private void waitForPendingRows(long timeoutMillis) throws Exception {
        AtomicInteger currentPendingRows = this.pendingRows;
        if (currentPendingRows == null) {
            return;
        }

View on GitHub (pinned to cf67b549a7)