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

Java中使用Apache Beam JDBCIO插入数据库表时如何提取错误记录

实现插入失败记录捕获的方案

你当前使用的默认JdbcIO.write()会在单条记录写入异常时直接抛出异常终止整个Pipeline,无法单独提取失败记录。可以通过以下两种方案实现失败记录的收集:


方案1:使用Beam官方内置死信队列(推荐,适用于Beam 2.30及以上版本)

新版本JdbcIO提供了原生的死信输出能力,不需要自行封装写入逻辑,改造代码如下:

import org.apache.beam.sdk.io.jdbc.JdbcIO;
import org.apache.beam.sdk.values.PCollection;

PipelineOptions options = PipelineOptionsFactory.create();
options.setRunner(FlinkRunner.class);
Pipeline p = Pipeline.create(options);

// 模拟数据部分不变
Collection<Stock> stockList = Arrays.asList(
    new Stock("AAP", 2000,"Apple Inc"),
    new Stock("MSF", 3000, "Microsoft Corporation"),
    new Stock("NVDA", 4000, "NVIDIA Corporation"),
    new Stock("INT", 3200, "Intel Corporation")
);
PCollection<Stock> data = p.apply(Create.of(stockList)
    .withCoder(SerializableCoder.of(Stock.class)));

// 改造写入逻辑,用writeWriteRows替代原有write方法,开启死信输出
JdbcIO.WriteResult<Stock> writeResult = data.apply(JdbcIO.<Stock>writeWriteRows()
    .withDataSourceConfiguration(JdbcIO.DataSourceConfiguration
        .create("org.postgresql.Driver","jdbc:postgresql://localhost:5432/postgres")
        .withUsername("postgres").withPassword("sachin"))
    .withStatement("insert into stocks values(?, ?, ?)")
    .withPreparedStatementSetter(new JdbcIO.PreparedStatementSetter<Stock>() {
        private static final long serialVersionUID = 1L;
        public void setParameters(Stock element, PreparedStatement query) throws SQLException {
            query.setString(1, element.getSymbol());
            query.setLong(2, element.getPrice());
            query.setString(3, element.getCompany());
        }
    })
    // 开启死信队列,收集写入失败的记录
    .withDeadLetterOutput()
);

// 获取所有插入失败的记录,每个元素包含原始Stock对象、报错异常、执行的SQL等信息
PCollection<JdbcIO.WriteFailure<Stock>> failedRecords = writeResult.getFailedRecords();

// 可以对失败记录做后续处理,比如打印、写入本地文件、存入异常表等
// 示例:打印失败记录的符号和报错信息
failedRecords.apply(ParDo.of(new DoFn<JdbcIO.WriteFailure<Stock>, Void>() {
    @ProcessElement
    public void processElement(@Element JdbcIO.WriteFailure<Stock> failure) {
        System.out.println("插入失败股票:" + failure.getElement().getSymbol() 
            + ",错误原因:" + failure.getException().getMessage());
    }
}));

p.run().waitUntilFinish();

方案2:自定义ParDo封装写入逻辑(适用于低版本Beam)

如果你的Beam版本不支持原生死信队列,可以通过TupleTag分流成功/失败数据,自行捕获写入异常:

  1. 首先定义两个TupleTag用于数据分流
// 写入成功的输出Tag,不需要可省略
TupleTag<Void> successTag = new TupleTag<Void>(){};
// 写入失败的输出Tag,存储原始数据和异常信息
TupleTag<KV<Stock, Exception>> failedTag = new TupleTag<KV<Stock, Exception>>(){};
  1. 替换原有写入逻辑为自定义ParDo,在内部处理数据库写入和异常捕获
PCollectionTuple resultTuple = data.apply(ParDo.of(new DoFn<Stock, Void>() {
    private Connection conn;
    private PreparedStatement stmt;

    @Setup
    public void setup() throws SQLException {
        // 初始化数据库连接,每个DoFn实例复用一个连接
        conn = DriverManager.getConnection(
            "jdbc:postgresql://localhost:5432/postgres", 
            "postgres", 
            "sachin"
        );
        stmt = conn.prepareStatement("insert into stocks values(?, ?, ?)");
    }

    @ProcessElement
    public void processElement(@Element Stock element, MultiOutputReceiver out) {
        try {
            stmt.setString(1, element.getSymbol());
            stmt.setLong(2, element.getPrice());
            stmt.setString(3, element.getCompany());
            stmt.executeUpdate();
            out.get(successTag).output(null);
        } catch (SQLException e) {
            // 捕获写入异常,输出到失败Tag
            out.get(failedTag).output(KV.of(element, e));
        }
    }

    @Teardown
    public void teardown() throws SQLException {
        // 销毁时关闭连接
        if (stmt != null) stmt.close();
        if (conn != null) conn.close();
    }
}).withOutputTags(successTag, TupleTagList.of(failedTag)));

// 提取失败记录
PCollection<KV<Stock, Exception>> failedRecords = resultTuple.get(failedTag);

内容的提问来源于stack exchange,提问作者Sachin Rane

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.25 00:54:04