apache/seatunnel · critical · RowErrorHandlingFatalException

Error sink failed for stage

Error message

Error sink failed for stage [%s], plugin [%s]

What it means

Thrown as RowErrorHandlingFatalException from ErrorHandler.onError when the error-sink writer itself fails while trying to persist a bad row. Since the error sink is the last line of defense for row-level errors, its failure fails the whole job. The message includes the stage and plugin name to locate the failing pipeline component.

Solutions

  1. Check the log's 'Error sink failed for stage [..], plugin [..]' entry for the root sinkEx cause.
  2. Fix or restore the error-sink target and resubmit the job.
  3. Temporarily switch error-handler mode from ROUTE to IGNORE/FAIL to unblock the job while the sink is repaired.
  4. Validate the error-sink plugin config (plugin_name, credentials, schema) before enabling ROUTE mode.

Example fix

// before (broken sink config)
env { error-handler { mode = ROUTE } }
// after
env { error-handler { mode = ROUTE
  sink { plugin_name = "Console" } } }
Defensive patterns

Strategy: validation

Validate before calling

// before submit: verify error-sink config when mode=ROUTE
if (mode == "ROUTE" && !errorHandlerConfig.hasPath("sink.plugin_name")) throw new ConfigValidationError("ROUTE requires sink.plugin_name");

Try / catch

try { runJob(); } catch (RowErrorHandlingFatalException e) { if (e.getMessage().startsWith("Error sink failed")) { restoreSinkTarget(); resubmit(); } }

Prevention

When it happens

Trigger: A row-level error occurs with error-handler mode ROUTE; the error sink write (or its initialization) throws; ErrorHandler.onError catches the sink exception and rethrows fatal.

Common situations: Error-sink target unavailable (DB down, file unwritable) at the moment a bad row arrives; error-sink schema mismatch with captured error record; queue overflow combined with sink write failure.

Related errors


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

Appendix: source

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

        if (config.getMode() == ErrorHandlerMode.ROUTE && errorSinkWriter != null) {
            try {
                log.debug(
                        "Writing error row to sink. stage={}, plugin={}, tableId={}",
                        ctx.getStage(),
                        ctx.getPluginName(),
                        ctx.getTableId());
                boolean accepted = errorSinkWriter.writeAndCheckAccepted(ctx, row, t);
                result =
                        accepted
                                ? ErrorHandleResult.ROUTED_TO_ERROR_SINK
                                : ErrorHandleResult.DROPPED;
            } catch (Exception sinkEx) {
                log.error(
                        "Error sink failed for stage [{}], plugin [{}], failing the job",
                        ctx.getStage(),
                        ctx.getPluginName(),
                        sinkEx);
                throw new RowErrorHandlingFatalException(
                        String.format(
                                "Error sink failed for stage [%s], plugin [%s]",
                                ctx.getStage(), ctx.getPluginName()),
                        sinkEx);
            }
        }

        maybeThrowOnThreshold(ctx, currentErrorCount);
        return result;
    }

    private void maybeThrowOnThreshold(RowErrorContext ctx, long currentErrorCount) {
        if (config.getMaxErrorRecords() > 0 && currentErrorCount > config.getMaxErrorRecords()) {
            throw new RowErrorHandlingFatalException(
                    String.format(
                            "Too many row-level errors in stage [%s], plugin [%s]: %d records exceeded max_error_records=%d",
                            stageName(ctx),
                            pluginName(ctx),

View on GitHub (pinned to cf67b549a7)