Flink Streaming JdbcSink异常处理:捕获BatchUpdateException并继续处理
处理Flink JdbcSink插入Postgres的异常记录
方案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侧输出流,将失败记录发送到侧输出流,主流程继续处理正常数据。
步骤:
- 定义侧输出流标签:
private static final OutputTag<YourEvent> BAD_RECORDS_TAG = new OutputTag<YourEvent>("bad-records") {};
- 在异常捕获时发送侧输出:
// 在try-catch块中添加 context.output(BAD_RECORDS_TAG, failedEvent);
- 从主流获取侧输出流并处理:
DataStream<YourEvent> mainStream = ...; DataStream<YourEvent> badRecordsStream = mainStream.getSideOutput(BAD_RECORDS_TAG); // 比如将坏记录写入日志或存储 badRecordsStream.addSink(new LoggingSink<>("Bad record detected"));
内容的提问来源于stack exchange,提问作者Tany Fernandez
相关产品推荐
相关产品推荐

