- 大数据
- 批处理
- 流处理
- 数据工程
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
Apache Beam 2.21.0 是 2020 年 5 月 27 日发布的重要版本(发布日期见 beam-2.21.0.md),它在 Python SDK 的类型系统、Java SDK 的 Schema 体系以及批/流 I/O 连接器三个方向同时迈出了一大步。本文将逐条解析该版本的 I/O 变更、新特性、破坏性变更与弃用项,并结合当前仓库源码验证每个关键点的真实实现,帮助你在升级或迁移管线时准确评估影响面、掌握新 API 的正确用法。
一、版本总览与发布要点
2.21.0 的官方发布说明强调"包含改进与新功能"("This release includes both improvements and new functionality"),其核心亮点集中在:
- Python SDK:原生 Python 3 类型注解全面接管管线类型提示(type hints);废弃的 Datastore v1 模块被移除;BigQuery 批量写入支持 Avro 文件加载。
- Java SDK:Beam Schema 引入全新的 Options 概念(取代旧的 FieldType metadata);protobuf 扩展完全 Schema-aware;新增 Google Cloud AI 视频智能与自然语言处理 API 集成。
- 构建与发布:引入
docker-pull-licenses标签,可将第三方依赖的 License/Notice 打入 Docker 镜像。 - 破坏性变更:Dataflow Runner 强制要求
--region;HBaseIO.ReadAll输入类型变化;ProcessContext.updateWatermark被移除;Row 对象的 Coder 推断被禁用。
需要注意的是,该版本距今较久,文中涉及的 API 部分已在后续版本中进一步演进;本文所有结论均以当前仓库中对应源码为据,若你的项目仍基于 2.21.x 使用,可直接对照验证。
二、Python SDK:类型注解成为一等公民
2.1 背景:从装饰器到类型注解
2.21.0 之前,Python SDK 的管线类型提示主要依赖@beam.typehints.with_input_types/@beam.typehints.with_output_types等显式装饰器。本版本(对应 PR #10717 的变更)开始默认将 Python 3 的函数类型注解(type annotations)作为管线类型提示来源,这意味着你可以直接写:
import apache_beam as beam class SplitWordsFn(beam.DoFn): def process(self, element: str) -> list: return element.split()process()的element: str与返回注解-> list会自动被 Beam 的类型系统识别,用于推导 PCollection 的 Coder、校验上下游类型一致性,而无需再叠加装饰器。
2.2 回退与局部禁用机制
官方发布说明明确提示:如果怀疑该特性导致管线失败,可以在创建管线之前调用apache_beam.typehints.disable_type_annotations()完全关闭它;对特定函数(如process())则用@apache_beam.typehints.no_annotations装饰来单独禁用。
仓库源码印证了这两个 API 的实现位置与机制(sdks/python/apache_beam/typehints/decorators.py):
no_annotations(fn)通过setattr(fn, '_beam_no_annotations', True)给函数打上标记(见该文件 L170-L173);disable_type_annotations()则从模块层面全局关闭(见 L177 附近);- 在 L330 处的类型推断逻辑中会同时检查全局开关
_disable_from_callable与函数级标记getattr(fn, '_beam_no_annotations', False),命中任一即跳过注解推断。
对应测试 sdks/python/apache_beam/typehints/decorators_test.py 中的test_disable_type_annotations、test_no_annotations_on_same_function、test_no_annotations_on_diff_function等用例,覆盖了全局关闭、同函数/不同函数局部禁用等场景,可作为回归参考。
使用建议:
import apache_beam as beam from apache_beam.typehints import disable_type_annotations, no_annotations # 全局关闭(放在构建管线前) # disable_type_annotations() class LegacyDoFn(beam.DoFn): @no_annotations # 仅此函数不参与注解推断 def process(self, element): return [element]更多细节可参考仓库文档 python-type-safety 相关说明(路径见网站源码目录)。
三、Python I/O:Datastore v1 移除与 Spanner 批量写入增强
3.1 Datastore:v1模块移除,迁移至v1new
发布说明指出:apache_beam.io.gcp.datastore.v1模块因依赖的客户端过旧且不支持 Python 3,在 2.21.0 中被移除(对应 BEAM-9529)。迁移路径为apache_beam.io.gcp.datastore.v1new.datastoreio。
当前仓库中v1new包依然完整存在(sdks/python/apache_beam/io/gcp/datastore/v1new/),包含:
datastoreio.py:核心读写 Transform;datastore_write_it_pipeline.py/datastore_write_it_test.py:批量写入集成测试;datastoreio_test.py、query_splitter.py、rampup_throttling_fn.py等辅助模块。
典型的迁移后用法:
from apache_beam.io.gcp.datastore.v1new.datastoreio import ReadFromDatastore, WriteToDatastore from google.cloud.datastore import client迁移时重点核对三点:实体 key 的构造方式、查询构造 API、以及返回类型是否从旧 proto 实体变为新客户端实体。仓库中 datastore_wordcount 示例 之外的 Python 侧示例可在sdks/python/apache_beam/examples/cookbook/目录下查找datastore_wordcount.py作为参考。
3.2 Spanner:新增集成测试与批量写入更新
发布说明提及 Python SDK 为 Google Cloud Spanner Transform 新增了集成测试,并更新了批量写入(batch write)功能(对应 BEAM-8949)。这说明 Spanner 连接器的生产就绪度在该版本得到提升,尤其在大批量写入场景下,批量 API 相比逐条写入能显著减少 RPC 开销。
3.3 BigQuery:Avro 文件加载(File Loads)
2.21.0 为 Python 的 BigQuery 写入新增了通过 Avro 文件加载(file loads)写入 BigQuery的能力(对应 BEAM-8841)。
发布说明的核心要点:
- 默认仍是 JSON 文件加载,但可通过
temp_file_format参数切换为 AVRO; - Avro 加载的原理:将 Python 类型导出为 Avro 类型后再批量导入 BigQuery,因此切换后需要把 JSON 兼容类型(字符串形式的日期/时间戳、用字符串表示的大数值)改为 Python 原生类型(
date、datetime、decimal等); - 文件加载方式适合大吞吐批量写入,相比流式插入更经济高效。
当前仓库源码完整保留了这一参数链(sdks/python/apache_beam/io/gcp/bigquery.py):
- L2052 处
WriteToBigQuery签名中定义了temp_file_format=None; - L2206 处注释明确说明"
temp_file_format: The format to use for file loads into BigQuery"; - L2289 处默认值落地:
self._temp_file_format = temp_file_format or bigquery_tools.FileFormat.JSON; - L2407-L2412 处理 AVRO 分支,并提示若不匹配 JSON 则需显式指定
temp_file_format="NEWLINE_DELIMITED_JSON"。
同目录的 bigquery_file_loads.py 实现了文件加载的核心流水线(L966-L1291 多处将self._temp_file_format传递给文件写出与导入步骤),而 bigquery_test.py 与 bigquery_file_loads_test.py 中大量temp_file_format=bigquery_tools.FileFormat.AVRO的断言验证了 Avro 路径的行为。
实操示例:
import apache_beam as beam from apache_beam.io.gcp.bigquery import WriteToBigQuery from apache_beam.io.gcp.bigquery_tools import FileFormat with beam.Pipeline() as p: (p | beam.Create([{'dt': __import__('datetime').date(2020, 5, 27), 'amount': __import__('decimal').Decimal('12345.67')}]) | WriteToBigQuery( table='your-project:your-dataset.your_table', write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND, create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED, temp_file_format=FileFormat.AVRO, # 切换为 Avro 文件加载 method='FILE_LOADS'))注意:使用 AVRO 时,日期/时间戳必须使用date/datetime/decimal等 Python 原生类型,而不是字符串形式。
四、Java SDK:Beam Schema Options 与 protobuf 扩展
4.1 Schema Options:替代 FieldType metadata 的强类型机制
发布说明指出(对应 BEAM-9035):Java SDK 在 Beam Schema 中引入Options概念,为字段(Field)和整个 Schema 提供额外上下文,取代原先仅存在于FieldType中的 Beam metadata。Options完全类型化,甚至可以包含复杂的 Row 结构。
仓库源码印证了 Options 的完整落位(sdks/java/core/src/main/java/org/apache/beam/sdk/schemas/Schema.java):
- L365-L370:
Schema支持.withOptions(...)追加选项并生成新实例; - L1128:
Field抽象类中定义public abstract Options getOptions(); - L1483:
Schema自身也提供getOptions(); - SchemaTranslation.java(L101-L152 与 L316-L340)展示了 Schema/Field 的 Options 与 proto 表示之间的互转;
- SchemaUtils.java(L366-L374)在
toPrettyString输出中渲染fieldOptions与schemaOptions,便于调试。
典型使用方式:
import org.apache.beam.sdk.schemas.Schema; import org.apache.beam.sdk.schemas.Schema.Options; Schema.Options fieldOpts = Schema.Options.builder() .setOption("description", "订单金额") .setOption("nullable", true) .build(); Schema.Field field = Schema.Field.of("amount", Schema.FieldType.DECIMAL) .withOptions(fieldOpts); Schema schema = Schema.builder() .addField(field) .build() .withOptions(Schema.Options.builder().setOption("source", "orders_topic").build());注意:发布说明明确标注Schema aware 仍属实验性(experimental),Options 替代 metadata 的完整切换要到 2.23.0 才完成(见下文 Deprecations)。
4.2 protobuf 扩展全面 Schema-aware
发布说明指出(对应 BEAM-9044):protobuf 扩展已完全 Schema-aware,并支持将 protobuf 选项(custom options)转换为 Beam Schema Options。这意味着你可以直接对 proto 消息生成的类使用 Schema-aware 转换(例如通过ProtoCoder配合 Schema 相关 Transform 或 SQL),而无需手工映射字段。
仓库中的实现证据(sdks/java/extensions/protobuf/):
- ProtoSchemaTranslator.java:负责 proto
Descriptor↔ BeamSchema的双向翻译(getFieldNumber、isNullable等辅助逻辑可见于 ProtoBeamConverter.java L120-L127 与 L400 附近); - ProtoByteBuddyUtils.java:通过 Byte Buddy 生成 Row ↔ proto 消息的转换逻辑;
- ProtoSchemaLogicalTypes.java:将 proto 的 uint32、sint64、fixed64 等标量映射为 Beam LogicalType。
同样,该特性标注为实验性,生产使用需做好兼容性评估。
五、Java I/O 与云服务集成
5.1 Google Cloud AI 集成:视频智能与自然语言处理
2.21.0 为 Java SDK 新增两个 Google Cloud AI 服务集成:
- VideoIntelligence(视频智能,对应 BEAM-9147);
- Natural Language Processing API(自然语言处理,对应 BEAM-9634)。
这使 Beam 管线可以在批/流场景中直接调用 AI 服务进行视频内容分析(如镜头检测、标签标注)与文本 NLP(如实体识别、情感分析),是当时云 AI 能力与 Beam 统一编程模型结合的代表性实践。相关代码位于sdks/java/io/google-cloud-platform/的 GCP 连接器族中。
5.2docker-pull-licenses:镜像内嵌第三方许可证
发布说明引入docker-pull-licenses标签(对应 BEAM-9136):构建 Docker 镜像时若设置该标签,第三方依赖的 License/Notice 会被写入镜像内/opt/apache/beam/third_party_licenses/目录;默认不写入。这解决了容器化部署时开源合规(License 随镜像分发)的需求,对需要对外分发容器的团队尤为实用。
六、破坏性变更(Breaking Changes)逐条核对
6.1 Dataflow Runner 强制--region
自 2.21.0 起,Dataflow Runner 要求必须设置--region选项,除非环境中已配置默认值(对应 BEAM-9199)。这源于 Dataflow 服务按区域端点(regional endpoints)管理资源与配额。
仓库源码印证(runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/DataflowRunner.java):
- L375 处校验逻辑:缺失时向
missing列表加入"region",随后触发参数校验失败; - DataflowPipelineOptions.java L145-L155 定义了 region 选项的注释与 setter。
因此,升级到 2.21+ 后所有 Dataflow 任务必须显式带上区域,例如:
mvn compile exec:java \ -Dexec.mainClass=com.example.MyPipeline \ -Dexec.args="--runner=DataflowRunner \ --project=my-project \ --region=us-central1 \ --stagingLocation=gs://my-bucket/staging"6.2HBaseIO.ReadAll输入类型变更
发布说明指出:HBaseIO.ReadAll现在要求输入PCollection<HBaseIO.Read>(而非此前的HBaseQuery对象),对应 BEAM-9279。
仓库源码印证(sdks/java/io/hbase/src/main/java/org/apache/beam/sdk/io/hbase/HBaseIO.java):
- L404-L405:
public static class ReadAll extends PTransform<PCollection<Read>, PCollection<Result>>,类型参数明确为PCollection<Read>→PCollection<Result>。
迁移示例:
// 旧:PCollection<HBaseQuery> // 新:先构造 HBaseIO.Read PCollection<HBaseIO.Read> reads = p.apply(Create.of( HBaseIO.read().withConfiguration(conf).withTable("table1"), HBaseIO.read().withConfiguration(conf).withTable("table2"))); PCollection<Result> results = reads.apply(HBaseIO.readAll());6.3ProcessContext.updateWatermark移除
发布说明:ProcessContext.updateWatermark被移除,改用WatermarkEstimator(对应 BEAM-9430)。这是对 DoFn 自主水印管理能力的重构:状态处理型 DoFn 应通过WatermarkEstimator(配合@DoFn.WatermarkEstimator状态绑定)来汇报输出水印,而非直接调用updateWatermark。这属于 API 层面的行为变更,涉及自定义有状态 DoFn 的都需要改。
6.4 Row 对象 Coder 推断禁用
发布说明:PCollection of Row 的 Coder 推断被禁用(对应 BEAM-9569)。即 Beam 不再自动为Row类型的 PCollection 猜测 Coder(因为 Row 需要 Schema 才能编码,而 Schema 无法凭空推断),使用者必须显式指定,例如:
PCollection<Row> rows = input.apply(...).setCoder(RowCoder.of(schema));6.5 Go SDK Docker 镜像暂停发布
发布说明:Go SDK 的 Docker 镜像暂时停止发布("until further notice")。如果 CI/CD 依赖 Go SDK 官方容器镜像,需要评估替代方案或暂缓升级。
七、弃用项(Deprecations)
7.1FieldType.getMetadata弃用
发布说明:FieldType.getMetadata被弃用,由 Schema Options 取代,并将在2.23.0中移除(对应 BEAM-9704)。任何基于 metadata 的代码都应尽快迁移到 4.1 节介绍的Options体系。
7.2 Dataflow--zone弃用
发布说明:Dataflow Runner 的--zone选项弃用,改用--worker_zone(对应 BEAM-9716)。
仓库源码佐证(DataflowRunner.java L565-L576):worker_region与workerRegion、workerZone等实验之间存在互斥校验,说明 worker 级区域/可用区配置已成为新的规范入口。
# 旧:--zone=us-central1-a # 新:--worker_zone=us-central1-a八、升级与迁移检查清单
基于以上分析,从 2.20.x(或更早)升级到 2.21.0 时建议按以下顺序排查:
Python 侧
- 若使用
apache_beam.io.gcp.datastore.v1,改为v1new并重写实体/查询构造; - 运行测试集,若类型注解推断导致失败,用
disable_type_annotations()或@no_annotations回退; - BigQuery 大批量写入可评估
temp_file_format=FileFormat.AVRO的收益,同时把字符串日期/大数切换为原生 Python 类型。
- 若使用
Java 侧
- 将
FieldType.getMetadata迁移到 SchemaOptions(2.23.0 前完成); - 有状态 DoFn 若调用过
updateWatermark,改用WatermarkEstimator; RowPCollection 显式设置RowCoder;HBaseIO.ReadAll输入改为PCollection<HBaseIO.Read>;- Dataflow 任务补充
--region,--zone改--worker_zone。
- 将
发布与合规
- 需要随镜像分发第三方 License 时,在构建 Docker 镜像时设置
docker-pull-licenses标签; - Go SDK 用户关注 Docker 镜像恢复发布的公告。
- 需要随镜像分发第三方 License 时,在构建 Docker 镜像时设置
九、参考资源
- 官方发布说明原文:beam-2.21.0.md(本文所有变更条目均可在此追溯)
- Python 类型注解机制实现:decorators.py 与对应测试 decorators_test.py
- BigQuery Avro 文件加载:bigquery.py、bigquery_file_loads.py
- Datastore v1new 包:v1new/
- Java Schema Options:Schema.java、SchemaTranslation.java
- protobuf Schema-aware 扩展:sdks/java/extensions/protobuf/
- Dataflow region 校验:DataflowRunner.java、DataflowPipelineOptions.java
- HBaseIO.ReadAll 类型签名:HBaseIO.java
Apache Beam 2.21.0 是 Python 类型系统现代化与 Java Schema 体系重构进程中的关键节点,本文提供的源码级证据与迁移清单,可帮助你在实际升级中做到有的放矢、风险可控。
- 大数据
- 批处理
- 流处理
- 数据工程
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
相关推荐
Apache Beam 2.18.0 版本全解析:Spark Structured Streaming Runner、SQS/RabbitMQ I/O 与 SQL 能力升级
Apache Beam 2.18.0 版本全解析:Spark Structured Streaming Runner、SQS/RabbitMQ I/O 与 SQ
大数据批处理流处理数据工程Apache Beam 2.11.0 版本解析:依赖升级、新 I/O 能力与运行时改进全指南
Apache Beam 2.11.0 版本解析:依赖升级、新 I/O 能力与运行时改进全指南 Apache Beam 2.11.0 是该项目于 2019 年发布
大数据批处理流处理数据工程Apache Beam 2.48.0 版本解析:Experimental 注解移除、Kinesis 增强扇出与 Go SDK I/O 新能力
Apache Beam 2.48.0 版本解析:Experimental 注解移除、Kinesis 增强扇出与 Go SDK I/O 新能力 Apache Be
大数据批处理流处理数据工程
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考