☰
Apache Beam Java 实战:使用 JdbcIO 连接器向 JDBC 数据库写入数据
2026/9/29 8:16:03 网站建设 项目流程
  • 大数据
  • 批处理
  • 流处理
  • 数据工程

【免费下载链接】beam

Apache Beam is a unified programming model for Batch and Streaming data processing.

项目地址:https://gitcode.com/gh_mirrors/beam4/beam
点击查看免费下载

本指南以 Apache Beam 官方仓库中的代码生成示例(learning/prompts/code-generation/java/11_io_jdbc.md)为核心,完整讲解如何用 Java SDK 中的JdbcIO连接器把PCollection写入任意支持 JDBC 的关系型数据库(Oracle、PostgreSQL、MySQL 等)。你将掌握JdbcIO.DataSourceConfiguration的构建方式、JdbcIO.write()的完整参数用法、PipelineOptions 命令行参数化模式,以及底层批处理、重试与连接池机制,最终可以独立编写一个可运行、可配置的 JDBC Sink 写入管线。

一、JdbcIO 与 JDBC Sink 写入场景概述

Apache Beam 的 Java SDK 提供了一套统一的批流编程模型,而org.apache.beam.sdk.io.jdbc.JdbcIO正是连接这套模型与 JDBC 生态的官方 I/O 连接器。它既支持读取(JdbcIO.read()、readAll()、readRows()、readWithPartitions()),也支持写入(JdbcIO.write()、writeVoid()、writeWithResults()),因此可以把任意关系型数据库同时当作数据源和数据汇。

写入场景非常典型:从 Kafka、Pub/Sub、BigQuery 等上游取数,经过转换后落地到 Oracle / PostgreSQL / MySQL 等 JDBC 兼容数据库。JdbcIO.Write的核心机制是把PCollection中的每个元素(T)通过用户提供的PreparedStatementSetter绑定到一条PreparedStatement上,再按批(batch)执行并提交事务。

从源码结构看,整个连接器收敛在 sdks/java/io/jdbc/src/main/java/org/apache/beam/sdk/io/jdbc/JdbcIO.java 这一个类中,内部按职责拆分为DataSourceConfiguration(连接配置)、Write/WriteVoid(写入变换)、WriteFn(实际执行批量写入的 DoFn)、RetryConfiguration/RetryStrategy(容错重试)等组件,单元测试与集成测试位于 sdks/java/io/jdbc/src/test/java/org/apache/beam/sdk/io/jdbc/JdbcIOTest.java 与 sdks/java/io/jdbc/src/test/java/org/apache/beam/sdk/io/jdbc/JdbcIOIT.java。

二、完整示例:向 JDBC Sink 写入数据

以下代码即仓库中 11_io_jdbc.md 提供的标准示例,演示如何用JdbcIO连接器把一组样例数据写入 JDBC 数据库。它采用了 Beam 官方的PipelineOptions 模式:把表名、JDBC URL、驱动类名、用户名、密码全部抽象成命令行参数,既避免硬编码,又方便在 Direct Runner、Dataflow、Flink、Spark 等不同执行引擎间迁移。

