Apache Airflow Apache Beam Provider 完全指南:安装、依赖与 Python/Java/Go 流水线编排实战
【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow
Apache Airflow 的apache-airflow-providers-apache-beam(6.2.4)是连接 Airflow 与 Apache Beam 的官方 Provider 包,它允许你在 DAG 中以任务的形式启动 Beam 批处理与流式数据流水线,并支持 DirectRunner、DataflowRunner、SparkRunner、FlinkRunner 等多种 Runner。阅读本文后,你将掌握该 Provider 的安装与版本约束、跨 Provider 依赖配置、官方包校验方法,以及通过BeamRunPythonPipelineOperator、BeamRunJavaPipelineOperator、BeamRunGoPipelineOperator三种 Operator 在实际 DAG 中编排 Beam 流水线的完整方案,包括可延迟(deferrable)异步执行与 Google Cloud Dataflow 深度集成。
本文以 Provider 首页文档 为骨架,结合 Operators 指南、源码实现 与系统测试示例(tests/system/apache/beam)展开,全部结论均可回到当前仓库验证。
一、Apache Beam Provider 是什么
Apache Beam 是一个用于定义批处理与流式数据并行处理管道的统一开源模型:开发者使用 Beam SDK 编写定义管道的程序,再由 Beam 支持的分布式处理后端执行——包括 Apache Flink、Apache Spark 与 Google Cloud Dataflow。Airflow 中的apache.beamProvider 把"提交并跟踪一个 Beam 管道"封装成 Airflow 任务,让数据工程师可以在一个编排系统内统一调度 Beam 作业与其他上下游任务。
从 get_provider_info.py 可以看到该 Provider 对外暴露的组件全貌:
- Operators:模块
airflow.providers.apache.beam.operators.beam,提供三个流水线执行算子; - Hooks:模块
airflow.providers.apache.beam.hooks.beam,负责底层命令构造与子进程执行; - Triggers:模块
airflow.providers.apache.beam.triggers.beam,支撑可延迟异步执行。
所有类都打包在airflow.providers.apache.beamPython 包内,与包名apache-airflow-providers-apache-beam一一对应。
二、安装与版本要求
2.1 基础安装
在已安装 Airflow 的环境中,通过 pip 直接安装:
pip install apache-airflow-providers-apache-beam2.2 最低版本要求
当前仓库中该 Provider 的发布版本为6.2.4(见 src/airflow/providers/apache/beam/init.py),其对 Apache Airflow 的最低支持版本是2.11.0。这一点在运行时也会被强制校验:__init__.py在导入时会解析 Airflow 版本号,若低于2.11.0则直接抛出RuntimeError,提示需要 Airflow 2.11.0+。
Provider 的完整依赖矩阵如下(来自 index.rst 与 README.rst):
| PIP package | Version required |
|---|---|
apache-airflow | >=2.11.0 |
apache-airflow-providers-common-compat | >=1.12.0 |
apache-beam | >=2.76.0 |
pyarrow | >=16.1.0; python_version < "3.14" |
pyarrow | >=22.0.0; python_version >= "3.14" |
numpy | >=1.22.4; python_version < "3.11" |
numpy | >=1.23.2; python_version < "3.12" and python_version >= "3.11" |
numpy | >=1.26.0; python_version >= "3.12" and python_version < "3.14" |
numpy | >=2.4.3; python_version >= "3.14" |
可见依赖会随 Python 版本动态选择pyarrow与numpy的版本:例如 Python < 3.14 使用 pyarrow>=16.1.0,Python >= 3.14 使用 pyarrow>=22.0.0;numpy 则在 1.22.4 到 2.4.3 之间按解释器版本分档。README 同时说明该包支持 Python 3.10、3.11、3.12、3.13、3.14。
三、跨 Provider 依赖与可选依赖
3.1 跨 Provider 依赖(google extra)
要使用全部功能(尤其是 DataflowRunner 与 GCS 文件拉取),需要安装 Google Provider。安装时指定 extra:
pip install apache-airflow-providers-apache-beam[google]| Dependent package | Extra |
|---|---|
apache-airflow-providers-google | google |
3.2 可选依赖
googleextra 同时会引入 Beam 的 GCP 扩展:
| Extra | Dependencies |
|---|---|
google | apache-beam[gcp]>=2.76.0 |
3.3 缺少 Google Provider 时的行为
从源码看,Google Provider 的缺失并不会导致导入失败,而是在运行时按需报错:
- operators/beam.py 中
BeamDataflowMixin.__init__在未检测到 Google Provider 时抛出AirflowOptionalProviderFeatureException,提示需安装合适的 Google Provider 版本才能使用 Dataflow 服务; - Go 流水线在启动时会校验能否导入
airflow.providers.google.go_module_utils,否则同样抛出AirflowOptionalProviderFeatureException(见 hooks/beam.py 中start_go_pipeline的实现)。
因此,如果只跑本地 DirectRunner 且使用本地文件,可以不装 Google Provider;一旦涉及 GCS 文件、Dataflow 或 Go 流水线,就必须安装[google]extra。
四、官方发布包下载与校验
Provider 的正式发布包可从 Apache 官方下载站获取,并支持校验和(checksum)与签名(signature)验证:
- sdist 包:
apache_airflow_providers_apache_beam-6.2.4.tar.gz(附带.asc与.sha512文件) - wheel 包:
apache_airflow_providers_apache_beam-6.2.4-py3-none-any.whl(附带.asc与.sha512文件)
建议下载后先核对 SHA-512 校验和再用 GPG 验证.asc签名,确保包来源可信、未被篡改。项目内各组件均遵循 Apache License 2.0(见 providers/apache/beam/LICENSE)。
五、通过 Operator 运行 Beam 流水线
Apache Beam Operator 指南(operators.rst)提供了 Python、Java、Go 三种 SDK 的完整用法。三个 Operator 均继承自BeamBasePipelineOperator,共享runner、default_pipeline_options、pipeline_options、gcp_conn_id、dataflow_config等核心参数。
5.1 公共参数解析
从 BeamBasePipelineOperator 的 docstring 与实现可以归纳出公共参数语义:
runner:流水线执行后端,默认"DirectRunner"。可选值还包括DataflowRunner、SparkRunner、FlinkRunner、PortableRunner等,完整枚举见 hooks/beam.py 中的BeamRunnerType(还包含 SamzaRunner、NemoRunner、JetRunner、Twister2Runner);default_pipeline_options:默认管道选项,适合存放对 DAG 中所有 Beam Operator 通用的高层参数(如 project、zone);pipeline_options:管道选项字典,会与default_pipeline_options合并后传给 Beam。值类型决定生成的命令行参数:None值:该选项被显式跳过(避免被当作字符串'None'传给 Beam);True:生成无值的--key选项;False:选项被跳过;但对use_public_ips这类特殊标志,会生成对应的取反标志--no_use_public_ips(见 beam_options_to_args 中的_FLAG_THAT_SETS_FALSE_VALUE映射);list:为每个元素生成一个--key=value(例如['A','B']生成--key=A --key=B);dict:序列化为 JSON 字符串传入(常用于 labels);- 其他类型:使用 Python 文本表示。
gcp_conn_id:连接 GCS 使用的 Airflow 连接 ID,默认"google_cloud_default";dataflow_config:Dataflow 专用配置(DataflowConfiguration或等价 dict),仅在 runner 为DataflowRunner时生效;若 runner 不是 DataflowRunner 却配置了该项,Operator 会打印告警日志。
5.2 Dataflow 集成要点
当runner="DataflowRunner"时,BeamDataflowMixin 会做一系列自动化处理:
- 通过
DataflowHook构造 Dataflow 作业名(支持append_job_name自动追加后缀); - 将
serviceAccount、impersonateServiceAccount(支持列表拼接为逗号分隔字符串)、project、region等写入 pipeline options; - 自动为每个作业注入
airflow-versionlabel,便于在 Dataflow 控制台识别作业来源; - 从 Beam 输出日志中实时解析 Dataflow job id,并立即通过 XCom 以
dataflow_job_id为 key 推送(见 dataflow_job_id setter),让 Sensor 可以在作业结束前就拿到 job id; - 作业结束时通过
DataflowJobLink在 Airflow UI 中提供直达 Dataflow 监控页的链接; - 任务被 kill 时调用
cancel_job取消对应 Dataflow 作业(on_kill钩子)。
注意(官方指南明确提示):当 Beam 流水线运行在 Dataflow 服务上时,要求 Airflow Worker 节点安装
gcloud(Google Cloud SDK)命令行工具。
六、运行 Python 流水线:BeamRunPythonPipelineOperator
6.1 核心参数
py_file(必填,支持模板):要执行的 Beam Python 管道文件,可以是本地绝对路径,也可以是 GCS 上的gs://路径(Airflow 会通过GCSHook.provide_file下载为临时文件,见 execute);py_interpreter:执行管道所用的 Python 版本,默认python3。若 Airflow 实例运行在 Python 2 环境则指定python2并确保py_file为 Python 2 代码,官方建议优先使用 Python 3;py_options:额外 Python 选项,如["-m", "-v"];py_requirements:指定后将创建临时 Python 虚拟环境并在其中安装这些依赖,然后在此环境内运行管道——也可用来安装特定版本的apache-beam;py_system_site_packages:创建虚拟环境时是否包含 Airflow 实例的系统站点包。默认False,除非 Dataflow 作业确实需要,否则不建议开启;deferrable:是否以可延迟异步模式运行,默认读取operators.default_deferrable配置(默认 False)。
需要注意的参数组合约束(来自 hooks/beam.py 的start_python_pipeline):如果指定了py_requirements但列表为空、且py_system_site_packages=False,会抛出异常——因为此时虚拟环境中没有apache-beam包,作业无法执行。修复方法是:要么在系统安装apache-beam并设置py_system_site_packages=True,要么把apache-beam加入py_requirements列表。
6.2 完整 DAG 示例(DirectRunner)
以下示例取自系统测试 example_python.py:
from airflow import models from airflow.providers.apache.beam.operators.beam import BeamRunPythonPipelineOperator with models.DAG( "example_beam_native_python", start_date=START_DATE, schedule=None, # 按需覆盖 catchup=False, default_args=DEFAULT_ARGS, tags=["example"], ) as dag: # 本地文件 + DirectRunner(通过模块方式运行 wordcount 示例) start_python_pipeline_local_direct_runner = BeamRunPythonPipelineOperator( task_id="start_python_pipeline_local_direct_runner", py_file="apache_beam.examples.wordcount", py_options=["-m"], py_requirements=["apache-beam[gcp]==2.59.0"], py_interpreter="python3", py_system_site_packages=False, ) # GCS 文件 + DirectRunner start_python_pipeline_direct_runner = BeamRunPythonPipelineOperator( task_id="start_python_pipeline_direct_runner", py_file=GCS_PYTHON, # 形如 gs://bucket/path/to/pipeline.py py_options=[], pipeline_options={"output": GCS_OUTPUT}, py_requirements=["apache-beam[gcp]==2.59.0"], py_interpreter="python3", py_system_site_packages=False, )6.3 DataflowRunner 示例
from airflow.providers.google.cloud.operators.dataflow import DataflowConfiguration start_python_pipeline_dataflow_runner = BeamRunPythonPipelineOperator( task_id="start_python_pipeline_dataflow_runner", runner="DataflowRunner", py_file=GCS_PYTHON, pipeline_options={ "tempLocation": GCS_TMP, "stagingLocation": GCS_STAGING, "output": GCS_OUTPUT, }, py_options=[], py_requirements=["apache-beam[gcp]==2.59.0"], py_interpreter="python3", py_system_site_packages=False, dataflow_config=DataflowConfiguration( job_name="{{task.task_id}}", project_id=GCP_PROJECT_ID, location="us-central1" ), )job_name支持 Jinja 模板(示例中使用{{task.task_id}}让作业名与任务 ID 一致);dataflow_config也接受等价的 dict 写法。示例 DAG 中还展示了 SparkRunner、FlinkRunner 的用法,只要把runner替换并传入对应 pipeline options(如 Flink 的output、Spark 的endpoint)即可。
七、运行 Java 流水线:BeamRunJavaPipelineOperator
7.1 核心参数
jar(必填,支持模板):自执行的 Apache Beam JAR 包路径,可以是本地绝对路径或 GCSgs://路径;若在 GCS 上,Operator 会先下载到本地再执行(见 execute);job_class:要执行的 Beam 管道主类名,通常不是 JAR 清单中配置的 main class;- 其余公共参数(
runner、pipeline_options、dataflow_config、deferrable)与 Python 版本一致。
7.2 底层执行方式
从 start_java_pipeline 可以看到实际命令构造:
command_prefix = ["java", "-cp", jar, job_class] if job_class else ["java", "-jar", jar]即:指定job_class时执行java -cp <jar> <job_class>,否则执行java -jar <jar>。运行环境需要 Java 运行时可用。
7.3 完整 DAG 示例(DirectRunner 与 DataflowRunner)
取自 example_beam.py 与 example_java_dataflow.py:
from airflow.providers.apache.beam.operators.beam import BeamRunJavaPipelineOperator from airflow.providers.google.cloud.transfers.gcs_to_local import GCSToLocalFilesystemOperator # 先从 GCS 拉取 JAR 到本地(文件名可含模板) jar_to_local = GCSToLocalFilesystemOperator( task_id="jar_to_local_direct_runner", bucket=GCS_JAR_DIRECT_RUNNER_BUCKET_NAME, object_name=GCS_JAR_DIRECT_RUNNER_OBJECT_NAME, filename="/tmp/beam_wordcount_direct_runner_{{ ds_nodash }}.jar", ) # DirectRunner start_java_pipeline_direct_runner = BeamRunJavaPipelineOperator( task_id="start_java_pipeline_direct_runner", jar="/tmp/beam_wordcount_direct_runner_{{ ds_nodash }}.jar", pipeline_options={ "output": "/tmp/start_java_pipeline_direct_runner", "inputFile": GCS_INPUT, }, job_class="org.apache.beam.examples.WordCount", ) # DataflowRunner(dataflow_config 使用 dict 写法) start_java_pipeline_dataflow = BeamRunJavaPipelineOperator( task_id="start_java_pipeline_dataflow", runner="DataflowRunner", jar="/tmp/beam_wordcount_dataflow_runner_{{ ds_nodash }}.jar", pipeline_options={ "tempLocation": GCS_TMP, "stagingLocation": GCS_STAGING, "output": GCS_OUTPUT, }, job_class="org.apache.beam.examples.WordCount", dataflow_config={"job_name": "{{task.task_id}}", "location": "us-central1"}, ) jar_to_local >> start_java_pipeline_direct_runnerJava Operator 在 Dataflow 模式下还支持check_if_running检查:若同名 Dataflow 流式作业已处于 RUNNING 状态,可以选择不重复提交(CheckJobRunning.WaitForRun语义),避免流式作业重名冲突。
八、运行 Go 流水线:BeamRunGoPipelineOperator
8.1 核心参数
go_file:Beam 管道 Go 源码路径,如/local/path/to/main.go或gs://bucket/path/to/main.go;从 GCS 拉取时,Operator 会先执行go mod init初始化模块、再用go mod tidy安装依赖,随后以go run <go_file>运行(等价于本地文件的执行方式);launcher_binary:为启动平台编译的 Go 可执行二进制(本地路径或gs://路径);worker_binary:为 Worker 平台编译的二进制,用于启动平台与 Worker 平台 OS/架构不同的跨编译场景(参考 Beam 的 Go cross-compilation 文档);若未设置则默认取launcher_binary的值,且仅在设置了launcher_binary时才有意义;go_file与launcher_binary必须且只能提供一个,否则execute会抛出ValueError(见 execute)。
8.2 环境前置条件
Go 流水线要求运行环境安装go命令,否则 start_go_pipeline 会抛出AirflowConfigException提示安装 Go;同时需要 Google Provider 提供go_module_utils支持。另外,Go SDK 目前不支持impersonation_chain参数(Operator 会记录告警并跳过)。使用launcher_binary/worker_binary时,若二进制位于 GCS,Operator 会用线程池并行下载并为下载的二进制添加可执行权限(见 _GoBinary.download_from_gcs)。
8.3 完整 DAG 示例
取自 example_go.py:
from airflow.providers.apache.beam.operators.beam import BeamRunGoPipelineOperator from airflow.providers.google.cloud.operators.dataflow import DataflowConfiguration # 本地 Go 源码 + DirectRunner start_go_pipeline_local_direct_runner = BeamRunGoPipelineOperator( task_id="start_go_pipeline_local_direct_runner", go_file="files/apache_beam/examples/wordcount.go", ) # GCS Go 源码 + DirectRunner start_go_pipeline_direct_runner = BeamRunGoPipelineOperator( task_id="start_go_pipeline_direct_runner", go_file=GCS_GO, pipeline_options={"output": GCS_OUTPUT}, ) # GCS Go 源码 + DataflowRunner(需指定 WorkerHarnessContainerImage) start_go_pipeline_dataflow_runner = BeamRunGoPipelineOperator( task_id="start_go_pipeline_dataflow_runner", runner="DataflowRunner", go_file=GCS_GO, pipeline_options={ "tempLocation": GCS_TMP, "stagingLocation": GCS_STAGING, "output": GCS_OUTPUT, "WorkerHarnessContainerImage": "apache/beam_go_sdk:latest", }, dataflow_config=DataflowConfiguration( job_name="{{task.task_id}}", project_id=GCP_PROJECT_ID, location="us-central1" ), )注意:Go 管道跑 Dataflow 时需要在pipeline_options中显式指定WorkerHarnessContainerImage(示例使用官方镜像apache/beam_go_sdk:latest)。
九、可延迟(Deferrable)异步执行模式
三个 Operator 均支持deferrable=True的异步模式,其价值在于:当任务进入等待阶段时,Worker 槽位被释放,由 Trigger 在后台轮询状态,集群不再被空闲 Operator 或 Sensor 占用,显著降低资源浪费。
- 对于非 Dataflow 的本地流水线,Operator 会
defer到 BeamPythonPipelineTrigger(Python)或BeamJavaPipelineTrigger(Java),Trigger 内部使用BeamAsyncHook通过 asyncio 子进程异步执行并持续读取日志(见 hooks/beam.py 中的BeamAsyncHook); - 对于 Dataflow 作业,Operator 先启动作业并拿到 job id,随后委托给
DataflowJobStateCompleteTrigger(Google Provider 支持时)或DataflowJobStatusTrigger(期望状态JOB_STATE_DONE)等待作业完成; - 恢复执行时调用
execute_complete(context, event):若事件状态为error则抛出AirflowException,否则记录完成信息(见 execute_complete)。
异步模式示例(取自 example_python_async.py),与同步写法唯一区别是增加deferrable=True:
start_python_pipeline_local_direct_runner = BeamRunPythonPipelineOperator( task_id="start_python_pipeline_local_direct_runner", py_file="apache_beam.examples.wordcount", py_options=["-m"], py_requirements=["apache-beam[gcp]==2.59.0"], py_interpreter="python3", py_system_site_packages=False, deferrable=True, )十、底层执行机制:子进程与参数转换
无论哪种语言,最终都会走到 run_beam_command:使用subprocess.Popen启动命令(shell=False),通过select.select同时监听 stdout/stderr,把 stderr 以 warning 级别、stdout 以 info 级别写入任务日志,并在日志流中调用回调解析 Dataflow job id;进程退出码非 0 时抛出AirflowException。同步模式如此,异步模式则用asyncio.create_subprocess_shell配合readline任务实现等价行为。
参数转换的核心函数是 beam_options_to_args,其逻辑与 Apache Beam Python SDK 的pipeline_options.py保持兼容。所有命令最后都会统一追加--runner=<runner>参数。此外,start_python_pipeline在启动前会通过import apache_beam; print(apache_beam.__version__)探测 Beam 版本并记录日志,且当使用impersonate_service_account选项时会校验 Beam 版本必须 >= 2.39.0。
十一、单元测试与进一步阅读
- 单元测试覆盖 Hook、Operator、Trigger 三个层次:tests/unit/apache/beam(
hooks/test_beam.py、operators/test_beam.py、triggers/test_beam.py),可用于验证参数转换、命令构造与触发逻辑; - 系统测试示例 DAG 位于 tests/system/apache/beam,包含 Python/Java/Go 各 Runner 组合(
example_python.py、example_python_async.py、example_python_dataflow.py、example_beam.py、example_beam_java_flink.py、example_beam_java_spark.py、example_java_dataflow.py、example_go.py、example_go_dataflow.py),是编写真实 DAG 的最佳参考模板; - 包级变更记录见 changelog.rst,安全公告见 security.rst,Python API 参考可通过包首页的 References 导航进入;
- 若需从源码构建安装该 Provider,可参考 installing-providers-from-sources.rst。
十二、小结
apache-airflow-providers-apache-beam6.2.4 把 Apache Beam 的多语言管道能力无缝嵌入 Airflow 调度体系:安装上要求 Airflow >= 2.11.0、Beam >= 2.76.0,Dataflow 场景需通过[google]extra 引入 Google Provider 并在 Worker 上安装gcloud;使用上,Python/Java/Go 三种 SDK 各有专用 Operator 与参数约定,GCS 文件自动下载、Dataflow 作业 ID 自动解析与 UI 链接、deferrable异步执行、任务 kill 自动取消作业等机制均由 operators/beam.py、hooks/beam.py 与 triggers/beam.py 在源码层完整支撑。参考仓库内系统测试示例即可快速落地第一个由 Airflow 编排的 Beam 数据管道。
【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考