
【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam18/beam点击查看免费下载Cloud Spanner 是 Google Cloud 提供的全托管关系型数据库服务具备全球规模的强事务一致性、ANSI 2011含扩展SQL 支持以及自动同步复制的高可用能力。Apache Beam 内置了 SpannerIO 连接器可同时作为批处理与流处理管线的数据源source和数据汇sink。本文将结合本仓库中 SpannerIO 的 Java、Python跨语言与 Go 实现源码系统讲解其支持范围、参数语义与底层机制帮助你直接上手编写可运行的 Beam Spanner 读写管线。一、Cloud Spanner 与 Apache Beam 的对接概览Cloud Spanner 的核心能力包括全局规模的事务一致性在跨区域复制的数据上提供强一致读事务关系模型与 SQL支持 ANSI 2011 标准含扩展可直接使用SELECT等标准 SQL自动同步复制为高可用性提供自动、同步的多副本复制无需应用层处理数据分片或一致性协调。Apache Beam 通过内置的SpannerIO 连接器提供对 Cloud Spanner 的读写能力其关键特性如下能力维度支持情况管线类型批处理batch与流处理streaming均支持角色既可作为 Source读也可作为 Sink写Java 入口org.apache.beam.sdk.io.gcp.spanner.SpannerIO位于 SpannerIO.javaPython 入口跨语言变换apache_beam.io.gcp.spanner位于 spanner.pyGo 入口beam/io/spannerio包位于 read.go二、读取 Cloud Spanner三种语言的实战写法2.1 Python 跨语言读取示例原文档核心示例Apache Beam 在 Python 侧通过**跨语言变换cross-language transform**接入 Java 实现的 SpannerIO。以下代码来自原文档从example_row表读取全部数据并逐行打印class ExampleRow(NamedTuple): id: int name: str with beam.Pipeline(optionsoptions) as p: output (p | Read from table ReadFromSpanner( project_idoptions.project_id, instance_idoptions.instance_id, database_idoptions.database_id, row_typeExampleRow, sqlSELECT * FROM example_row ) | Map Data Map(lambda row: fId {row.id}, Name {row.name}) | Log Data Map(logging.info))关键点row_type必须是一个NamedTuple用于声明返回行的 schema。在 spanner.py 的ReadFromSpanner中该类型会被序列化为 Beam schemanamed_tuple_to_schema(row_type).SerializeToString()随 URNbeam:transform:org.apache.beam:spanner_read:v1一起发送给 Java 扩展服务sql与table二选一sql执行任意查询table模式会按row_type中的全部字段自动生成SELECT列变换默认通过默认扩展服务展开。在 spanner.py 中说明了两种配置方式Option 1使用默认扩展服务Python SDK 自动下载或构建beam-sdks-java-io-google-cloud-platform-expansion-serviceshadowJar要求本机可执行java命令Option 2手动启动自定义扩展服务并通过expansion_service参数传入地址Flink Runner 的 Job Server 默认在 8097 端口暴露内建扩展服务。2.2 读取参数时间戳边界Timestamp Bound与批量模式ReadFromSpanner在 spanner.py 中暴露了timestamp_bound_mode、read_timestamp、staleness、time_unit、batching等参数语义如下timestamp_bound_mode定义只读事务或单次读/查询如何选择时间戳枚举TimestampBoundMode支持五类取值源码 spanner.pySTRONG在全部已提交事务可见的时间点执行读写最强一致性READ_TIMESTAMP在指定的read_timestamp时间点执行读写MIN_READ_TIMESTAMP选择至少晚于read_timestamp的时间点执行读写EXACT_STALENESS使用精确的过期时间staleness值时间戳在读取启动后不久选定MAX_STALENESS选择“最多过期staleness时间”的时间点执行读写staleness必须搭配time_unitTimeUnit枚举NANOSECONDS、MICROSECONDS、MILLISECONDS、SECONDS、HOURS、DAYS见 spanner.py使用read_timestamp仅在READ_TIMESTAMP/MIN_READ_TIMESTAMP模式下使用batching默认使用 Cloud Spanner 的Batch APIPartitionQuery并行分片读取当底层查询不可根分区non-root-partitionable时需设置batchingFalse。Java 侧对应的一致性保证在 SpannerIO.java 的Read consistency一节有更完整的说明读取变换保证基于只读事务的一致快照执行可通过withTimestampBound/withTimestamp控制数据新鲜度多个PCollection还能共享同一个事务——先用SpannerIO.createTransaction()惰性创建事务再通过withTransaction(tx)传入多个读取变换从而在单个一致快照上完成多表联合读取。2.3 Java 直接读取Query、Table 与 Index 三种模式Java 侧最直接的使用方式是SpannerIO.read()SpannerIO.java// 1) SQL 查询 PCollectionStruct rows p.apply( SpannerIO.read() .withInstanceId(instanceId) .withDatabaseId(dbId) .withQuery(SELECT id, name, email FROM users)); // 2) 整表读取 列裁剪 PCollectionStruct rows p.apply( SpannerIO.read() .withInstanceId(instanceId) .withDatabaseId(dbId) .withTable(users) .withColumns(id, name, email)); // 3) 通过二级索引读取 PCollectionStruct rows p.apply( SpannerIO.read() .withInstanceId(instanceId) .withDatabaseId(dbId) .withTable(users) .withIndex(users_by_name) .withColumns(id, name, email));注意read()默认走 PartitionQuery API 并行分片若查询不支持分区用withBatching(false)退化为非分区读取此时若数据量很大建议在输出后追加Reshuffle.viaRandomKey()让下游变换并行执行。此外SpannerIO.readAll()可对PCollectionReadOperation中的多个查询/表做一致批量读取但不应在流式管线中使用——同一只读事务只创建一次数据会过期且超过 1 小时无读取会被 Spanner 服务端自动关闭事务导致后续读取失败SpannerIO.java。2.4 Go 读取结构体 tag 驱动的类型安全读取Go SDK 的spannerio包提供Read与Query两个变换read.go// T 必须是导出的结构体且字段带有 spanner tag type User struct { ID int64 spanner:id Name string spanner:name } // 整表读取自动按 tag 推断列等价于投影查询 users : spannerio.Read(s, projects/p/instances/i/databases/d, users, reflect.TypeOf(User{})) // 自定义查询 users : spannerio.Query(s, projects/p/instances/i/databases/d, SELECT id, name FROM users, reflect.TypeOf(User{}))其实现通过structx.InferFieldNames(t, spannerTag)从结构体 tag 生成列清单再拼装SELECT语句默认同样使用 Spanner 的分区读取能力将结果拆分为多个 bundle若底层查询不可根分区可通过QueryOptionFn如UseBatching相关选项见 query_options.go关闭批量分区。三、写入 Cloud Spanner批处理与流处理的调优策略3.1 写入变换家族Python 侧通过WriteToSpannerSchema定义了一批写变换spanner.py变换语义对应 URNSpannerInsert插入主键冲突即失败beam:transform:org.apache.beam:spanner_insert:v1SpannerUpdate更新已存在行不存在则失败beam:transform:org.apache.beam:spanner_update:v1SpannerInsertOrUpdate存在则更新、不存在则插入最常用beam:transform:org.apache.beam:spanner_insert_or_update:v1SpannerReplace整体替换会删除未提供的列beam:transform:org.apache.beam:spanner_replace:v1SpannerDelete按主键删除输入行类型为键类型beam:transform:org.apache.beam:spanner_delete:v1一个典型的写入管线来自 spanner.py 的类文档示例class ExampleRow(NamedTuple): id: int name: str coders.registry.register_coder(ExampleRow, coders.RowCoder) with Pipeline() as p: _ ( p | Impulse beam.Impulse() | Generate beam.FlatMap(lambda x: range(num_rows)) | To row beam.Map(lambda n: ExampleRow(n, str(n))) | Write to Spanner SpannerInsertOrUpdate( instance_idyour_instance, database_idexisting_database, project_idyour_project_id, tableyour_table))Java 侧对应SpannerIO.write()接收PCollectionMutation.grouped()变体接收PCollectionMutationGroup以保证组内变更在同一事务提交SpannerIO.java。写变换返回SpannerWriteResult其中包含写失败的MutationGroup集合以及可用于批处理管线Wait.OnSignal的完成信号PCollection流式管线下该信号永不触发因为输入无界且其位于 GlobalWindow。3.2 批量Batching与分组Grouping机制写入时变换会把 Mutation 分组为批batch以减少发送给 Spanner 的事务数相关参数及默认值源码文档 spanner.py参数默认值含义max_batch_size_bytes10485761MB每批最多变更的字节数max_number_mutations5000每批最多变更的 cell 数max_number_rows500每批最多变更的行数grouping_factor1000按 key 排序进行批量选取的放大倍数排序占用 worker 本地内存过大易 OOMcommit_deadline15 秒Commit API 调用截止时间DEADLINE_EXCEEDED触发退避重试直到该上限max_cumulative_backoff900 秒15 分钟对DEADLINE_EXCEEDED累计退避上限超时后按failure_mode处理failure_modeFAIL_FAST失败即抛异常REPORT_FAILURES则继续处理错误使用可能造成数据丢失需要特别留意的约束SpannerIO.java单事务上限Cloud Spanner 单个事务最多 20000 个变更 cell含索引中的 cell。若索引较多并遇到INVALID_ARGUMENT: The transaction contains too many mutations异常需要调小MaxNumMutations三种配置模式分组 批量批处理管线默认按表与主键排序后成批写入吞吐最高但写入延迟也最高仅批量不分组withGroupingFactor(1)关闭分组流式管线默认凑满一个批即写入延迟与吞吐折中完全不批量withBatchSizeBytes(0)收到即写延迟最低内存与延迟权衡每个 worker 需要能容纳GroupingFactor × MaxBatchSizeBytes的 Mutation调大批大小时应相应调小分组因子分组批量在流式场景带来的延迟往往不可接受这正是流式默认禁用分组的原因。3.3 写入监控与 schema 就绪信号Java 写入变换提供了一系列监控计数器SpannerIO.javabatchable_mutation_groups/unbatchable_mutation_groups被批量写入与无法批量过大或范围删除而单独提交的变更组数量mutation_group_batches_received/..._write_success/..._write_failed批处理数量REPORT_FAILURES模式下失败批会被拆分、单个变更组分别重试mutation_groups_received/..._write_success/..._write_fail单个变更组处理计数spanner_write_success/spanner_write_fail/spanner_write_retries/spanner_write_timeouts写入成功、失败、重试与超时次数大量超时通常说明 Spanner 实例过载spanner_write_total_latency_ms写入总耗时毫秒。此外Write 变换在管线启动时会读取数据库 schema 以确定各表/索引的主键排序方式。如果同一管线内同时创建新表/索引会存在竞争条件导致 schema 读取早于建表完成、排序批量退化为非最优。此时应使用withSchemaReadySignal(PCollection)传入建表变换的输出信号借助Wait.OnSignal暂停 Write 直到 schema 就绪SpannerIO.java。3.4 事务语义边界务必知晓SpannerIO.write()不提供与 Cloud Spanner 相同的事务保证SpannerIO.java单个 Mutation 原子提交但所有 Mutation 并不在同一事务中提交每个 Mutation 至少应用一次at-least-once可能出现重复应用若管线意外停止已应用的 Mutation 不会回滚。需要将一小批变更捆绑进同一事务时使用MutationGroupSpannerIO.write().grouped()但要保证单个MutationGroup不超过 Spanner 事务上限。四、测试与快速上手资源仓库内为 SpannerIO 提供了完整的集成测试与性能测试可作为动手验证与二次开发的参考PythonReadFromSpanner与各写变换的集成/性能测试位于 spannerio_read_it_test.py、spannerio_read_perf_test.py、spannerio_write_it_test.py、spannerio_write_perf_test.pyGo读写与查询选项的单测位于 read_test.go、write_test.go、query_options_test.goJava连接器核心源码集中在 spanner 包除SpannerIO外还包括SpannerConfig配置封装、MutationGroup、ReadOperation、SpannerWriteResult、Transaction等核心类型以及changestreams子包Change Streams 读取支持。五、总结Apache Beam 对 Cloud Spanner 的支持是完整且多语言一致的Java 提供原生SpannerIOPython 通过跨语言变换复用同一套 Java 实现Go 提供 tag 驱动的类型安全读取三者在批处理与流式场景下均可作为 source/sink。读取侧的关键是理解 PartitionQuery 分区读取与时间戳边界TimestampBound的语义写入侧的关键则是围绕批量/分组因子、单事务 20000 cell 上限与至少一次投递语义进行吞吐与延迟的权衡。以本仓库的源码和测试为参照你可以快速搭建起生产可用的 Beam Spanner 数据处理管线。赞分享【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam18/beam点击查看免费下载相关推荐Apache Beam 与 Cloud Spanner 集成指南SpannerIO 连接器的读、写与变更流实践Apache Beam 与 Cloud Spanner 集成指南SpannerIO 连接器的读、写与变更流实践 Cloud Spanner 是 Google大数据批处理流处理数据工程Apache Beam 集成 Cloud Spanner 实战SpannerIO 跨语言读写、批量写入与一致性控制全解析Apache Beam 集成 Cloud Spanner 实战SpannerIO 跨语言读写、批量写入与一致性控制全解析 Cloud Spanner 是 Go大数据批处理流处理数据工程Apache Beam 实战使用 SpannerIO 将数据写入 Google Cloud Spanner 表Java 示例详解Apache Beam 实战使用 SpannerIO 将数据写入 Google Cloud Spanner 表Java 示例详解 本指南以 Apache大数据批处理流处理数据工程上一篇鸣潮自动化脚本解放双手的智能游戏助手终极指南下一篇终极指南深度解析MelonLoader初始化失败的3个关键步骤与解决方案创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考