☰
Apache Beam 0.6.0 中的 Python SDK:Beam 编程模型的第二种实现与实战入门
2026/10/9 1:38:43 网站建设 项目流程
  • 大数据
  • 批处理
  • 流处理
  • 数据工程

【免费下载链接】beam

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

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

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 / FlatMapcore.py面向元素的并行处理原语,__all__中导出ParDo、Map、FlatMap、Filter等
GroupByKey / CombineGloballycore.py按键分组与全局合并,GroupByKey、CombinePerKey、CombineValues均在导出列表中
Windowingwindow.py提供GlobalWindows、FixedWindows、SlidingWindows、Sessions等窗口函数
Create / Impulsecore.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_newlinesTrue每个元素后是否追加换行符
num_shards0输出分片数;为 0 时由执行引擎自动决定,不建议手动约束
shard_name_template'-SSSSS-of-NNNNN'分片命名模板,S与N分别替换为分片序号与总数;''表示单文件输出
coderToBytesCoder()每行编码使用的 Coder
compression_typeCompressionTypes.AUTO压缩类型,AUTO时按文件扩展名自动识别
header/footerNone文件头部/尾部字符串(配合append_trailing_newlines会追加\n)
max_records_per_shard/max_bytes_per_shardNone单个分片的记录数/字节数上限

这些参数让 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):

  1. Direct Runner:在本地机器上直接执行整个管道图。当前源码 direct_runner.py 中,SwitchingDirectRunner会在 FnApiRunner(批处理高吞吐)与 BundleBasedDirectRunner(支持流式执行及部分原语)之间自动切换,因此本地调试体验持续演进;
  2. Dataflow Runner:提交到 Google Cloud Dataflow 托管服务执行,源码位于 dataflow/。

由于当时两个 Runner 都只支持有界 PCollection,Python SDK 的流式能力尚不可用;公告预告"即将到来的特性会让 Python SDK 支持更多 Runner"——这与后续 Beam 推出跨语言 Fn API 的路线完全吻合。

五、Roadmap 回顾:从有界批处理到统一模型

公告最后披露了 Python SDK 当时的两大路线图目标:

  1. 突破有界限制:当时 Runner 仅支持 bounded PCollections,团队计划扩展以支持 unbounded PCollections(即流式处理)。从当前仓库看,这一目标已实现:Direct Runner 与 Dataflow Runner 均支持流式管道,pubsub.py 提供了流式场景的 Pub/Sub IO;
  2. 扩展 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.

项目地址:https://gitcode.com/gh_mirrors/beam4/beam
点击查看免费下载
上一篇:Windows 11 免重装去臃肿,10 分钟搞定:Win11Debloat 新手入门指南
下一篇:Playball请求限流机制:保护API服务的措施

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

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询