ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

Apache Beam Java SDK 实战:用 SpannerIO.read() 从 Cloud Spanner 批量读取表数据并逐行加工

Apache Beam Java SDK 实战:用 SpannerIO.read() 从 Cloud Spanner 批量读取表数据并逐行加工 大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载Apache Beam 的SpannerIO是官方为 Google Cloud Spanner 提供的高层 I/O 连接器封装了 Spanner 的 Read/Query API 与读写事务语义。本文以代码解释型文档 learning/prompts/code-explanation/java/05_io_spanner.md 中的ReadSpannerTable示例为骨架完整讲解如何用 Java SDK 声明 Spanner 读取选项、构建只读事务管道、按列批量读取表数据并用ParDo将 Spanner 的Struct行对象转换为自定义输出。读完本文你将能独立编写一个可运行、可参数化的 Spanner 读取 Pipeline并理解其底层 API 的约束与进阶用法。一、示例全景一个读取 Spanner 表的完整 Pipeline原始示例ReadSpannerTable演示了 Beam 读取 Spanner 最典型的四条主线用自定义PipelineOptions子接口声明运行期可配置参数实例、数据库、表、项目用PipelineOptionsFactory.fromArgs(args).withValidation()解析命令行参数并构造Pipeline用SpannerIO.read().withInstanceId(...).withDatabaseId(...).withTable(...).withColumns(...)构造一次整表读取用ParDoDoFnStruct, String逐行提取字段、打印日志并输出格式化字符串。该示例对应的完整代码如下摘自原文档保持原样以便对照学习public class ReadSpannerTable { private static final Logger LOG LoggerFactory.getLogger(ReadTableSpanner.class); public interface ReadSpannerTableOptions extends DataflowPipelineOptions { Description(Spanner instance) Default.String(test-instance) String getInstanceName(); void setInstanceName(String value); Description(Spanner table to read from ) Default.String(singers) String getTableName(); void setTableName(String value); Description(Spanner Database ID) Default.String(example-db) String getDatabaseName(); void setDatabaseName(String value); Nullable Description(Project ID) String getSpannerProjectName(); void setSpannerProjectName(String value); } public static void main(String[] args) { ReadSpannerTableOptions options PipelineOptionsFactory.fromArgs(args).withValidation().as(ReadSpannerTableOptions.class); Pipeline p Pipeline.create(options); String project (options.getSpannerProjectName() null) ? options.getProject() : options.getSpannerProjectName(); p .apply(SpannerIO.read() .withInstanceId(options.getInstaneName()) .withDatabaseId(options.getDatabaseName()) .withTable(options.getTableName()) .withColumns(SingerId, FirstName, LastName) .withProjectId(project) ) .apply(Process Row, ParDo.of(new DoFnStruct, String() { ProcessElement public void processElement(ProcessContext c) { Struct struct c.element(); Long singerId struct.getLong(SingerId); String firstName struct.getString(FirstName); String lastName struct.getString(LastName); String row String.format(ID %d, First name %s, Last name %s, singerId, firstName, lastName); LOG.info(row); c.output(row); } }) ); p.run(); } }整条链路的执行逻辑是SpannerIO.read()返回一个PCollectionStruct其中每个Struct元素对应表中一行随后的ParDo以该Struct为输入取出SingerId、FirstName、LastName三列拼成字符串并输出最终p.run()提交执行。需要注意的是DoFn直接操作的是com.google.cloud.spanner.Struct类型getLong/getString按列名取值列名必须与withColumns中声明的列一致。二、逐块拆解选项接口、读取构造与行处理2.1 ReadSpannerTableOptions用注解声明 Pipeline 参数原文档指出ReadSpannerTableOptions接口用于声明 Spanner 实例、表和数据库public interface ReadSpannerTableOptions extends DataflowPipelineOptions { Description(Spanner instance) Default.String(test-instance) String getInstanceName(); void setInstanceName(String value); Description(Spanner table to read from ) Default.String(singers) String getTableName(); void setTableName(String value); Description(Spanner Database ID) Default.String(example-db) String getDatabaseName(); void setDatabaseName(String value); Nullable Description(Project ID) String getSpannerProjectName(); void setSpannerProjectName(String value); }其要点如下Description为每个选项提供说明文案--help时会展示Default.String设置默认值例如未指定时实例默认为test-instance、表默认为singers、数据库默认为example-dbgetSpannerProjectName()是可选参数用Nullable标记且无默认值——这正是示例中需要兜底逻辑的原因接口继承DataflowPipelineOptions因此天然拥有--project、--region、--runner等 Dataflow 运行参数其中--project由 GCP 核心扩展提供见 GcpOptions.java。选项值最终以命令行形式注入例如--runnerDataflowRunner \ --projectmy-gcp-project \ --regionus-central1 \ --instanceNametest-instance \ --databaseNameexample-db \ --tableNamesingers2.2 main()解析参数并创建 Pipeline原文档给出了参数解析与管道创建的骨架ReadSpannerTableOptions options PipelineOptionsFactory.fromArgs(args).withValidation().as(ReadSpannerTableOptions.class); Pipeline p Pipeline.create(options);fromArgs(args)从命令行解析选项withValidation()触发参数校验如必填项缺失会报错as(ReadSpannerTableOptions.class)将解析结果转换为自定义选项类型。随后Pipeline.create(options)依据选项中的 Runner 配置构建管道。项目 ID 的取值逻辑值得单独说明String project (options.getSpannerProjectName() null) ? options.getProject() : options.getSpannerProjectName();即若用户显式传入--spannerProjectName则优先使用否则回退到通用--project。这种设计让 Spanner 可以位于与作业执行项目不同的 GCP 项目中是跨项目读取的常见手法。2.3 SpannerIO.read()构造一次整表读取p .apply(SpannerIO.read() .withInstanceId(options.getInstaneName()) .withDatabaseId(options.getDatabaseName()) .withTable(options.getTableName()) .withColumns(SingerId, FirstName, LastName) .withProjectId(project) )SpannerIO.read()是 Apache Beam 提供的Read变换入口返回PCollectionStruct。结合 SpannerIO.java 源码约 L889-L1076本例使用的链式方法含义如下方法作用底层实现源码依据withProjectId(String)指定 Spanner 所属 GCP 项目写入SpannerConfig.withProjectIdL889-L897withInstanceId(String)指定 Spanner 实例 ID写入SpannerConfig.withInstanceIdL900-L908withDatabaseId(String)指定数据库 ID写入SpannerConfig.withDatabaseIdL911-L919withTable(String)指定要整表读取的表名通过ReadOperation.withTable记录表名L1042-L1044withColumns(String...)指定要读取的列可传多个列名通过ReadOperation.withColumns记录列清单L1050-L1056SpannerIO.read()对“表读取”模式的校验在expand()中完成L1106-L1132必须通过withTimestampBound或withTimestamp显式设置只读事务的时间戳约束withQuery与withTable不能同时设置表读取必须提供非空的列清单withColumnswithQuery与withTable至少设置其一。按表读取时源码会在createSourceDef()L1097-L1103中构造SpannerTableSourceDef底层走 Cloud Spanner 的 Batch Read API支持并行分区读取。2.4 ParDo 行处理从 Struct 到格式化字符串原文档给出的DoFn实现.apply(Process Row, ParDo.of(new DoFnStruct, String() { ProcessElement public void processElement(ProcessContext c) { Struct struct c.element(); Long singerId struct.getLong(SingerId); String firstName struct.getString(FirstName); String lastName struct.getString(LastName); String row String.format(ID %d, First name %s, Last name %s, singerId, firstName, lastName); LOG.info(row); c.output(row); } }) );这里ParDo是 Beam 最核心的逐元素变换每个Struct一行都会触发一次ProcessElement。c.element()取当前元素Struct.getLong(...)/getString(...)按列名取类型化字段值最后c.output(row)把拼接结果输出到下游PCollectionString。LOG.info(row)便于在 Runner 日志中观测读取进度c.output(row)则让结果可供后续 Sink如写 GCS、BigQuery继续使用。2.5 p.run()提交执行p.run();p.run()是 Beam Pipeline 的执行入口在 Direct Runner 下于本地串行/并行执行在 Dataflow Runner 下则打包并提交到云端执行。执行与否、采用何种 Runner均由--runner选项控制。三、源码与测试佐证从示例到可验证的实现事实3.1 官方单元测试验证同款用法Beam 官方测试 SpannerIOReadTest.java 中runBatchReadTestWithProjectIdL407-L415与示例的构造方式高度一致SpannerIO.read() .withSpannerConfig(spannerConfig) .withTable(TABLE_ID) .withColumns(id, name) .withTimestampBound(TIMESTAMP_BOUND);可见“withTablewithColumns 时间戳约束”是官方认可的批量表读取标准组合。测试还覆盖了不指定项目、项目为 null、withHighPriority()优先级、Data Boost 等场景L417-L471证明示例中的“项目 ID 可选”设计有对应的运行时行为。3.2 集成测试中的真实表读取端到端集成测试 SpannerReadIT.java 同样使用withTable(options.getTable()).withColumns(Key, Value)的组合读取真实 Spanner 表约 L188-L214并包含对错误表名的失败路径验证L272-L273。这说明示例面向的“singers歌手表”只是可替换的表名占位实际使用时把withTable/withColumns换成你自己的表与列即可。3.3 从expand()校验看示例的隐含前提回顾 SpannerIO.java L1106-L1132 的校验逻辑示例能正常运行依赖两个隐含前提时间戳约束SpannerIO.read()强制要求设置withTimestampBound如TimestampBound.strong()或withTimestamp否则直接抛异常。示例虽未显式调用但在实际生产中必须补上例如SpannerIO.read() .withInstanceId(instance) .withDatabaseId(database) .withTable(table) .withColumns(SingerId, FirstName, LastName) .withTimestampBound(TimestampBound.strong()) // 强一致快照表读取必须指定非空列清单只调用withTable而不调用withColumns会在expand()阶段报错因此示例中withColumns(SingerId, FirstName, LastName)是必不可少的。四、进阶扩展同一 Read 变换的更多玩法Read变换远不止整表读取一种形态SpannerIO.java 还提供了以下常用能力可在示例基础上自由组合SQL 查询读取withQuery(String sql)或withQuery(Statement)替代withTable返回查询结果行。官方 JavadocL180-L189示例PCollectionStruct rows p.apply( SpannerIO.read() .withInstanceId(instanceId) .withDatabaseId(dbId) .withQuery(SELECT id, name, email FROM users));二级索引读取withIndex(users_by_name)配合withTable与withColumns按索引读表L212-L223只读事务共享SpannerIO.createTransaction()生成事务PCollectionView多个read()通过withTransaction(tx)在同一一致性快照下读取多张表L233-L257批处理开关withBatching(boolean)控制是否使用 Cloud Spanner Batch APIL1019-L1022。默认走 PartitionQuery/PartitionRead 并行分区若查询不支持分区可设withBatching(false)降级为单次读取时间戳与一致性withTimestamp(Timestamp)与withTimestampBound(TimestampBound)控制读取快照的新鲜度L1034-L1040RPC 优先级withLowPriority()/withHighPriority()设置 Spanner RPC 优先级L1087-L1095低优先级适合后台批处理避免抢占线上流量多表/多查询批量一致读取SpannerIO.readAll()接收PCollectionReadOperation对多张表/多个查询在同一个只读事务内完成一致性读取官方 Javadoc L259-L280注意该变换不适合流式管道因为只读事务创建一次后 1 小时会被 Spanner 服务端超时关闭。五、运行前提与注意事项依赖需要引入beam-sdks-java-io-google-cloud-platform模块其中包含org.apache.beam.sdk.io.gcp.spanner.SpannerIO凭据需要配置具备 Spanner 读取权限的 GCP 凭据如GOOGLE_APPLICATION_CREDENTIALS指向服务账号 JSON并在 GCP 上预先创建好实例、数据库与表项目回退逻辑示例中getSpannerProjectName()为可选未设置时回退到options.getProject()options.getProject()来自DataflowPipelineOptions对应--project参数表读取的必填项withTable 非空withColumns 时间戳约束三者缺一不可这是 SpannerIO.java 的expand()校验强制的行为示例中的笔误原示例withInstanceId(options.getInstaneName())中getInstaneName与方法声明getInstanceName拼写不一致实际运行时请统一为getInstanceName()运行方式本地调试用--runnerDirectRunner生产环境用--runnerDataflowRunner配合--project、--region、--tempLocation等 Dataflow 参数提交到 Google Cloud Dataflow。六、总结ReadSpannerTable示例浓缩了 Beam Spanner 读取的全部关键要素注解驱动的选项接口、SpannerIO.read()的链式配置、Struct行对象到业务字符串的ParDo转换以及项目 ID 的优雅回退。结合 SpannerIO.java 源码与 SpannerIOReadTest.java、SpannerReadIT.java 测试用例你可以放心地将其改造为生产级管道补上withTimestampBound、替换真实表名列名、按需切换到withQuery或readAll()并利用withLowPriority、Data Boost 等选项在吞吐与成本之间取得平衡。赞分享大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载相关推荐Apache Beam Java Kata 实战用 TextIO.read() 从文本文件读取 PCollectionApache Beam Java Kata 实战用 TextIO.read 从文本文件读取 PCollection 本篇技术指南以 Apache Beam 官大数据批处理流处理数据工程Apache Beam CdapIO 实战构建从 ServiceNow 批量拉取数据并落盘 TXT 的管道示例Apache Beam CdapIO 实战构建从 ServiceNow 批量拉取数据并落盘 TXT 的管道示例 本文基于 Apache Beam 仓库中 ex大数据批处理流处理数据工程Apache Beam Go SDK 实战用 Side Input 在 ParDo 中注入运行时附加数据Apache Beam Go SDK 实战用 Side Input 在 ParDo 中注入运行时附加数据 Side Input侧输入是 Apache Be大数据批处理流处理数据工程上一篇3种方案解决Zotero PDF Translate插件版本兼容性问题下一篇卡牌批量生成神器3分钟搞定100张桌游卡牌效率提升300%创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进