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
相关产品推荐
相关产品推荐

