news 2026/9/28 3:30:20

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

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Apache Beam Java SDK 实战:用 SpannerIO.read() 从 Cloud Spanner 批量读取表数据并逐行加工
  • 大数据
  • 批处理
  • 流处理
  • 数据工程

【免费下载链接】beam

Apache 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 最典型的四条主线:

  1. 用自定义PipelineOptions子接口声明运行期可配置参数(实例、数据库、表、项目);
  2. 用PipelineOptionsFactory.fromArgs(args).withValidation()解析命令行参数并构造Pipeline;
  3. 用SpannerIO.read().withInstanceId(...).withDatabaseId(...).withTable(...).withColumns(...)构造一次整表读取;
  4. 用ParDo+DoFn<Struct, 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 DoFn<Struct, 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()返回一个PCollection<Struct>,其中每个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-db;
  • getSpannerProjectName()是可选参数,用@Nullable标记且无默认值——这正是示例中需要兜底逻辑的原因;
  • 接口继承DataflowPipelineOptions,因此天然拥有--project、--region、--runner等 Dataflow 运行参数,其中--project由 GCP 核心扩展提供(见 GcpOptions.java)。

选项值最终以命令行形式注入,例如:

--runner=DataflowRunner \ --project=my-gcp-project \ --region=us-central1 \ --instanceName=test-instance \ --databaseName=example-db \ --tableName=singers

2.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变换入口,返回PCollection<Struct>。结合 SpannerIO.java 源码(约 L889-L1076),本例使用的链式方法含义如下:

方法作用底层实现(源码依据)
withProjectId(String)指定 Spanner 所属 GCP 项目写入SpannerConfig.withProjectId(L889-L897)
withInstanceId(String)指定 Spanner 实例 ID写入SpannerConfig.withInstanceId(L900-L908)
withDatabaseId(String)指定数据库 ID写入SpannerConfig.withDatabaseId(L911-L919)
withTable(String)指定要整表读取的表名通过ReadOperation.withTable记录表名(L1042-L1044)
withColumns(String...)指定要读取的列(可传多个列名)通过ReadOperation.withColumns记录列清单(L1050-L1056)

SpannerIO.read()对“表读取”模式的校验在expand()中完成(L1106-L1132):

  • 必须通过withTimestampBound或withTimestamp显式设置只读事务的时间戳约束;
  • withQuery与withTable不能同时设置;
  • 表读取必须提供非空的列清单(withColumns);
  • withQuery与withTable至少设置其一。

按表读取时,源码会在createSourceDef()(L1097-L1103)中构造SpannerTableSourceDef,底层走 Cloud Spanner 的 Batch Read API,支持并行分区读取。

2.4 ParDo 行处理:从 Struct 到格式化字符串

原文档给出的DoFn实现:

.apply("Process Row", ParDo.of(new DoFn<Struct, 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)把拼接结果输出到下游PCollection<String>。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 中,runBatchReadTestWithProjectId(L407-L415)与示例的构造方式高度一致:

SpannerIO.read() .withSpannerConfig(spannerConfig) .withTable(TABLE_ID) .withColumns("id", "name") .withTimestampBound(TIMESTAMP_BOUND);

可见“withTable+withColumns+ 时间戳约束”是官方认可的批量表读取标准组合。测试还覆盖了不指定项目、项目为 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 的校验逻辑,示例能正常运行依赖两个隐含前提:

  1. 时间戳约束:SpannerIO.read()强制要求设置withTimestampBound(如TimestampBound.strong())或withTimestamp,否则直接抛异常。示例虽未显式调用,但在实际生产中必须补上,例如:
SpannerIO.read() .withInstanceId(instance) .withDatabaseId(database) .withTable(table) .withColumns("SingerId", "FirstName", "LastName") .withTimestampBound(TimestampBound.strong()) // 强一致快照
  1. 表读取必须指定非空列清单:只调用withTable而不调用withColumns会在expand()阶段报错,因此示例中withColumns("SingerId", "FirstName", "LastName")是必不可少的。

四、进阶扩展:同一 Read 变换的更多玩法

Read变换远不止整表读取一种形态,SpannerIO.java 还提供了以下常用能力,可在示例基础上自由组合:

  • SQL 查询读取:withQuery(String sql)或withQuery(Statement)替代withTable,返回查询结果行。官方 Javadoc(L180-L189)示例:
PCollection<Struct> 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 API(L1019-L1022)。默认走 PartitionQuery/PartitionRead 并行分区;若查询不支持分区,可设withBatching(false)降级为单次读取;
  • 时间戳与一致性:withTimestamp(Timestamp)与withTimestampBound(TimestampBound)控制读取快照的新鲜度(L1034-L1040);
  • RPC 优先级:withLowPriority()/withHighPriority()设置 Spanner RPC 优先级(L1087-L1095),低优先级适合后台批处理,避免抢占线上流量;
  • 多表/多查询批量一致读取:SpannerIO.readAll()接收PCollection<ReadOperation>,对多张表/多个查询在同一个只读事务内完成一致性读取(官方 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();
  • 运行方式:本地调试用--runner=DirectRunner,生产环境用--runner=DataflowRunner(配合--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 等选项在吞吐与成本之间取得平衡。

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

【免费下载链接】beam

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

项目地址:https://gitcode.com/gh_mirrors/beam4/beam
点击查看免费下载
上一篇:3种方案解决Zotero PDF Translate插件版本兼容性问题
下一篇:卡牌批量生成神器:3分钟搞定100张桌游卡牌,效率提升300%

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

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/28 3:26:35

大麦网抢票脚本:Selenium 只管登录,下单走 requests 直连接口

大麦网抢票脚本&#xff1a;Selenium 只管登录&#xff0c;下单走 requests 直连接口 【免费下载链接】Automatic_ticket_purchase 大麦网抢票脚本 项目地址: https://gitcode.com/GitHub_Trending/au/Automatic_ticket_purchase 这是一个大麦网抢票脚本&#xff08;V2.…

作者头像 李华