尧图网络科技YAOTU DIGITAL 获取报价
获取报价
首页 / 资讯中心 / 文章详情

Apache Beam Java 实战:使用 JdbcIO 连接器向 JDBC 数据库写入数据

发布时间:2026/9/29 8:16:07

资讯中心
01
ARTICLE

Apache Beam Java 实战:使用 JdbcIO 连接器向 JDBC 数据库写入数据

Apache Beam Java 实战:使用 JdbcIO 连接器向 JDBC 数据库写入数据
大数据批处理流处理数据工程【免费下载链接】beamApache 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 ListExampleRow 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.ExampleRowwrite() .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 \ --tableNamesales \ --jdbcSinkUrljdbc:postgresql://localhost:5432/mydb \ --driverClassNameorg.postgresql.Driver \ --sinkUsernamebeam_user \ --sinkPasswordsecret其中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.javacreate(DataSource dataSource)直接传入一个已构建好的javax.sql.DataSource要求可序列化适用于你已经持有连接池等定制化 DataSource 的场景create(String driverClassName, String url)仅凭驱动类名与 URL 构建底层在buildDatasource()中借助 Apache Commons DBCP2 的BasicDataSource组装JdbcIO.java。除了示例中用到的三个方法DataSourceConfigurationJdbcIO.java还提供了一批可选的链式配置适用于更复杂的生产场景方法作用注意事项withConnectionProperties(String)以[propertyNameproperty;]*格式向driver.connect(...)传递连接属性user/password无需重复设置直接使用withUsername/withPassword即可withConnectionInitSqls(CollectionString)设置连接初始化 SQL如SET ...仅 MySQL / MariaDB 支持其他数据库会抛出 SQL 异常withMaxConnections(Integer)连接池最大连接数传负数表示不限制withQueryTimeout(Integer)连接默认查询超时秒级作用于 DBCP2 的defaultQueryTimeoutwithDriverClassLoader(ClassLoader)指定加载 JDBC 驱动的 ClassLoader不指定时使用默认 ClassLoaderwithDriverJars(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 中只构建一次 DataSourceJdbcIO.java若担心默认行为下每个执行线程各拿一个 DataSource 导致连接数过大可使用JdbcIO.PoolableDataSourceProvider.of(config)显式启用 DBCP2 连接池JdbcIO.java。四、JdbcIO.write()写入变换的完整参数体系示例中组装了JdbcIO.ExampleRowwrite()并设置了四个核心方法下面结合 JdbcIO.javaWrite门面类与 JdbcIO.javaWriteVoid实际实现逐一说明withDataSourceConfiguration(config)绑定连接配置等价于把配置包装成SerializableFunctionVoid, DataSource也可用withDataSourceProviderFn(...)直接提供自定义 DataSource 工厂函数。withStatement(String)写入 SQL 模板?占位符由 PreparedStatementSetter 填充。示例中使用insert into %s values(?, ?, ?)并按表名参数动态拼接。该方法是必选项除非使用withTable走 schema 自动生成路径见下文。withPreparedStatementSetter(PreparedStatementSetterT)定义元素 → PreparedStatement 参数的映射逻辑对应org.apache.beam.sdk.io.jdbc.JdbcIO.PreparedStatementSetter函数式接口JdbcIO.java也是必选项。withBatchSize(long)每个批次最多包含的 SQL 语句数默认值为 1000DEFAULT_BATCH_SIZEJdbcIO.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.javawithTable(String)无当输入 PCollection 带 Beam Schema 时可省略withStatement/withPreparedStatementSetter由连接器自动比对目标表结构并生成INSERT INTO table(col1, col2, ...) VALUES(?, ?, ...)语句JdbcUtil.javawithResults()/withWriteResults(RowMapper)无返回PCollectionVoid或逐行写入结果可与Wait.on(...)配合实现写库完成后才触发下游的跨库编排批处理与事务的底层实现从源码可以清晰看到写入的执行链路WriteVoid.expand()先调用batchElements()把元素聚合成IterableT批次——有界输入用 DoFn 在 bundle 内累积列表无界输入则走GroupIntoBatches.ofSize(batchSize).withMaxBufferingDuration(...)JdbcIO.java随后WriteFnJdbcIO.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多数数据库的死锁状态码或40P01PostgreSQL 专用死锁码命中即重试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可能造成重复记录或主键冲突。官方注释明确建议改用数据库支持的MERGEupsert语句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 管线。赞分享大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载相关推荐Lance 表格式中的 Data Overlay Files不重写基文件的低成本单元格级更新机制Lance 表格式中的 Data Overlay Files不重写基文件的低成本单元格级更新机制 Data Overlay Files数据覆盖文件是 La大数据批处理流处理数据工程Zero 邮件批量发送完整指南To/Cc/Bcc 三步群发教程Zero 邮件批量发送完整指南To/Cc/Bcc 三步群发教程 ZeroMail0是一款开源邮件客户端把多个收件人写进同一封信的能力直接内置在撰写界面里大数据批处理流处理数据工程Apache Beam 实战使用 BigQueryIO 向 Google BigQuery 写入数据的 Java 指南Apache Beam 实战使用 BigQueryIO 向 Google BigQuery 写入数据的 Java 指南 导读 本文以 Apache Beam大数据批处理流处理数据工程上一篇Foundry 安全加固forge script 敏感缓存文件RPC URL的 Unix 权限控制下一篇Kubernetes 批处理工作组 2022 年度报告解读Job API 增强、Kueue 与拓扑感知调度创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
02
RELATED NEWS

相关资讯

更多网站建设与数字化升级内容

03
WHY YAOTU

想打造同款高转化官网?

懂行业、懂生意,从建站到增长一站式陪跑

◈

场景化定制

不做模板站,围绕你的业务场景量身设计,小众不撞款。

◐

营销型架构

以转化目标组织内容与路径,让官网真正带来询盘。

▲

全周期服务

设计、开发、运营、运维一体,上线只是开始。

免费获取你的建站方案

留下需求,专属顾问 24 小时内为你输出方案建议。