- 大数据
- 批处理
- 流处理
- 数据工程
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
Apache Beam 0.6.0 版本首次将 Beam 统一编程模型带到 Python 语言,使其与 Java SDK 一起构成 Beam 模型的两大官方实现。本文以官方发布公告为主线,梳理 Python SDK 的能力边界、Pi 估算实战示例,并结合当前仓库源码(sdks/python/下 1600+ Python 文件)讲解 ParDo、GroupByKey、Windowing 等核心原语的底层实现与可扩展 IO 架构,最后回顾其当时的技术局限与演进路线。
一、发布背景:Beam 0.6.0 与 Python SDK 的诞生
Apache Beam 0.6.0(2017 年 3 月发布)是一个具有里程碑意义的版本——它第一次把 Beam 的编程模型带到了 Python 生态。在这之前,Beam 只有 Java SDK,而 0.6.0 引入了 Python SDK,作为该模型第二个官方实现,让数据工程师可以用 Python 编写批处理管道,同时复用 Beam 统一的管道抽象。
该公告的完整内容保存在仓库的 python-sdk-release.md,是理解 Python SDK 早期能力与设计哲学的第一手资料。本文所有代码与源码引用均来自当前 Beam 仓库(gh_mirrors/beam4/beam)。
二、Python SDK 的核心能力:完整继承 Beam 编程模型
公告明确说明,Python SDK 完整纳入了 Beam 模型的全部核心概念,包括 ParDo、GroupByKey、Windowing 等。这些原语在当前仓库中都有成熟实现:
| 核心概念 | 源码位置 | 说明 |
|---|---|---|
| ParDo / Map / FlatMap | core.py | 面向元素的并行处理原语,__all__中导出ParDo、Map、FlatMap、Filter等 |
| GroupByKey / CombineGlobally | core.py | 按键分组与全局合并,GroupByKey、CombinePerKey、CombineValues均在导出列表中 |
| Windowing | window.py | 提供GlobalWindows、FixedWindows、SlidingWindows、Sessions等窗口函数 |
| Create / Impulse | core.py | 从内存数据或空脉冲创建 PCollection 的源头变换 |
从源码结构看,window.py 中每种窗口函数都定义了清晰的区间公式:FixedWindows将每个元素映射到[N * size + offset, (N+1) * size + offset)时间区间;SlidingWindows使用[N * period + offset, N * period + offset + size);Sessions则按指定的gap_size把间隔小于该值的连续事件聚合成会话。这些正是 Python SDK 宣称"包含 Windowing 等全部主要概念"的底层支撑。
2.1 可扩展的 IO API:有界 Source 与 Sink
公告指出 Python SDK 提供了可扩展的 IO API,用于编写有界(bounded)的 Source 与 Sink。当前仓库的 io/ 目录就是这套体系的最佳注脚:
- 文本读写:textio.py 中的
ReadFromText/WriteToText - Avro 读写:avroio.py
- TensorFlow Record:tfrecordio.py
- Google BigQuery:bigquery.py
- Google Cloud Datastore:datastore/
以 textio.py 为例,_TextSource继承FileBasedSource,按'\n'/'\r\n'将文件切分为元素;WriteToText的构造参数非常丰富,完整参数如下(来自 textio.py):
| 参数 | 默认值 | 作用 |
|---|---|---|
file_path_prefix | 必填 | 输出文件路径前缀,后接分片标识与file_name_suffix |
file_name_suffix | '' | 输出文件扩展名 |
append_trailing_newlines | True | 每个元素后是否追加换行符 |
num_shards | 0 | 输出分片数;为 0 时由执行引擎自动决定,不建议手动约束 |
shard_name_template | '-SSSSS-of-NNNNN' | 分片命名模板,S与N分别替换为分片序号与总数;''表示单文件输出 |
coder | ToBytesCoder() | 每行编码使用的 Coder |
compression_type | CompressionTypes.AUTO | 压缩类型,AUTO时按文件扩展名自动识别 |
header/footer | None | 文件头部/尾部字符串(配合append_trailing_newlines会追加\n) |
max_records_per_shard/max_bytes_per_shard | None | 单个分片的记录数/字节数上限 |
这些参数让 Python SDK 从发布之初就具备生产级的文件 IO 能力,而coder参数则体现了 Beam 对数据编码(Coder)的一等公民支持。
三、实战入门:安装、运行与 Pi 估算示例
3.1 安装与启动
公告给出的安装方式非常简单,从 PyPI 安装apache-beam包即可:
$ pip install apache-beam $ python这条命令会安装当前 Python SDK 及其依赖。当前仓库的 setup.py 与 pyproject.toml 记录了完整打包信息,读者也可以直接从源码构建。
3.2 完整示例:用蒙特卡洛方法估算 Pi
公告以一个纪念 Pi Day 的趣味示例展示 SDK 用法——向单位正方形随机投掷飞镖,统计落入单位圆内的比例,从而估算 π。下面是公告中的原始代码:
import random import apache_beam as beam def run_trials(count): """Throw darts into unit square and count how many fall into unit circle.""" inside = 0 for _ in xrange(count): x, y = random.uniform(0, 1), random.uniform(0, 1) inside += 1 if x*x + y*y <= 1.0 else 0 return count, inside def combine_results(results): """Given all the trial results, estimate pi.""" total, inside = sum(r[0] for r in results), sum(r[1] for r in results) return total, inside, 4 * float(inside) / total if total > 0 else 0 p = beam.Pipeline() (p | beam.Create([500] * 10) # Create 10 experiments with 500 samples each. | beam.Map(run_trials) # Run experiments in parallel. | beam.CombineGlobally(combine_results) # Combine the results. | beam.io.WriteToText('./pi_estimate.txt')) # Write PI estimate to a file. p.run()运行后查看估算结果:
$ cat pi_estimate.txt*这个例子虽短,却串联起了 Beam 模型的三个关键原语:
beam.Create:从内存列表创建有界 PCollection(10 个"500 次试验"的任务);beam.Map:对每个元素并行执行run_trials,即 ParDo 的简单形态;beam.CombineGlobally:把各任务结果汇总,用combine_results合并出 π 的估计值;WriteToText:把结果写出为文本文件。
3.3 仓库中的完整版本:EstimatePiTransform
当前仓库保留了该示例的完整工程化版本,位于 estimate_pi.py,它比公告中的演示代码更加严谨:
- 类型标注:使用
@beam.typehints.with_input_types/with_output_types为run_trials和combine_results声明输入输出类型,从而在管道构建期做类型检查; - combiner 输入输出类型一致性:
run_trials返回(runs, inside_runs, 0)三元组,最后一个 0 用于保证 combiner 函数输入输出类型相同(Beam 对 combiner 的硬性要求,源码中有明确注释); - 自定义 PTransform:
EstimatePiTransform(beam.PTransform)把"创建 100 个各含 10 万次试验的任务 → Map → CombineGlobally"封装为可复用变换,默认tries_per_work_item=100000,即共 1000 万次投掷; - 自定义 Coder:
JsonCoder将结果序列化为 JSON 字节串,作为WriteToText的coder参数; save_main_session:通过SetupOptions.save_main_session = True保存主模块上下文,确保分布式执行时DoFn能引用模块级全局状态(源码注释明确说明该设置的必要性)。
运行完整版示例的方式:
$ python sdks/python/apache_beam/examples/complete/estimate_pi.py --output ./pi_estimate.json四、执行器现状:Direct Runner 与 Dataflow Runner
公告指出,Python SDK 发布之初有两个可用的执行器(Runner),且均仅支持批处理(batch execution):
- Direct Runner:在本地机器上直接执行整个管道图。当前源码 direct_runner.py 中,
SwitchingDirectRunner会在 FnApiRunner(批处理高吞吐)与 BundleBasedDirectRunner(支持流式执行及部分原语)之间自动切换,因此本地调试体验持续演进; - Dataflow Runner:提交到 Google Cloud Dataflow 托管服务执行,源码位于 dataflow/。
由于当时两个 Runner 都只支持有界 PCollection,Python SDK 的流式能力尚不可用;公告预告"即将到来的特性会让 Python SDK 支持更多 Runner"——这与后续 Beam 推出跨语言 Fn API 的路线完全吻合。
五、Roadmap 回顾:从有界批处理到统一模型
公告最后披露了 Python SDK 当时的两大路线图目标:
- 突破有界限制:当时 Runner 仅支持 bounded PCollections,团队计划扩展以支持 unbounded PCollections(即流式处理)。从当前仓库看,这一目标已实现:Direct Runner 与 Dataflow Runner 均支持流式管道,pubsub.py 提供了流式场景的 Pub/Sub IO;
- 扩展 Runner 支持:计划通过即将推出的Fn API将 Python SDK 带到更多执行引擎。从仓库结构看,这一路线也已落地——runners/portability/ 目录承载跨 Runner 的可移植执行支持,runners/flink/ 等目录表明 Python 管道如今已能运行在 Flink 等更多引擎上。
从 0.6.0 至今,Python SDK 正是沿着这两条主线,逐步兑现了 Beam 的使命宣言——"一个统一的批处理与流式数据处理编程模型,可运行于任意执行引擎之上"。
六、总结
Apache Beam 0.6.0 的 Python SDK 是 Beam 生态的重要转折点:它以完整的 ParDo/GroupByKey/Windowing 原语、可扩展的 IO 体系(Text/Avro/TFRecord/BigQuery/Datastore)和两个可用的 Runner,向 Python 开发者开放了统一的批处理编程模型。公告中那个"投掷飞镖估算 Pi"的小例子,至今仍可在仓库 estimate_pi.py 中找到其工程化版本——这正是理解 Beam 管道"构建 → 变换 → 合并 → 写出"工作流的最佳起点。对想要深入学习的读者,建议依次阅读 core.py(变换原语)、textio.py(IO 实现)与 direct_runner.py(执行模型),即可完整理解从管道定义到本地执行的全链路。
- 大数据
- 批处理
- 流处理
- 数据工程
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
相关推荐
Apache Beam 入门第一课:Hello Beam Kata 实战与源码级解析
Apache Beam 入门第一课:Hello Beam Kata 实战与源码级解析 本篇技术指南聚焦于 Apache Beam 官方互动式训练营(Kata)中
Video2X 视频超分辨率教程:老视频放大到 4K 的完整指南
Video2X 视频超分辨率教程:老视频放大到 4K 的完整指南 老视频放大后为什么全是马赛克?360P 素材又该怎么变成 4K?Video2X 就是一个免费开
音视频视频处理图像处理深度学习Apache Beam Kotlin 入门第一课:用 Create 构造 "Hello Beam" 管道(Katas 实战)
Apache Beam Kotlin 入门第一课:用 Create 构造 "Hello Beam" 管道(Katas 实战) Apache Beam 是开源的统
大数据批处理流处理数据工程
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考