package jdbc; import org.apache.beam.sdk.Pipeline; import org.apache.beam.sdk.io.jdbc.JdbcIO; import org.apache.beam.sdk.options.Default; import org.apache.beam.sdk.options.Description; import org.apache.beam.sdk.options.PipelineOptions; import org.apache.beam.sdk.options.PipelineOptionsFactory; import org.apache.beam.sdk.options.Validation; import org.apache.beam.sdk.transforms.Create; import java.io.Serializable; import java.util.Arrays; import java.util.List; // Pipeline to write data to a JDBC sink using the Apache Beam JdbcIO connector public class WriteJdbcSink { // Class representing the data to be written to the JDBC sink public static class ExampleRow implements Serializable { private int id; private String month; private String amount; public ExampleRow() {} public ExampleRow(int id, String month, String amount) { this.id = id; this.month = month; this.amount = amount; } public int getId() { return id; } public String getMonth() { return month; } public String getAmount() { return amount; } } // Pipeline options for writing data to the JDBC sink public interface WriteJdbcSinkOptions extends PipelineOptions { @Description("Table name to write to") @Validation.Required String getTableName(); void setTableName(String tableName); @Description("JDBC sink URL") @Validation.Required String getJdbcSinkUrl(); void setJdbcSinkUrl(String jdbcSinkUrl); @Description("JDBC driver class name") @Default.String("org.postgresql.Driver") String getDriverClassName(); void setDriverClassName(String driverClassName); @Description("DB Username") @Validation.Required String getSinkUsername(); void setSinkUsername(String username); @Description("DB password") @Validation.Required String getSinkPassword(); void setSinkPassword(String password); } // Main method to run the pipeline public static void main(String[] args) { // Parse the pipeline options from the command line WriteJdbcSinkOptions options = PipelineOptionsFactory.fromArgs(args).withValidation().as(WriteJdbcSinkOptions.class); // Create the JDBC sink configuration using the provided options JdbcIO.DataSourceConfiguration config = JdbcIO.DataSourceConfiguration.create(options.getDriverClassName(), options.getJdbcSinkUrl()) .withUsername(options.getSinkUsername()) .withPassword(options.getSinkPassword()); // Create the pipeline Pipeline p = Pipeline.create(options); // Create sample rows to write to the JDBC sink List<ExampleRow> rows = Arrays.asList( new ExampleRow(1, "January", "$1000"), new ExampleRow(2, "February", "$2000"), new ExampleRow(3, "March", "$3000") ); // // Create PCollection from the list of rows p.apply("Create collection of records", Create.of(rows)) // Write the rows to the JDBC sink .apply( "Write to JDBC Sink", JdbcIO.<ExampleRow>write() .withDataSourceConfiguration(config) .withStatement(String.format("insert into %s values(?, ?, ?)", options.getTableName())) .withBatchSize(10L) .withPreparedStatementSetter( (element, statement) -> { statement.setInt(1, element.getId()); statement.setString(2, element.getMonth()); statement.setString(3, element.getAmount()); })); // Run the pipeline p.run(); } }

运行该程序时,通过命令行传入参数即可,例如(PostgreSQL 场景):

java -cp beam-sdks-java-io-jdbc.jar:postgresql-42.x.x.jar:beam-runners-direct-java.jar \ jdbc.WriteJdbcSink \ --tableName=sales \ --jdbcSinkUrl=jdbc:postgresql://localhost:5432/mydb \ --driverClassName=org.postgresql.Driver \ --sinkUsername=beam_user \ --sinkPassword=secret

