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

Flink技术问询:如何修改JDBCOutputFormat写入Row列表至数据库

解决Flink单条输入生成多条记录写入PostgreSQL的问题

你原来的代码已经能正常把每行数据转成单条记录写入PostgreSQL了,现在要改成单条输入生成两条记录,并且数据流变成了DataStream<List<Row>>,其实核心不需要修改JDBCOutputFormat的配置,而是需要先把List<Row>的数据流展开成单个Row的数据流,再用原来的输出格式写入。当然如果想做批量写入优化,也可以调整成批量模式的Sink,下面分两种情况说明:

方案一:展开List为单Row流(兼容原有JDBCOutputFormat)

因为JDBCOutputFormat是针对单个Row设计的,直接给它List<Row>的数据流会报错。所以先通过flatMap把每个List里的Row都拆出来,得到普通的DataStream<Row>,之后就能复用原来的JDBCOutputFormat配置了。

修改后的完整代码大概是这样:

String strQuery = "INSERT INTO public.alarm (id, name, marks) VALUES (?, ?, ?)";
JDBCOutputFormat jdbcOutput = JDBCOutputFormat.buildJDBCOutputFormat()
        .setDrivername("org.postgresql.Driver")
        .setDBUrl("jdbc:postgresql://localhost:5432/postgres?user=michel&password=polnareff")
        .setQuery(strQuery)
        .setSqlTypes(new int[] { Types.INTEGER, Types.VARCHAR, Types.INTEGER})
        .finish();

// 生成包含两条Row的List
DataStream<List<Row>> rows = FilterStream
        .map((tuple)-> {
            List<Row> rowList = new ArrayList<>();
            // 第一条记录:name后缀加-1
            Row row1 = new Row(3);
            row1.setField(0, tuple.f0);
            row1.setField(1, tuple.f1 + "-1");
            row1.setField(2, tuple.f2);
            rowList.add(row1);
            // 第二条记录:name后缀加-2
            Row row2 = new Row(3);
            row2.setField(0, tuple.f0);
            row2.setField(1, tuple.f1 + "-2");
            row2.setField(2, tuple.f2);
            rowList.add(row2);
            return rowList;
        });

// 关键步骤:将List<Row>展开为单个Row的数据流
DataStream<Row> flatRows = rows.flatMap((List<Row> rowList, Collector<Row> out) -> {
    for (Row row : rowList) {
        out.collect(row);
    }
});

// 继续用原来的JDBCOutputFormat写入
flatRows.writeUsingOutputFormat(jdbcOutput);
env.execute();

方案二:使用批量写入的JdbcSink(新版本Flink推荐)

如果你的Flink版本是1.11及以上,更推荐使用JdbcSink的批量插入模式,这样可以直接处理批量数据,减少数据库连接的开销,效率更高。这种情况下就不需要用JDBCOutputFormat了,换成JdbcSink:

// 定义批量插入的语句
String insertQuery = "INSERT INTO public.alarm (id, name, marks) VALUES (?, ?, ?)";

// 配置JdbcSink的批量模式
JdbcSink<Row> jdbcSink = JdbcSink.sink(
        insertQuery,
        (ps, row) -> {
            // 为每条Row设置SQL参数
            ps.setInt(1, (Integer) row.getField(0));
            ps.setString(2, (String) row.getField(1));
            ps.setInt(3, (Integer) row.getField(2));
        },
        JdbcExecutionOptions.builder()
                .withBatchSize(2) // 匹配单输入生成2条记录的场景
                .withBatchIntervalMs(1000)
                .withMaxRetries(3)
                .build(),
        new JdbcConnectionOptions.JdbcConnectionOptionsBuilder()
                .withUrl("jdbc:postgresql://localhost:5432/postgres?user=michel&password=polnareff")
                .withDriverName("org.postgresql.Driver")
                .build()
);

// 展开List<Row>后接入批量Sink
rows.flatMap((List<Row> rowList, Collector<Row> out) -> {
    for (Row row : rowList) {
        out.collect(row);
    }
}).addSink(jdbcSink);
env.execute();

总结

如果你坚持要保留JDBCOutputFormat,那不需要修改它的任何配置,只需要把DataStream<List<Row>>通过flatMap转换成DataStream<Row>即可。如果想优化写入性能,建议换成新版本的JdbcSink批量模式,代码更简洁且效率更高。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.12 04:42:27