- 大数据
- 批处理
- 流处理
- 数据工程
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
本指南以 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内置了两层容错:
- RetryStrategy(是否值得重试):默认的
DefaultRetryStrategy判断SQLException.getSQLState()是否为40001(多数数据库的死锁状态码)或40P01(PostgreSQL 专用死锁码),命中即重试(JdbcIO.java)。 - 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构造工具,可直接参考其写法搭建本地验证环境。
七、最佳实践与注意事项
- 谨慎使用
INSERT语句:Beam runner 为容错可能重放部分JdbcIO.Write执行(at-least-once 语义),直接INSERT可能造成重复记录或主键冲突。官方注释明确建议改用数据库支持的MERGE(upsert)语句(JdbcIO.java)。 - 善用 PipelineOptions 参数化:示例中的
@Description、@Validation.Required、@Default.String注解组合,让表名、URL、凭据全部可在命令行注入,避免把敏感信息写死在代码里;如需更强的凭据安全,请结合withSecretManager。 - 根据数据规模调节批次:
withBatchSize与withMaxBatchBufferingDuration需要结合数据库负载调优——批次过大增加单次事务时长与死锁概率,过小则放大网络与提交开销。 - 流式写入注意分片:无界数据流写入时若不启用
withAutoSharding(),所有元素会先WithKeys("")聚集到单一 key 上,可能形成热点;withAutoSharding()仅对流式管线生效,源码中对此有显式校验(JdbcIO.java)。 - 连接池与并发控制:默认 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.
相关推荐
Lance 表格式中的 Data Overlay Files:不重写基文件的低成本单元格级更新机制
Lance 表格式中的 Data Overlay Files:不重写基文件的低成本单元格级更新机制 Data Overlay Files(数据覆盖文件)是 La
大数据批处理流处理数据工程Zero 邮件批量发送完整指南:To/Cc/Bcc 三步群发教程
Zero 邮件批量发送完整指南:To/Cc/Bcc 三步群发教程 Zero(Mail0)是一款开源邮件客户端,把多个收件人写进同一封信的能力直接内置在撰写界面里
大数据批处理流处理数据工程Apache Beam 实战:使用 BigQueryIO 向 Google BigQuery 写入数据的 Java 指南
Apache Beam 实战:使用 BigQueryIO 向 Google BigQuery 写入数据的 Java 指南 导读 本文以 Apache Beam
大数据批处理流处理数据工程
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考