news 2026/10/9 5:11:54

Apache Beam 2.27.0 版本深度解析:ReadAllFromBigQuery 动态读取、MongoDB Atlas 支持与升级变更指南

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Apache Beam 2.27.0 版本深度解析:ReadAllFromBigQuery 动态读取、MongoDB Atlas 支持与升级变更指南
  • 大数据
  • 批处理
  • 流处理
  • 数据工程

【免费下载链接】beam

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

项目地址:https://gitcode.com/gh_mirrors/beam4/beam
点击查看免费下载

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)为:

参数类型默认值说明
querystrNone要执行的 SQL 查询;与table二选一
use_standard_sqlboolTrue是否使用 BigQuery 标准 SQL 方言;设为False时使用 Legacy SQL。仅对查询输入生效,表输入时忽略
tablestr或TableReferenceNone要读取的表 ID,需包含 project 与 dataset(形如'PROJECT:DATASET.TABLE')
flatten_resultsboolFalse是否展平查询结果中的嵌套与重复字段

关键约束在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_locationNone导出表数据的 GCS 桶路径;为None时回退到temp_location
validateFalse初始化时是否做各种检查(如表是否存在);导出方式较慢时可关闭
kms_keyNone实验性:创建临时表时使用的 Cloud KMS 密钥
temp_datasetNone临时数据集
bigquery_job_labelsNone附加到导出作业的标签
query_priorityBATCH查询优先级

其执行链路分为三步:

  1. 对每个请求运行_BigQueryReadSplit(ParDo),通过导出作业将表快照导出到 GCS,同时产出“待清理位置”(location_to_cleanup)旁路输出;
  2. 用SDFBoundedSourceReader读取导出的文件;
  3. 通过_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 优化,这一版本在“云上动态数据摄取”方向上迈出了重要一步。

升级清单建议如下:

  1. 为使用HBaseIO的模块显式补充hbase-shaded-client依赖;
  2. 全局检索--region并替换为--awsRegion;
  3. 评估ReadAllFromBigQuery是否适用于自身 runner(Portable / Dataflow v2)与流式场景,并规划临时导出目录的清理;
  4. 使用 MongoDB Atlas 时开启bucket_auto=True,并确认展示数据中不再包含连接串口令;
  5. 若依赖 Hadoop 生态,可参考hadoop-common的多版本矩阵规划 Hadoop 3 迁移。
  • 大数据
  • 批处理
  • 流处理
  • 数据工程

【免费下载链接】beam

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

项目地址:https://gitcode.com/gh_mirrors/beam4/beam
点击查看免费下载

相关推荐

上一篇:7步精通INAV飞控:从零搭建到精准导航的完整指南
下一篇:八大网盘下载加速终极指南:一键获取真实直链告别限速烦恼

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

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

2026年AP组网设备清单:从选型到部署的完整指南

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/10/9 5:10:46

linux中find查找

linux常用命令 find查找 find 查找范围 匹配条件(范围要尽量小,这样查找起来才快) #匹配条件: -name: 按照文件的名称-type: 文件类型(l,d,f)-size: 文件大小 &#xff…

作者头像 李华
网站建设 2026/10/9 5:08:58

GRE备考作业化:从目标拆解到每日清单的高效执行方案

1. 把GRE备考当成“作业”来经营:从目标到任务的翻译过程第一次翻开GRE官方指南的人,十有八九会和我当初一样,在目录面前坐半小时不动笔。整本书的章节、题型、评分规则铺在眼前,那种感觉不是“难”,而是“不知道自己该…

作者头像 李华
网站建设 2026/10/9 5:08:49

5G OTA测试全解析:从空口测量原理到暗室搭建与波束验证

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/10/9 5:08:28

自托管AI助手实战:Docker+NAS部署多Agent协作与定时任务

1. 从标题拆解这个自托管AI助手的真实价值1.1 这个项目到底解决了什么问题第一次看到这个标题的时候,我脑子里冒出来的第一个念头是:又一个套壳聊天界面?但仔细拆开看,它其实踩中了三个很实际的需求点。第一是自托管,数…

作者头像 李华