- 大数据
- 批处理
- 流处理
- 数据工程
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
Apache Beam 2.27.0(发布于 2020-12-22)是一次同时包含新功能与改进的版本:它带来了可在管道运行时动态构造 BigQuery 读取请求的全新 transformReadAllFromBigQuery,正式发布了 Java 11 SDK 容器镜像,并为 MongoDB Atlas 场景补齐了读取能力。本文以官方发布说明为主线,结合当前仓库中的 Python/Java 源码与测试,逐项拆解 2.27.0 的新特性、I/O 改进与两条需要升级注意的破坏性变更,帮助读者在升级或迁移到 2.27.0 时做出准确判断。
1. 版本概览与获取方式
Apache Beam 2.27.0 的官方博客位于website/www/site/content/en/blog/beam-2.27.0.md,对应的下载入口见仓库中的下载页面(章节2.27.0 (2020-12-22)),其中提供官方源码包下载链接、SHA-512 校验值与签名文件,并回链到该发布博客。该版本的核心亮点可归纳为两点:
- Java 11 容器镜像随版本正式发布:所有 Beam 发布版现在都会同步发布 Java 11 的 SDK 容器镜像。
- 新增
ReadAllFromBigQuerytransform:允许在管道运行时接收多个 BigQuery 读取请求(表或查询),是动态/流式刷新场景的重要补充(对应 issue BEAM-9650)。
2. 新核心 transform:ReadAllFromBigQuery
2.1 它解决什么问题
传统的ReadFromBigQuery在管道构建阶段就需要确定读取哪张表、执行哪条查询。而ReadAllFromBigQuery的输入是ReadFromBigQueryRequest元素集合——这意味着读取请求可以在管道运行时由上游数据动态生成,例如由定时触发、消息驱动或经过beam.Map转换而来。
仓库中该 transform 的完整文档位于 bigquery.py,其基本用法如下:
read_requests = p | beam.Create([ ReadFromBigQueryRequest(query='SELECT * FROM mydataset.mytable'), ReadFromBigQueryRequest(table='myproject.mydataset.mytable')]) results = read_requests | ReadAllFromBigQuery()2.2 请求载体:ReadFromBigQueryRequest
每个读取请求由ReadFromBigQueryRequest定义,其构造参数(见 bigquery.py)为:
| 参数 | 类型 | 默认值 | 说明 |
|---|---|---|---|
query | str | None | 要执行的 SQL 查询;与table二选一 |
use_standard_sql | bool | True | 是否使用 BigQuery 标准 SQL 方言;设为False时使用 Legacy SQL。仅对查询输入生效,表输入时忽略 |
table | str或TableReference | None | 要读取的表 ID,需包含 project 与 dataset(形如'PROJECT:DATASET.TABLE') |
flatten_results | bool | False | 是否展平查询结果中的嵌套与重复字段 |
关键约束在validate()方法中强制执行:query与table必须且只能指定其一,二者同时给出或都为空都会抛出ValueError。此外,每个请求对象在创建时会生成一个基于时间戳与随机 token 的内部对象 ID(obj_id),用于生成 BigQuery 导出目录和作业名,保证并发请求互不冲突。
2.3 流式管道中的典型应用:定时刷新 side input
该 transform 官方文档中给出的“杀手级场景”是:在流式管道中周期性地重新拉取 BigQuery 数据作为 side input。示例(见 bigquery.py):
side_input = ( p | 'PeriodicImpulse' >> PeriodicImpulse(first_timestamp, last_timestamp, interval, True) | 'MapToReadRequest' >> beam.Map( lambda x: ReadFromBigQueryRequest(table='dataset.table')) | beam.io.ReadAllFromBigQuery()) main_input = ( p | 'MpImpulse' >> beam.Create(sample_main_input_elements) | 'MapMpToTimestamped' >> beam.Map(lambda src: TimestampedValue(src, src)) | 'WindowMpInto' >> beam.WindowInto(window.FixedWindows(main_input_windowing_interval))) result = ( main_input | 'ApplyCrossJoin' >> beam.FlatMap( cross_join, rights=beam.pvalue.AsIter(side_input)))其核心思想是:用PeriodicImpulse周期性产生元素,经Map转换成ReadFromBigQueryRequest,再由ReadAllFromBigQuery拉取数据,最终以AsIter形式作为另一路主输入的 side input,实现类似“维度表定时刷新 + 流式事实表关联”的宽表/维表 join 模式。
2.4 源码实现剖析
从实现上看(bigquery.py),ReadAllFromBigQuery接受以下构造参数:
| 参数 | 默认值 | 说明 |
|---|---|---|
gcs_location | None | 导出表数据的 GCS 桶路径;为None时回退到temp_location |
validate | False | 初始化时是否做各种检查(如表是否存在);导出方式较慢时可关闭 |
kms_key | None | 实验性:创建临时表时使用的 Cloud KMS 密钥 |
temp_dataset | None | 临时数据集 |
bigquery_job_labels | None | 附加到导出作业的标签 |
query_priority | BATCH | 查询优先级 |
其执行链路分为三步:
- 对每个请求运行
_BigQueryReadSplit(ParDo),通过导出作业将表快照导出到 GCS,同时产出“待清理位置”(location_to_cleanup)旁路输出; - 用
SDFBoundedSourceReader读取导出的文件; - 通过
_PassThroughThenCleanup在读取完成后尝试清理导出的临时文件。
从代码注释可以确认两点重要限制:每次导出会使用请求对象里的 UUID 在 GCS 上新建子目录;该 transform仅受 Portable 与 Dataflow v2 runner 支持,且当前版本不会自动清理执行期间创建的临时数据集(见 issue BEAM-11359),官方注释也建议不要在 GlobalWindow 上用于流式作业,以免快照无法清理。这些限制在选用该 transform 时需纳入容量与成本评估。对应测试可参考 bigquery_test.py 与 bigquery_read_it_test.py。
3. I/O 改进:MongoDB
3.1 支持 MongoDB Atlas(Python)
2.27.0 起,Python 的ReadFromMongoDB可用于 MongoDB Atlas(issue BEAM-11266)。仓库中 mongodbio.py 的模块文档给出了明确解释:
- 传统分片读取依赖 MongoDB 的
splitVector命令,这是一个高权限命令,Atlas 不允许赋予任何用户; - 因此需要开启
bucket_auto=True,改用 MongoDB 聚合管道中的@bucketAuto阶段完成分片定位。
Atlas 场景用法示例:
pipeline | ReadFromMongoDB(uri='mongodb+srv://user:pwd@cluster0.mongodb.net', db='testdb', coll='input', bucket_auto=True)对应的 Java 实现同样支持该选项:Read.withBucketAuto(boolean)(见 MongoDbIO.java),其内部正是二选一逻辑——启用bucketAuto时走$bucketAuto聚合,否则构造splitVector命令获取 split keys(见同文件#L508-L536)。
3.2 display_data 中的密码掩码
针对密码泄露风险,2.27.0 为 Python 的ReadFromMongoDB/WriteToMongoDB增加了 display_data 密码掩码处理(issue BEAM-11444)。从当前源码的display_data()实现(mongodbio.py)可以看到,返回的展示数据仅包含database、collection、filter、projection、bucket_auto等字段,并不包含 uri 或密码字段,从而避免连接串中的明文口令被透传到管道展示/监控信息中。
4. 新特性与改进
4.1 Hadoop 3 兼容性验证
2.27.0 开始对依赖 Hadoop 的 Beam 模块做 Hadoop 3 兼容性测试(issue BEAM-8569;Hive/HCatalog 模块当时仍待验证)。仓库中 hadoop-common/build.gradle 的版本矩阵清晰反映了这一点,其中既包含 Hadoop 2.10.2,也包含 3.2.4 与 3.3.6:
def hadoopVersions = [ "2102": "2.10.2", "324": "3.2.4", "336": "3.3.6", // "341": "3.4.1", // tests already exercised on the default version ]可以推断:该模块采用多版本配置(configurations.create("hadoopVersion$kv.key"))分别挂载不同 Hadoop 版本进行测试,兼顾老版本兼容与新版迁移验证。
4.2 Java 11 SDK 容器镜像进入发布流程
发布 Java 11 SDK 容器镜像已作为 Apache Beam 发布流程的一部分被正式支持(issue BEAM-8106),这也是“Highlights”中“Java 11 Containers 随所有发布版发布”的工程基础。
4.3 Beam SQL 新增 Cloud Bigtable Provider
2.27.0 为 Beam SQL 添加了 Cloud Bigtable Provider 扩展(issues BEAM-11173、BEAM-11373),使 Bigtable 可以像其他表一样通过 Beam SQL 查询/写入。当前仓库中该提供器位于 bigtable 包,核心类包括BigtableTableProvider(表提供器注册)、BigtableTable(表元数据与读写实现)与BigtableFilter(谓词下推过滤),并有BigtableTableWithRowsTest、BigtableFilterTest等测试覆盖读写与过滤行为。
4.4 Thrift 数据的 Schema Provider
本版本新增了针对 thrift 数据的 schema provider(issue BEAM-11338),让 Beam SQL/ schema 体系可以识别 thrift 定义的数据结构。仓库中 thrift 相关模块位于 sdks/java/io/thrift,其测试资源(如thrift_test.thrift、payload.thrift)用于验证 thrift 编解码与 schema 推导逻辑。
4.5 Dataflow runner 的 Combiner Packing 优化
2.27.0 为 Dataflow runner 增加了 combiner packing 管道优化(issue BEAM-10641)。其底层逻辑可以在 DataflowPipelineTranslator.java 中看到:翻译器会判断allowCombinerLifting(是否允许将 Combine 操作上提/合并),其中不仅检查 transform 是否fewKeys()(少量 key 场景),还通过TriggerCombinerLiftingCompatibility(见同文件#L149-L206)遍历触发器树检查兼容性——例如 count 类触发器不支持 combiner lifting。最终通过DISALLOW_COMBINER_LIFTING属性传递给 Dataflow 服务端,决定是否启用该优化。理解这一点有助于评估:当管道使用非兼容触发器或大量 key 时,combiner 不会被合并,行为与预期一致。
4.6 新增 Kafka → Pub/Sub 摄取示例
本版本新增了将 Apache Kafka 数据摄取到 Google Pub/Sub 的示例(issue BEAM-11065)。仓库的 Java 示例目录 examples/java 下存在KafkaStreaming、KafkaWordCountJson、RateLimitedPubSubReader等 Kafka/Pub/Sub 相关示例,可直接作为从 Kafka 消费并写入 Pub/Sub 的参考起点。
5. 破坏性变更与升级注意
5.1 HBaseIO:hbase-shaded-client 依赖改为由用户提供
从 2.27.0 起,HBaseIO不再打包传递hbase-shaded-client,该依赖需要由使用方自行提供(issue BEAM-9278)。这意味着升级后,凡是使用HBaseIO的项目必须在构建文件中显式声明 HBase shaded client 依赖,否则会出现运行时类缺失错误。这是本次升级最需要优先处理的兼容性问题。
5.2 AWS:--region 正式更名为 --awsRegion
amazon-web-services2模块的--region参数被替换为--awsRegion(issue BEAM-11331)。仓库中 AwsOptions.java 是这一变更的落点:getAwsRegion()/setAwsRegion()成为标准 pipeline option,其默认值工厂AwsRegionFactory会尝试通过 AWS SDK 的DefaultAwsRegionProviderChain自动加载默认区域,加载失败时返回null。ClientBuilderFactory在创建各 AWS 服务客户端时统一解析该区域(见 ClientBuilderFactory.java),并对缺失区域做显式校验。模块的 build.gradle 中也已出现'--awsRegion=us-west-2'的用法示例。
因此升级时需要:
- 将命令行参数
--region=...改写为--awsRegion=...; - 检查代码中是否通过旧名设置区域,统一改用
AwsOptions#setAwsRegion。
6. 贡献者
根据git shortlog统计,2.27.0 版本共有约 80 位贡献者参与,包括 Ahmet Altay、Alexey Romanenko、Brian Hulette、Chamikara Jayalath、Emily Ye、Kenneth Knowles、lostluck、Maximilian Michels、Pablo Estrada、Robert Bradshaw、Robert Burke、Reuven Lax、Valentyn Tymofieiev、Udi Meiri 等(完整名单见发布博客原文)。若你计划向 Beam 贡献代码,可参考仓库根目录的 CONTRIBUTING.md 与 CI.md 了解流程与验证要求。
7. 总结与升级建议
Apache Beam 2.27.0 的核心价值在于把 BigQuery 读取从“构建期静态定义”推进到“运行时动态请求”,ReadAllFromBigQuery+ReadFromBigQueryRequest为流式 side input 刷新、多表批量拉取提供了官方范式(注意其仅支持 Portable/Dataflow v2,且临时数据集需自行治理)。配合 MongoDB Atlas 读取支持、Hadoop 3 兼容矩阵与 Dataflow combiner packing 优化,这一版本在“云上动态数据摄取”方向上迈出了重要一步。
升级清单建议如下:
- 为使用
HBaseIO的模块显式补充hbase-shaded-client依赖; - 全局检索
--region并替换为--awsRegion; - 评估
ReadAllFromBigQuery是否适用于自身 runner(Portable / Dataflow v2)与流式场景,并规划临时导出目录的清理; - 使用 MongoDB Atlas 时开启
bucket_auto=True,并确认展示数据中不再包含连接串口令; - 若依赖 Hadoop 生态,可参考
hadoop-common的多版本矩阵规划 Hadoop 3 迁移。
- 大数据
- 批处理
- 流处理
- 数据工程
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
相关推荐
Meteor 2.11 版本升级指南:MongoDB 6.x 支持、驱动升级与破坏性变更全解析
Meteor 2.11 版本升级指南:MongoDB 6.x 支持、驱动升级与破坏性变更全解析 导读 :本文以 Meteor 仓库 v2.11.0 变更日志 h
后端前端开发工具移动开发Apache Beam 2.36.0 版本深度解析:Kafka 停止读取时间、cloudpickle 序列化与破坏性变更全指南
Apache Beam 2.36.0 版本深度解析:Kafka 停止读取时间、cloudpickle 序列化与破坏性变更全指南 Apache Beam 2.36
大数据批处理流处理数据工程Apache Beam 2.23.0 版本深度解读:Twister2 Runner、Python 3.8 支持与 I/O 扩展全解析
Apache Beam 2.23.0 版本深度解读:Twister2 Runner、Python 3.8 支持与 I/O 扩展全解析 本文基于 Apache B
大数据批处理流处理数据工程
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考