You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Flink Streaming JdbcSink异常处理:捕获BatchUpdateException并继续处理

方案1:在JdbcStatementBuilder中捕获单条记录异常

直接在构建JDBC语句的环节对单条记录的插入逻辑做try-catch,精准捕获单条记录的插入失败,记录日志后跳过该记录,不影响批次内其他正常记录的插入。

关键是关闭批量重试,避免因单条坏记录重试整个批次,导致正常记录重复处理。

示例代码:

import org.apache.flink.connector.jdbc.JdbcConnectionOptions;
import org.apache.flink.connector.jdbc.JdbcExecutionOptions;
import org.apache.flink.connector.jdbc.JdbcSink;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

// 日志实例
private static final Logger LOG = LoggerFactory.getLogger(YourFlinkJob.class);

// 构建JdbcSink
JdbcSink.sink(
    "INSERT INTO your_table (col1, col2) VALUES (?, ?)",
    (statement, event) -> {
        try {
            statement.setString(1, event.getCol1());
            statement.setInt(2, event.getCol2());
            statement.addBatch();
        } catch (SQLException e) {
            // 记录包含事件详情的错误日志
            LOG.error("Failed to insert record: {}", event, e);
            // 不抛出异常,跳过该记录的批量添加
        }
    },
    JdbcExecutionOptions.builder()
        .withBatchSize(500)
        .withBatchIntervalMs(200)
        .withMaxRetries(0) // 关闭批量重试
        .build(),
    new JdbcConnectionOptions.JdbcConnectionOptionsBuilder()
        .withUrl("jdbc:postgresql://host:port/db_name")
        .withDriverName("org.postgresql.Driver")
        .withUsername("db_user")
        .withPassword("db_pass")
        .build()
);

方案2:自定义SinkFunction实现批量异常处理

如果需要更灵活的错误处理(比如将坏记录写入死信队列),可以自定义SinkFunction,在批量执行阶段捕获BatchUpdateException,解析异常定位失败记录,单独处理后继续执行后续批次。

示例代码:

import org.apache.flink.configuration.Configuration;
import org.apache.flink.streaming.api.functions.sink.SinkFunction;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import java.sql.*;
import java.util.ArrayList;
import java.util.List;

public class CustomPostgresSink<T> implements SinkFunction<T> {
    private static final Logger LOG = LoggerFactory.getLogger(CustomPostgresSink.class);
    private transient Connection conn;
    private transient PreparedStatement stmt;
    private final String insertSql;
    private final JdbcConnectionOptions connOptions;
    private final int batchSize;
    private List<T> batchBuffer = new ArrayList<>();

    public CustomPostgresSink(String insertSql, JdbcConnectionOptions connOptions, int batchSize) {
        this.insertSql = insertSql;
        this.connOptions = connOptions;
        this.batchSize = batchSize;
    }

    @Override
    public void open(Configuration parameters) throws Exception {
        // 初始化JDBC连接
        conn = DriverManager.getConnection(
            connOptions.getUrl(),
            connOptions.getUsername(),
            connOptions.getPassword()
        );
        conn.setAutoCommit(false);
        stmt = conn.prepareStatement(insertSql);
    }

    @Override
    public void invoke(T event, Context context) throws Exception {
        batchBuffer.add(event);
        if (batchBuffer.size() >= batchSize) {
            flushBatch();
        }
    }

    private void flushBatch() throws Exception {
        if (batchBuffer.isEmpty()) return;

        try {
            // 填充批量参数
            for (T event : batchBuffer) {
                // 替换为你的事件参数设置逻辑
                stmt.setString(1, ((YourEvent) event).getCol1());
                stmt.setInt(2, ((YourEvent) event).getCol2());
                stmt.addBatch();
            }
            stmt.executeBatch();
            conn.commit();
            batchBuffer.clear();
        } catch (BatchUpdateException e) {
            conn.rollback();
            int[] updateCounts = e.getUpdateCounts();
            // 遍历批次,定位失败记录
            for (int i = 0; i < updateCounts.length; i++) {
                if (updateCounts[i] == Statement.EXECUTE_FAILED) {
                    T failedEvent = batchBuffer.get(i);
                    LOG.error("Batch insert failed for record: {}", failedEvent, e);
                    // 可选:将失败记录发送到死信队列
                    // deadLetterSink.invoke(failedEvent);
                }
            }
            batchBuffer.clear();
        } catch (SQLException e) {
            LOG.error("JDBC batch operation failed", e);
            conn.rollback();
            batchBuffer.clear();
        }
    }

    @Override
    public void close() throws Exception {
        // 关闭前刷写剩余批次
        flushBatch();
        if (stmt != null) stmt.close();
        if (conn != null) conn.close();
    }
}

使用自定义Sink:

stream.addSink(new CustomPostgresSink<>(
    "INSERT INTO your_table (col1, col2) VALUES (?, ?)",
    new JdbcConnectionOptions.JdbcConnectionOptionsBuilder()
        .withUrl("jdbc:postgresql://host:port/db_name")
        .withDriverName("org.postgresql.Driver")
        .withUsername("db_user")
        .withPassword("db_pass")
        .build(),
    500
));

方案3:结合侧输出流分离坏记录

如果需要单独收集坏记录做后续处理,可以利用Flink侧输出流,将失败记录发送到侧输出流,主流程继续处理正常数据。

步骤:

  1. 定义侧输出流标签:
private static final OutputTag<YourEvent> BAD_RECORDS_TAG = new OutputTag<YourEvent>("bad-records") {};
  1. 在异常捕获时发送侧输出:
// 在try-catch块中添加
context.output(BAD_RECORDS_TAG, failedEvent);
  1. 从主流获取侧输出流并处理:
DataStream<YourEvent> mainStream = ...;
DataStream<YourEvent> badRecordsStream = mainStream.getSideOutput(BAD_RECORDS_TAG);
// 比如将坏记录写入日志或存储
badRecordsStream.addSink(new LoggingSink<>("Bad record detected"));

内容的提问来源于stack exchange,提问作者Tany Fernandez

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.15 15:30:57