- 大数据
- 批处理
- 流处理
- 数据工程
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
本篇技术指南围绕 Apache Beam 的核心概念——Runner(执行引擎)展开,结合当前仓库(gh_mirrors/beam4/beam)的 Python SDK 源码与官方示例,系统讲解 Runner 是什么、为什么选择 Runner 是流水线开发的关键一步、如何通过--runner参数指定执行引擎,以及如何使用 DirectRunner 在本机完成测试与调试。读完本文,你将掌握在本地 Direct Runner 与云端 Dataflow Runner 之间切换运行 Beam 流水线的完整方法。
Runner 是什么:统一编程模型与执行引擎的解耦
在 Apache Beam 中,Runner 是负责执行流水线的执行引擎。Beam 的核心设计哲学是"一套代码,多种执行引擎":用户用统一的 Beam 编程模型(PCollection、PTransform、DoFn 等)编写流水线,而 Runner 负责把流水线翻译或适配成可以在某个大规模并行数据处理系统上运行的形式,例如 Apache Flink、Apache Spark、Google Cloud Dataflow 等。
从当前仓库的源码可以看到这一抽象的直接体现。Python SDK 中定义了抽象基类PipelineRunner,位于 sdks/python/apache_beam/runners/runner.py:
run():同步执行给定 PTransform,阻塞直到流水线完成(内部调用run_async后wait_until_finish);run_async():异步启动执行,返回可查询进度的PipelineResult;run_portable_pipeline():接收标准化的beam_runner_api_pb2.Pipeline协议消息执行整个流水线,是各 Runner 需要覆写的关键方法;run_pipeline():将 Python 端的Pipeline对象转换为 Runner API 协议后交给run_portable_pipeline。
这意味着不同 Runner 共享同一套流水线描述(Beam Runner API),只是执行后端不同。也正因如此,"选择 Runner"才成为流水线开发流程中一个独立的、重要的环节。
为什么选择 Runner 是关键一步
你选择的 Runner 决定了流水线在哪里运行、以何种方式运行。同一个 WordCount 程序,在 DirectRunner 上运行就是本地单机进程;在 DataflowRunner 上运行就是提交到 Google Cloud 的托管服务;在 FlinkRunner / SparkRunner 上运行则是提交到对应的集群。
不同 Runner 对 Beam Model(统一批流模型)的各个特性支持程度不同——例如窗口、触发、状态、定时器、Splittable DoFn 等能力的完备程度。为此,Apache Beam 官方维护了一份Beam Capability Matrix(能力矩阵),用于横向对比各 Runner 对各类特性的支持情况。该页面在仓库中的源文件为 website/www/site/content/en/documentation/runners/capability-matrix/_index.md,其中明确说明:Beam 的可移植 API 层允许流水线"跨多种执行引擎执行",而每个 Runner 对 Beam Model 的实现程度各不相同,能力矩阵正是为澄清这一点而设计。
因此,在选择 Runner 时应重点评估:
- 运行环境:本地开发机、自建 Flink/Spark 集群,还是云上托管服务;
- 特性支持:流水线用到的窗口、触发、状态等特性在该 Runner 上是否完备(对照能力矩阵);
- 生产需求:吞吐、容错、资源弹性、运维成本等。
如何指定 Runner:--runner 参数
在执行流水线时,通过--runner标志指定 Runner。Beam 的 Python SDK 中,StandardOptions负责解析该参数,源码位于 sdks/python/apache_beam/options/pipeline_options.py。
该参数接受的有效值由ALL_KNOWN_RUNNERS定义(当前仓库版本)包括:
| Runner 名称 | 对应实现类(模块路径) |
|---|---|
DataflowRunner | apache_beam.runners.dataflow.dataflow_runner.DataflowRunner |
DirectRunner | apache_beam.runners.direct.direct_runner.DirectRunner |
BundleBasedDirectRunner | apache_beam.runners.direct.direct_runner.BundleBasedDirectRunner |
SwitchingDirectRunner | apache_beam.runners.direct.direct_runner.SwitchingDirectRunner |
InteractiveRunner | apache_beam.runners.interactive.interactive_runner.InteractiveRunner |
FlinkRunner | apache_beam.runners.portability.flink_runner.FlinkRunner |
FnApiRunner | apache_beam.runners.portability.fn_api_runner.FnApiRunner |
KafkaStreamsRunner | apache_beam.runners.portability.kafka_streams_runner.KafkaStreamsRunner |
PortableRunner | apache_beam.runners.portability.portable_runner.PortableRunner |
PrismRunner | apache_beam.runners.portability.prism_runner.PrismRunner |
SparkRunner | apache_beam.runners.portability.spark_runner.SparkRunner |
除上述名称外,--runner也接受任何PipelineRunner子类的完整限定名。若未指定,默认值为DirectRunner(见StandardOptions.DEFAULT_RUNNER)。
从源码可以看出,Runner 名称到实例的解析逻辑位于create_runner()函数(sdks/python/apache_beam/runners/runner.py):它先把小写名称映射到_RUNNER_MAP中的模块路径,再通过importlib动态导入并实例化;若导入失败,还会给出针对性提示(例如 Dataflow Runner 需要pip install apache_beam[gcp],Interactive Runner 需要apache_beam[interactive])。
实战:在 Google Cloud Dataflow 上运行 WordCount
官方文档给出的在 Dataflow 上运行 WordCount 的完整命令如下:
python -m apache_beam.examples.wordcount \ --region DATAFLOW_REGION \ --input gs://dataflow-samples/shakespeare/kinglear.txt \ --output gs://STORAGE_BUCKET/results/outputs \ --runner DataflowRunner \ --project PROJECT_ID \ --temp_location gs://STORAGE_BUCKET/tmp/各参数含义:
| 参数 | 说明 |
|---|---|
--region DATAFLOW_REGION | Dataflow 作业运行的 Google Cloud 区域,例如us-central1 |
--input gs://... | 输入文件的 Cloud Storage 路径(此处为公开示例数据kinglear.txt) |
--output gs://... | 输出结果写入的 Cloud Storage 前缀 |
--runner DataflowRunner | 指定执行引擎为 Google Cloud Dataflow |
--project PROJECT_ID | 提交作业的 Google Cloud 项目 ID |
--temp_location gs://STORAGE_BUCKET/tmp/ | 临时文件(staging 文件、中间结果)存放位置,Dataflow 必需 |
注意:--project与--temp_location在 DataflowRunner 下属于必填项;--runner以外的参数大多可由PipelineOptions统一解析,无需在代码中单独声明。该示例程序的源码位于仓库 sdks/python/apache_beam/examples/wordcount.py,它演示了从文本读取、正则分词、按单词计数到写出结果的完整流水线,是理解 Runner 行为的最小可运行样例。
DirectRunner:本地执行引擎
Apache Beam Direct Runner 在本地机器上执行流水线,非常适合测试与调试。它不需要任何分布式集群或云环境,安装 Beam SDK 后即可直接运行。
从源码看,Direct Runner 并非单一实现(sdks/python/apache_beam/runners/direct/direct_runner.py):
BundleBasedDirectRunner:基于 Bundle 机制的本地执行器,支持流式执行以及某些 FnApiRunner 尚未实现的原语;SwitchingDirectRunner:根据流水线特性自动在 FnApiRunner(批处理吞吐高)与 BundleBasedDirectRunner(支持流式)之间切换,是本地执行的"智能"入口;DirectRunner:默认的别名入口,作为DEFAULT_RUNNER直接可用。
python -m apache_beam.examples.wordcount \ --input /path/to/input.txt \ --output /tmp/wordcount-output不指定--runner时,Beam 默认使用 DirectRunner,上面的命令即会在本地完成整个 WordCount 计算,输出结果写入/tmp/wordcount-output前缀对应的分片文件。本地调试时可以省略--project、--region、--temp_location等云端参数,这让"先本地跑通、再上云跑批"的开发节奏非常顺畅。
DirectRunner 的价值在于:流水线中的每个变换(ParDo、GroupByKey、Combine 等)都在本地真实执行,配合 Beam 的日志、Metrics 与PipelineResult查询接口,可以快速定位逻辑错误、验证窗口与触发语义。
Runner 与流水线生命周期:PipelineResult 与 PipelineState
运行流水线后得到的PipelineResult对象(sdks/python/apache_beam/runners/runner.py)提供三个核心能力:
state:查询流水线当前状态;wait_until_finish(duration):阻塞等待流水线结束并返回最终状态(可指定毫秒级超时);cancel()/metrics():取消执行或获取运行指标。
PipelineState枚举(sdks/python/apache_beam/runners/runner.py)是各 Runner 状态的并集,包含STARTING、RUNNING、DONE、FAILED、CANCELLED、DRAINING、DRAINED等状态,其中DONE、FAILED、CANCELLED、UPDATED、DRAINED是终态。这在调试 DirectRunner 时特别有用:wait_until_finish()的返回值能明确告诉你流水线是成功(DONE)还是失败(FAILED),配合异常堆栈即可完成问题定位。
如何选择 Runner 并继续深入学习
面对多个 Runner,建议的决策路径是:
- 开发阶段:一律使用默认的
DirectRunner,快速验证逻辑正确性; - 评估特性:对照能力矩阵(仓库源文件:website/www/site/content/en/documentation/runners/capability-matrix/_index.md)确认所需特性在目标 Runner 上的支持程度;
- 生产运行:根据基础设施情况选择 Dataflow(云托管)、Flink、Spark、Prism 等 Runner,并通过
--runner参数切换,代码本体无需改动; - 调优与排查:利用
PipelineResult的状态查询与 Metrics 接口,结合各 Runner 的控制台监控定位执行问题。
进一步上手,可参考仓库内各 SDK 的 Quickstart 与示例:Java 示例见 sdks/java/core 相关模块,Python 示例见 sdks/python/apache_beam/examples,Go 示例见 sdks/go/examples。按照各自的 Quickstart 配置开发环境与 Runner 后,即可动手实践本文中的命令。
- 大数据
- 批处理
- 流处理
- 数据工程
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
相关推荐
Apache Beam 执行引擎(Runner)核心概念详解
Apache Beam 执行引擎 Runner 核心概念详解 一、Apache Beam 执行引擎概述 Apache Beam 作为一个统一的大数据处理编程模型
大数据批处理流处理数据工程Apache Beam Runner选型指南:DirectRunner、Flink、Spark与Dataflow如何选?
Apache Beam Runner选型指南:DirectRunner、Flink、Spark与Dataflow如何选? Apache Beam 是一个统一的批
大数据批处理流处理数据工程Apache Beam 入门指南:统一批流处理模型、多语言 SDK 与 Runner 执行体系
Apache Beam 入门指南:统一批流处理模型、多语言 SDK 与 Runner 执行体系 Apache Beam 是一个用于定义批处理和流处理数据并行管道
大数据批处理流处理数据工程
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考