【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
Apache Beam 的统一编程模型覆盖批处理与流处理,其发布说明文件 CHANGES.md 记录了自 2.19.0 以来每一个版本的高光特性、I/O 连接器变化、破坏性变更、安全修复与已知问题,是理解项目技术演进、规划升级路线和排查兼容性问题的第一手权威资料。读完本文,你将能够快速解读 Beam 版本号背后的技术脉络,掌握各版本间的迁移要点(如 Avro 模块迁移、Runner v2 强制化、Go SDK 模块化改造),并知道如何在仓库源码中逐一验证这些变更。
CHANGES.md 是什么:一份结构化的版本发布档案
CHANGES.md 位于仓库根目录,采用"按版本号倒序、每个版本按固定分类"的组织方式。文件头部保留了官方发布说明模板,定义了每个版本条目必须覆盖的八大分类:
- Highlights:该版本最值得关注的能力突破
- I/Os:新增或改进的连接器(Source/Sink)
- New Features / Improvements:SDK 与 Runner 的功能增强
- Breaking Changes:破坏性变更,通常伴随迁移指引
- Deprecations:即将移除的功能预告
- Bugfixes:缺陷修复
- Security Fixes:CVE 安全修复
- Known Issues:已知问题及规避方案
模板中预留的[#X]、[BEAM-X]占位符表明:每个条目都应关联对应的 GitHub PR/Issue 或 JIRA 工单,方便追溯具体讨论与实现细节。例如 2.55.0(Unreleased)中的 Enrichment Transform 关联了 PR #30001。
从版本时间线看,Beam 保持约每月一次的发布节奏(2.19.0 于 2020-01-31,2.53.0 于 2024-01-04),且仓库中存在 2.54.0(Cut, Unreleased)与 2.55.0(Unreleased)两个尚未正式发布的候选条目,说明该文件是发布流程的"活文档"。
版本节奏与语言运行环境支持基线
CHANGES.md 逐版本记录了 SDK 运行环境的支持变化,这是评估升级成本的关键依据:
| 版本 | 运行环境变更 |
|---|---|
| 2.20.0 | 最后一次完整支持 Python 2 / 3.5 之前的过渡期 |
| 2.24.0 | 明确为"最后一个支持 Python 2 和 Python 3.5 的版本" |
| 2.25.0 | 正式移除 Python 2 与 Python 3.5 |
| 2.29.0 | 多数 Runner(Dataflow、Flink、Spark)官方支持 Java 11 |
| 2.33.0 | Go SDK 不再标记为 experimental,正式纳入发布流程;最低 Go 1.16 |
| 2.34.0 | Python DataFrame API 从 experimental 转正 |
| 2.37.0 | Dataflow 支持 Java 17;Python 3.9 支持 |
| 2.39.0 | Python 3.6 不再支持 |
| 2.40.0 | Go SDK 最低要求 1.18(以支持泛型) |
| 2.43.0 | Python 3.10 支持 |
| 2.45.0 | Beam 要求pyarrow>=3、pandas>=1.4.3(与numpy==1.24.0兼容) |
| 2.46.0 | Go SDK 要求 Go 1.19 |
| 2.47.0 | Python 3.11 支持 |
| 2.49.0 | Python 3.7 支持移除 |
| 2.50.0 | Go SDK 要求 Go 1.20 |
| 2.52.0 | 发布 Java 21 SDK 容器镜像(Direct 与 Dataflow Runner 实验性支持) |
其中 2.52.0 明确提示:Java 21 在 Direct Runner 与 Dataflow Runner 上为实验性支持,而 Flink、Spark、Samza 等其他 Runner 的支持状态取决于各 Runner 项目自身的进展——这类措辞体现了官方对跨 Runner 能力差异的谨慎界定。
Runner 体系的演进:Splittable DoFn 默认化、Prism 与 Runner v2
Read 变换全面切换到 Splittable DoFn
2.25.0 起,Java 系 Runner(Direct、Flink、Jet、Samza、Twister2)执行 Read 变换时默认使用 Splittable DoFn;2.26.0 扩展至 Spark。期望输出不变,但可通过--experiments=use_deprecated_read临时回退到旧实现。官方在 2.25.0 中明确征求社区反馈,计划永久化该变更,体现了 SDF 作为新一代数据读取模型的地位。
Go SDK 本地 Runner:Prism 取代 Direct Runner
2.46.0 引入 Prism 的原生 Go 实现,2.50.0 正式将其设为 Go SDK 的默认本地 Runner,并宣布 Go Direct Runner 进入弃用状态。仓库中的 Prism README 给出了更完整的定位:
- Prism 是使用 Beam FnAPI 与任意语言 SDK 通信的可移植本地 Runner,可在单机快速启动、便于测试;
- 默认执行时不做 Fusion,逐个独立执行每个 Transform,便于在开发期精准定位错误;
- 支持 Side Inputs、多输出 DoFn、Flatten、GBK/CoGBK(含 Session 窗口)、Combine、Splittable DoFn 展开与 Process Continuations 等特性;
- 当前限制包括:仅限测试使用、未实现 Docker 容器执行(因此暂不能运行 Java/Python SDK 的跨语言变换)、纯内存执行。
对于 Go SDK 用户,迁移路径在 README 中表述得很清楚:短期将 runner 设为"prism",中期切换默认值,长期将"direct"别名指向"prism"并删除旧 Direct Runner。入口实现见 prism.go 的Execute函数。
Dataflow Runner v2 强制化
2.45.0 起,可移植 Java 管道、Go 管道、Python 流处理管道以及可移植 Python 批处理管道在 Dataflow 上必须使用 Runner v2,disable_runner_v2、disable_runner_v2_until_2023、disable_prime_runner_v2等实验选项会在管道构建阶段直接报错。2.50.0 进一步移除了 Dataflow 对 legacy runner 的支持,并停止在--staging_location阶段从 PyPI 打包 Beam SDK——自定义容器镜像若不基于 Beam 默认镜像,必须自行包含 Apache Beam 安装。
Flink 与 Spark 版本支持节奏
- Flink:2.31.0 支持 1.13,2.37.0 支持 1.13.2/1.12.5/1.11.4,2.38.0 支持 1.14.x,2.47.0 支持 1.16.x;2.39.0 起 Flink 1.11 不再支持,更早的 1.8/1.9/1.10 在 2.30.0/2.31.0 被相继移除。
- Spark:2.29.0 起 Spark 3 获得官方支持,2.43.0 起 PortableRunner 默认假设 Spark 3 主版本,2.41.0 预告 Spark 2.4.x 弃用,2.46.0 正式移除 Spark 2 的 SparkRunner,2.50.0 将 Spark 3.2.2 设为默认版本。
- 2.52.0 为 Flink Runner 新增
UseDataStreamForBatch管道选项:设为 true 时批处理作业改用 DataStream API,默认 false(仍走 DataSet API)。
I/O 连接器矩阵的持续扩张
CHANGES.md 的 I/Os 分类展示了一个庞大连接器生态的演进:既有老牌连接器的能力增强,也有新连接器的持续加入。
大数据生态:BigQuery、Kafka、GCS
- BigQuery:2.47.0 起 Python SDK 通过跨语言支持 Storage Write API;2.54.0 支持用 Storage Write API 写入动态目标表;2.34.0 为
ReadFromBigQuery引入method=DIRECT_READ(Storage Read API)与use_native_datetime参数,并默认以 BATCH 优先级运行查询;2.55.0(未发布)支持写入分簇且非时间分区的表(Java)。 - Kafka:2.50.0 的 Java KafkaIO 支持通过
topicPattern动态发现主题;2.53.0 支持坏记录处理(bad records handling);2.20.0 支持使用 Confluent Schema Registry 解析 schema。 - GCS:2.53.0 起 Python GCSIO 用 GCP GCS Client 取代 apitools 实现;2.51.0 修复了 GCS 连接器的异常链问题。
坏记录处理成为统一能力
2.54.0 集中为 FileIO、TextIO、AvroIO、BigtableIO 增加了坏记录处理能力,2.53.0 则将其带给了 KafkaIO。这与 2.34.0 中 Python ParDo(Map、FlatMap 等)新增的with_exception_handling选项一脉相承,标志着"死信模式"(dead letter pattern)成为 Beam 的标准实践。配套的 Error Handler 框架在 2.53.0 的 Java SDK 中落地(PR #29164)。
新连接器一览
Go SDK 在 I/O 上发力明显:2.44.0 增加 Bigtable sink 与 S3 文件系统,2.45.0 增加 MongoDB IO,2.53.0 增加 NATS IO。Java 侧新增了 GoogleAdsIO(2.50.0)、Cosmos DB Core SQL API(2.50.0)、SingleStoreDB(2.44.0)、Neo4j(2.38.0)、PulsarIO(2.39.0)、Firestore(2.32.0)、InfluxDbIO(2.25.0)等。AWS 生态则经历了 amazon-web-services2 模块从 2.38.0 达到功能对等、到 2.48.0 移除 AWS 2 client providers、SnsIO.writeAsync与旧 coders 的完整生命周期。
机器学习与数据科学能力的崛起
2.40.0 是 ML 能力的里程碑:正式引入RunInference——一个框架无关的推理变换,首发支持 PyTorch 与 Scikit-learn。此后版本持续扩充:
- 模型处理器:2.46.0 支持 ONNX runtime 与 Tensorflow Model Handler;2.50.0 加入 Hugging Face Model Handler 与 Hugging Face Pipelines,Vertex AI 支持私有端点;2.51.0 支持用
KeyedModelHandler在同一变换中加载多个模型。 - 侧输入与模型热更新:2.46.0 起 RunInference 接受模型路径作为 SideInputs,并新增
WatchFilePattern变换用于按文件模式监听模型更新;2.52.0 修复了watch_file_pattern参数此前无效的问题,并指引旧行为改用WatchFilePatternPTransform。 - 前处理/后处理与死信:2.48.0 为 RunInference 增加死信队列支持与 pre/post-processing 操作定义;2.51.0 的
VertexAIModelHandlerJSON支持透传inference_args到 Vertex 端点。 - MLTransform:2.50.0 加入 MLTransform,覆盖常见 ML 预处理/后处理操作;2.53.0 支持为 Vertex AI 与 Hugging Face Hub 模型生成文本 embeddings;2.52.0 修复其丢弃重复元素的问题(#29600),并改为输出人类可读的 min/max/quantiles 等工件。
仓库中 RunInference 的实现位于 sdks/python/apache_beam/ml/inference/base.py,模块注释明确说明其职责是处理指标收集、线程间模型共享与元素批处理等标准推理功能,用户通过实现ModelHandler扩展任意 ML 框架。完整的推理示例代码集中在 sdks/python/apache_beam/examples/inference,覆盖 PyTorch 图像分类、Sklearn 回归、TensorFlow、ONNX、Hugging Face、TensorRT、Vertex AI、XGBoost 等场景。
DataFrame API 与交互式分析
2.32.0 宣布 DataFrame API 不再是 experimental,官方给出生产可用建议。此前的 2.25.0 提供了初版预览,2.29.0 支持GroupBy.apply,2.34.0 建议通过pip install apache-beam[dataframe]安装,2.37.0 适配 pandas 1.4.x,2.39.0 实现DataFrame.unstack()、pivot()与Series.unstack()。2.47.0 还允许 Schema 化的 PTransform 像 PCollection 一样直接作用于 DataFrame,并建议多操作时显式链式连接(df | (Transform1 | Transform2 | ...))以避免反复转换。
跨语言与可移植性:FnAPI 与 Beam YAML
- 2.53.0 起,本地运行多语言管道不再需要 Docker——用于展开(expansion)的子进程可同时充当跨语言 worker。
- Go SDK 在 2.37.0 为 JDBCIO、Debezium、BeamSQL、BigQuery、KafkaIO 提供跨语言包装并支持自动启动 expansion service;2.43.0 为 Go 增加 DataFrame 包装(自动 expansion service)。
- 2.35.0 采用新的跨语言变换 URN 约定,可能影响使用自定义 expansion service 连接不同版本 Java/Python SDK 的高级场景。
- 2.52.0 宣布Beam YAML稳定发布:管道可用 YAML 编写,并享受包含初版 I/O 与 turnkey 变换的 YAML 框架。
破坏性变更与迁移指南(升级必读)
CHANGES.md 的价值很大程度上体现在对破坏性变更的明确预警与迁移指引,以下为影响面较大的几条:
Java Avro 模块迁移(2.46.0 起):Avro 相关类在
beam-sdks-java-core中弃用,迁移至beam-sdks-java-extensions-avro,如将org.apache.beam.sdk.coders.AvroCoder换成org.apache.beam.sdk.extensions.avro.coders.AvroCoder(包路径与类层级保持一致以简化迁移)。2.52.0 完成了对旧类的最终移除,CountingSource.CounterMark改用自定义CounterMarkCoder作为默认 coder。Go SDK 模块化(2.33.0):迁移到 Go Modules 后,import 路径需加入
v2,例如.../sdks/go/...变为.../sdks/v2/go/...,go.mod需声明require github.com/apache/beam/sdks/v2。Python CoGroupByKey 类型提示变化(2.43.0):分组值的类型从
List改为Iterable,下游变换若期望List需要调整。BatchElements 批处理行为(2.46.0):默认更激进的批处理(上限从 1 秒改为 10 秒),可用
target_batch_duration_secs_including_fixed_cost=1恢复旧行为。确定性编码强制化(2.29.0):GroupByKey 与有状态 DoFn 强制确定性编码,可注册
FakeDeterministicFastPrimitivesCoder回退或使用allow_non_deterministic_key_coders选项。WriteToBigQuery 行为(2.24.0):需通过
custom_gcs_temp_location或--temp_location提供 GCS 临时位置,或显式传method="STREAMING_INSERTS"。参数更名:Java 侧
--region→--awsRegion(2.27.0)、host→firestoreHost(2.52.0);Python 侧单横杠开头的未解析命令行参数将被忽略(2.47.0)。
安全修复与供应链加固
Security Fixes 分类记录了大量 CVE 修复,主要包括三类:
- 语言工具链升级:使用 go 1.21.1(2.51.0)与 go 1.21.5(2.53.0)构建,分别修复 CVE-2023-39320、CVE-2023-45285 与 CVE-2023-39326;2.52.0 修复 CVE-2023-39325(Java/Python/Go)并缓解 CVE-2023-47248。
- 容器镜像加固:2.55.0(未发布)将 Go SDK 基础容器迁移到
distroless/base-nossl-debian12,把易受攻击面收敛到内核与 glibc;2.51.0 更新 Python 容器修复 7 个 aom 相关 CVE;2.35.0 将测试套件的 Log4j 升级至 2.17.0,并明确 Beam 发行版本身不依赖受 CVE-2021-44228 影响的 log4j-core。 - 依赖 BOM 升级:2.28.0 声明 Guava 30.1-jre,2.35.0 升级 GCP Libraries BOM 至 24.0.0,2.30.0 为 20.0.0——依赖冲突时可通过 Maven
dependencyManagement或 Gradleforce指定版本。
已知问题与升级建议
Known Issues 分类是升级决策的重要输入,值得注意的跨版本问题包括:
- Python 长时运行管道内存泄漏(#28246):自 2.50.0 起被反复提及,直至 2.52.0 宣布修复。
- MLTransform 丢弃重复元素(#29600):2.52.0 修复,此前的 2.50.0/2.51.0 均列为已知问题。
- Python 流式管道偶发崩溃(#27330):影响 2.47.0 及更新版本,2.53.0 修复并强烈建议升级。
- 跨语言 Bigtable sink 时间戳处理(#28632):写入前需为所有记录设置显式时间戳。
- Dataflow 升级兼容性:2.41.0 起默认开启的 Projection Pushdown 优化器可能破坏升级兼容性,可用
--experiments=disable_projection_pushdown禁用。
如何在仓库中验证这些变更
CHANGES.md 的每条记录都能在仓库中找到对应实现:
- Runner 能力:如 Prism 的 README 与 prism.go,Go SDK 验证脚本在 sdks/go/README.md 中通过
./gradlew :sdks:go:test:prismValidatesRunner校验。 - ML 变换:RunInference 核心在 sdks/python/apache_beam/ml/inference/base.py,示例代码见 sdks/python/apache_beam/examples/inference。
- IO 连接器:Java 连接器位于 sdks/java/io,扩展与跨语言组件见 sdks/java/extensions;Python 侧集中在 sdks/python/apache_beam/io。
- Runner 实现:Flink 在 runners/flink,Spark 在 runners/spark,Dataflow 在 runners/google-cloud-dataflow-java。
- 发布流程:版本发布与分支切割的操作细节可参考 contributor-docs/release-guide.md,其中包含 tag RC 提交与发布分支创建的可视化流程说明。
结语
CHANGES.md 远不止是一份流水账式的变更列表:它通过统一的分类模板,完整呈现了 Apache Beam 从"统一批流编程模型"走向"多语言、可移植、ML 原生"的技术演进路径。对开发者而言,它既是升级前必读的兼容性检查清单,也是按图索骥定位源码实现的导航图。理解本文梳理的版本脉络与迁移要点,将帮助你在升级 Beam 时从容应对破坏性变更,并充分利用 RunInference、DataFrame API、Beam YAML 等新能力。
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
相关推荐
httpx认证配置实战:5种认证策略与密钥文件批量搞定授权目标
httpx认证配置实战:5种认证策略与密钥文件批量搞定授权目标 httpx 是一款快速、多用途的 HTTP 探测工具箱(HTTP toolkit),除了批量探测
网络安全CLIAngular CLI 版本变更全解析:从 CHANGELOG.md 读懂 Angular CLI 的演进与升级路径
Angular CLI 版本变更全解析:从 CHANGELOG.md 读懂 Angular CLI 的演进与升级路径 导读 CHANGELOG.md 是 Ang
CLI开发工具前端构建构建工具代码生成前端highlight.js 版本演进全解析:从 CHANGES.md 解读 v9 到 v11 的架构变迁与升级路径
highlight.js 版本演进全解析:从 CHANGES.md 解读 v9 到 v11 的架构变迁与升级路径 本文以开源仓库 highlight.js ht
前端
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考