资讯动态

Apache Beam 0.6.0 中的 Python SDK:Beam 编程模型的第二种实现与实战入门

发布时间:2026/10/9 1:38:46 来源:尧图企业网站定制
大数据批处理流处理数据工程【免费下载链接】beamApache 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.02017 年 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, (N1) * 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/WriteToTextAvro 读写avroio.pyTensorFlow Recordtfrecordio.pyGoogle BigQuerybigquery.pyGoogle Cloud Datastoredatastore/以 textio.py 为例_TextSource继承FileBasedSource按\n/\r\n将文件切分为元素WriteToText的构造参数非常丰富完整参数如下来自 textio.py参数默认值作用file_path_prefix必填输出文件路径前缀后接分片标识与file_name_suffixfile_name_suffix输出文件扩展名append_trailing_newlinesTrue每个元素后是否追加换行符num_shards0输出分片数为 0 时由执行引擎自动决定不建议手动约束shard_name_template-SSSSS-of-NNNNN分片命名模板S与N分别替换为分片序号与总数表示单文件输出coderToBytesCoder()每行编码使用的 Codercompression_typeCompressionTypes.AUTO压缩类型AUTO时按文件扩展名自动识别header/footerNone文件头部/尾部字符串配合append_trailing_newlines会追加\nmax_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从内存列表创建有界 PCollection10 个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 的硬性要求源码中有明确注释自定义 PTransformEstimatePiTransform(beam.PTransform)把创建 100 个各含 10 万次试验的任务 → Map → CombineGlobally封装为可复用变换默认tries_per_work_item100000即共 1000 万次投掷自定义 CoderJsonCoder将结果序列化为 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 executionDirect Runner在本地机器上直接执行整个管道图。当前源码 direct_runner.py 中SwitchingDirectRunner会在 FnApiRunner批处理高吞吐与 BundleBasedDirectRunner支持流式执行及部分原语之间自动切换因此本地调试体验持续演进Dataflow Runner提交到 Google Cloud Dataflow 托管服务执行源码位于 dataflow/。由于当时两个 Runner 都只支持有界 PCollectionPython 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.pyIO 实现与 direct_runner.py执行模型即可完整理解从管道定义到本地执行的全链路。赞分享大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载相关推荐Apache Beam 入门第一课Hello Beam Kata 实战与源码级解析Apache Beam 入门第一课Hello Beam Kata 实战与源码级解析 本篇技术指南聚焦于 Apache Beam 官方互动式训练营Kata中Video2X 视频超分辨率教程老视频放大到 4K 的完整指南Video2X 视频超分辨率教程老视频放大到 4K 的完整指南 老视频放大后为什么全是马赛克360P 素材又该怎么变成 4KVideo2X 就是一个免费开音视频视频处理图像处理深度学习Apache Beam Kotlin 入门第一课用 Create 构造 Hello Beam 管道Katas 实战Apache Beam Kotlin 入门第一课用 Create 构造 Hello Beam 管道Katas 实战 Apache Beam 是开源的统大数据批处理流处理数据工程上一篇Windows 11 免重装去臃肿10 分钟搞定Win11Debloat 新手入门指南下一篇Playball请求限流机制保护API服务的措施创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价 →
↑