其中@Validation.Required标注的参数缺失时,withValidation()会在启动阶段直接报错;driverClassName因带有@Default.String("org.postgresql.Driver")默认值,即使不传也会使用 PostgreSQL 驱动。若目标库是 Oracle 或 MySQL,仅需把驱动类名与 JDBC URL 一并替换(如oracle.jdbc.OracleDriver/jdbc:oracle:thin:@//host:1521/service或com.mysql.cj.jdbc.Driver/jdbc:mysql://host:3306/mydb)。

三、DataSourceConfiguration:连接配置的构建与可选参数

示例中通过JdbcIO.DataSourceConfiguration.create(driverClassName, url)创建配置,随后链式调用.withUsername(...)与.withPassword(...)。对应源码中提供了两种创建入口(JdbcIO.java):

  • create(DataSource dataSource):直接传入一个已构建好的javax.sql.DataSource(要求可序列化),适用于你已经持有连接池等定制化 DataSource 的场景;
  • create(String driverClassName, String url):仅凭驱动类名与 URL 构建,底层在buildDatasource()中借助 Apache Commons DBCP2 的BasicDataSource组装(JdbcIO.java)。

除了示例中用到的三个方法,DataSourceConfiguration(JdbcIO.java)还提供了一批可选的链式配置,适用于更复杂的生产场景:

方法作用注意事项
withConnectionProperties(String)以[propertyName=property;]*格式向driver.connect(...)传递连接属性user/password无需重复设置,直接使用withUsername/withPassword即可
withConnectionInitSqls(Collection<String>)设置连接初始化 SQL(如SET ...)仅 MySQL / MariaDB 支持,其他数据库会抛出 SQL 异常
withMaxConnections(Integer)连接池最大连接数传负数表示不限制
withQueryTimeout(Integer)连接默认查询超时(秒级)作用于 DBCP2 的defaultQueryTimeout
withDriverClassLoader(ClassLoader)指定加载 JDBC 驱动的 ClassLoader不指定时使用默认 ClassLoader
withDriverJars(String)逗号分隔的驱动 Jar 路径,跨文件系统(如gs://bucket/driver.jar,gs://bucket/driver2.jar)底层会把远端 Jar 下载到本地后用URLClassLoader加载
withSecretManager(String)指定密钥管理器提供商(GoogleCloudSecretManager、GoogleCloudHsmGeneratedSecretManager)配合withPassword传入 JSON 格式密钥规格,避免明文密码落盘

关于密码与密钥管理,源码注释明确指出:withPassword既可以传明文密码,也可以配合withSecretManager传入 JSON 密钥规格(如{"name": "my-db-secret", "project": "my-project"}),由密钥管理器在buildDatasource()阶段实时拉取真实密码(JdbcIO.java),这是避免敏感凭据明文存放的推荐做法。

此外,连接器内部对 DataSource 做了进程内单例缓存:DataSourceProviderFromDataSourceConfiguration使用ConcurrentHashMap保证同一份配置在整个 pipeline 中只构建一次 DataSource(JdbcIO.java);若担心默认行为下每个执行线程各拿一个 DataSource 导致连接数过大,可使用JdbcIO.PoolableDataSourceProvider.of(config)显式启用 DBCP2 连接池(JdbcIO.java)。

四、JdbcIO.write():写入变换的完整参数体系

示例中组装了JdbcIO.<ExampleRow>write()并设置了四个核心方法,下面结合 JdbcIO.java(Write门面类)与 JdbcIO.java(WriteVoid实际实现)逐一说明:

  • withDataSourceConfiguration(config):绑定连接配置,等价于把配置包装成SerializableFunction<Void, DataSource>;也可用withDataSourceProviderFn(...)直接提供自定义 DataSource 工厂函数。
  • withStatement(String):写入 SQL 模板,?占位符由 PreparedStatementSetter 填充。示例中使用insert into %s values(?, ?, ?)并按表名参数动态拼接。该方法是必选项(除非使用withTable走 schema 自动生成路径,见下文)。
  • withPreparedStatementSetter(PreparedStatementSetter<T>):定义"元素 → PreparedStatement 参数"的映射逻辑,对应org.apache.beam.sdk.io.jdbc.JdbcIO.PreparedStatementSetter函数式接口(JdbcIO.java),也是必选项。
  • withBatchSize(long):每个批次最多包含的 SQL 语句数,默认值为 1000(DEFAULT_BATCH_SIZE,JdbcIO.java)。达到该上限或超过最大缓冲时长即触发一次executeBatch()+commit()。示例中设为10L,适合小批量演示。

Write门面还透传了以下进阶参数(均委托给WriteVoid):

方法默认值说明
withMaxBatchBufferingDuration(long)200(毫秒,DEFAULT_MAX_BATCH_BUFFERING_DURATION)批量提交前的最大缓冲时长,与batchSize二选一先到先触发
withAutoSharding()关闭仅适用于**流式(无界)**管线,使用动态分片键(sharded key)避免单 key 热点导致批次倾斜
withRetryStrategy(RetryStrategy)DefaultRetryStrategy自定义"哪些 SQLException 值得重试"的判断逻辑
withRetryConfiguration(RetryConfiguration)create(5, null, Duration.standardSeconds(5))指数退避重试参数:最大 5 次尝试、初始退避 1 秒、累计退避上限 1000 天(JdbcIO.java)
withTable(String)无当输入 PCollection 带 Beam Schema 时,可省略withStatement/withPreparedStatementSetter,由连接器自动比对目标表结构并生成INSERT INTO table(col1, col2, ...) VALUES(?, ?, ...)语句(JdbcUtil.java)
withResults()/withWriteResults(RowMapper)无返回PCollection<Void>或逐行写入结果,可与Wait.on(...)配合实现"写库完成后才触发下游"的跨库编排

批处理与事务的底层实现

从源码可以清晰看到写入的执行链路:WriteVoid.expand()先调用batchElements()把元素聚合成Iterable<T>批次——有界输入用 DoFn 在 bundle 内累积列表,无界输入则走GroupIntoBatches.ofSize(batchSize).withMaxBufferingDuration(...)(JdbcIO.java);随后WriteFn(JdbcIO.java)在executeBatch()中对每个元素调用PreparedStatementSetter、addBatch(),最后统一executeBatch()并commit()。值得注意的是WriteFn会显式connection.setAutoCommit(false)并自行管理提交(JdbcIO.java),同时通过RECORDS_PER_BATCH、MS_PER_BATCH两个 Beam Metrics 分布指标持续上报每批记录数与耗时(JdbcIO.java),便于在运行监控中观测写入吞吐。

五、容错与重试机制

分布式环境下数据库瞬时故障(尤其是死锁)不可避免,JdbcIO.Write内置了两层容错:

  1. RetryStrategy(是否值得重试):默认的DefaultRetryStrategy判断SQLException.getSQLState()是否为40001(多数数据库的死锁状态码)或40P01(PostgreSQL 专用死锁码),命中即重试(JdbcIO.java)。
  2. RetryConfiguration(如何重试):基于FluentBackoff的指数退避,create(maxAttempts, maxDuration, initialDuration)三个参数均可配,传入null或零值时回落到默认值——初始退避 1 秒、累计上限 1000 天(JdbcIO.java)。

重试流程在executeBatch()中可见:捕获 SQLException 后先调用retryStrategy.apply(exception)判断,命中则clearBatch()+connection.rollback()清理批次状态,再按退避策略休眠后重放整个批次(JdbcIO.java)。集成测试 JdbcIOExceptionHandlingParameterizedTest.java 与单元测试 JdbcIOTest.java 对该链路均有覆盖。

六、运行验证与测试证据

仓库对 JDBC 写入提供了完备的测试支撑,可作为验证与学习素材:

  • JdbcIOTest.java:基于内存 H2 数据库的单元测试,覆盖write()、writeVoid()、withBatchSize(10L)、withRetryConfiguration、schema 自动生成 INSERT 等路径(如 JdbcIOTest.java#L776-L790 所示);
  • JdbcIOIT.java、JdbcIOPostgresIT.java:针对真实数据库的集成测试;
  • 测试辅助类 JdbcTestHelper.java 提供 H2 建表与DataSource构造工具,可直接参考其写法搭建本地验证环境。

七、最佳实践与注意事项

  1. 谨慎使用INSERT语句:Beam runner 为容错可能重放部分JdbcIO.Write执行(at-least-once 语义),直接INSERT可能造成重复记录或主键冲突。官方注释明确建议改用数据库支持的MERGE(upsert)语句(JdbcIO.java)。
  2. 善用 PipelineOptions 参数化:示例中的@Description、@Validation.Required、@Default.String注解组合,让表名、URL、凭据全部可在命令行注入,避免把敏感信息写死在代码里;如需更强的凭据安全,请结合withSecretManager。
  3. 根据数据规模调节批次:withBatchSize与withMaxBatchBufferingDuration需要结合数据库负载调优——批次过大增加单次事务时长与死锁概率,过小则放大网络与提交开销。
  4. 流式写入注意分片:无界数据流写入时若不启用withAutoSharding(),所有元素会先WithKeys("")聚集到单一 key 上,可能形成热点;withAutoSharding()仅对流式管线生效,源码中对此有显式校验(JdbcIO.java)。
  5. 连接池与并发控制:默认 DataSource 按执行线程请求连接,高并发下可能打爆数据库连接数,生产环境优先使用PoolableDataSourceProvider或自定义共享单例 DataSource。

至此,你已经可以对照 11_io_jdbc.md 的示例与 JdbcIO.java 的源码,快速搭建并调优属于自己的 JDBC Sink 管线。

  • 大数据
  • 批处理
  • 流处理
  • 数据工程

【免费下载链接】beam

Apache Beam is a unified programming model for Batch and Streaming data processing.

项目地址:https://gitcode.com/gh_mirrors/beam4/beam
点击查看免费下载

相关推荐

上一篇:Foundry 安全加固:`forge script` 敏感缓存文件(RPC URL)的 Unix 权限控制
下一篇:Kubernetes 批处理工作组 2022 年度报告解读:Job API 增强、Kueue 与拓扑感知调度

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询