资讯动态

Apache Airflow 集成 Apache Beam 管道:三大 Pipeline Operator 的完整实战指南

发布时间:2026/9/13 6:14:07 来源:尧图企业网站定制
Apache Airflow 集成 Apache Beam 管道三大 Pipeline Operator 的完整实战指南【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflowApache Airflow 通过apache-airflow-providers-apache-beam提供器将 Apache Beam 的数据并行处理管道接入 DAG 调度体系。本文基于该提供器的 operators.rst 官方指南结合 operator 源码 与 系统测试示例系统讲解BeamRunPythonPipelineOperator、BeamRunJavaPipelineOperator、BeamRunGoPipelineOperator的用法、关键参数、DirectRunner/DataflowRunner 双运行模式与 deferrable 异步执行帮助你在真实 DAG 中直接落地 Beam 批流管道。背景Beam 与 Airflow 的集成方式Apache Beam 是开源的统一批处理与流处理数据并行管道编程模型。开发者使用 Beam SDK 编写一段管道程序再由 Beam 支持的分布式执行后端如 Apache Flink、Apache Spark、Google Cloud Dataflow来运行。在 Airflow 中Beam 提供器位于providers/apache/beam/src/airflow/providers/apache/beam/其核心组件包括Operatorsoperators/beam.pyBeamRunPythonPipelineOperator、BeamRunJavaPipelineOperator、BeamRunGoPipelineOperator以及共享逻辑基类BeamBasePipelineOperator与BeamDataflowMixinHookhooks/beam.pyBeamHook负责拼接命令行并启动管道子进程Triggerstriggers/beam.pyBeamPythonPipelineTrigger、BeamJavaPipelineTrigger支撑 deferrable 异步模式。重要前提当 Beam 管道运行在 Google Cloud Dataflow 服务上时Airflow Worker 节点上必须安装gcloud命令行工具Google Cloud SDK。若仅在本地以 DirectRunner 运行则无此要求。三大 Operator 的公共执行模型三个 Operator 都继承自抽象基类BeamBasePipelineOperator共享以下核心执行逻辑runner参数指定运行器默认值为DirectRunner。从 hooks/beam.py 中的BeamRunnerType可以看到支持的全部运行器DataflowRunner、DirectRunner、SparkRunner、FlinkRunner、SamzaRunner、NemoRunner、JetRunner、Twister2Runner。default_pipeline_options与pipeline_options两个字典会在执行时合并default_pipeline_options适合存放作用于 DAG 内所有 Beam Operator 的高层选项如 project、zonepipeline_options存放单任务特有选项见 operator 源码中的_init_pipeline_options。gcp_conn_id默认为google_cloud_default用于连接 GCS 下载管道文件。当runnerDataflowRunner时还需通过dataflow_configDataflowConfiguration对象或字典提供 Dataflow 专属配置此时还会为任务附加DataflowJobLink快捷链接并在执行期间自动为 pipeline 打上airflow-version标签。pipeline_options 的值类型语义BeamHook中的beam_options_to_args负责把选项字典转成命令行参数其转换规则值得特别注意值为None该选项会被显式跳过不产生--keyNone之类的错误参数值为False选项被跳过但有两个例外——use_public_ips会转成--no_use_public_ipsPython SDKusePublicIps会转成--usePublicIpsfalseJava SDK以确保显式关闭能力生效值为True追加单参数--key无值值为列表为每个元素追加一个选项如key[A,B]生成--keyA --keyB值为字典序列化为 JSON 追加--key{a:1}可用于 labels 等复合选项其他类型以 Python 文本表示形式追加--keyvalue。一、运行 Python 管道BeamRunPythonPipelineOperatorBeamRunPythonPipelineOperator用于启动用 Python 编写的 Beam 管道其完整定义与文档字符串见 operator 源码。核心参数说明参数说明默认值py_file必填。指向待执行 Beam 管道的 Python 文件。可以是 GCS 上的gs://对象Airflow 会自动下载也可以是本地文件系统的绝对路径支持模板渲染无py_interpreter执行管道时使用的 Python 解释器python3若 Airflow 实例运行在 Python 2 环境可指定python2但官方建议使用 Python 3 以获得最佳效果py_options附加的 Python 选项如[-m, -v]配合py_fileapache_beam.examples.wordcount可按模块方式运行内置示例[]py_requirements附加 Python 包列表。一旦指定Operator 会创建一个临时虚拟环境并安装这些依赖再在虚拟环境中运行管道也可借此安装或指定特定版本的apache_beamNonepy_system_site_packages是否让虚拟环境在指定了py_requirements时继承 Airflow 实例的所有 Python 包。除非 Dataflow 作业确有需要否则建议保持默认关闭Falsedeferrable是否以 deferrable异步模式运行默认读取配置项operators.default_deferrableFalse从 execute 方法 可以看到若py_file以gs://开头Operator 会用GCSHook将文件下载到本地临时文件后再启动管道pipeline_options中的requirements_file同样支持gs://路径会被下载为本地文件后交给 Beam。DirectRunner本地运行 Python 管道以下示例来自 example_python.py分别展示本地模块运行与 GCS 文件运行两种方式from airflow import models from airflow.providers.apache.beam.operators.beam import BeamRunPythonPipelineOperator # 方式一本地 Python 模块 py_options[-m] start_python_pipeline_local_direct_runner BeamRunPythonPipelineOperator( task_idstart_python_pipeline_local_direct_runner, py_fileapache_beam.examples.wordcount, py_options[-m], py_requirements[apache-beam[gcp]2.59.0], py_interpreterpython3, py_system_site_packagesFalse, ) # 方式二GCS 上的 py_file经 pipeline_options 传入输出位置 start_python_pipeline_direct_runner BeamRunPythonPipelineOperator( task_idstart_python_pipeline_direct_runner, py_fileGCS_PYTHON, # 形如 gs://bucket/wordcount_debugging.py py_options[], pipeline_options{output: GCS_OUTPUT}, py_requirements[apache-beam[gcp]2.59.0], py_interpreterpython3, py_system_site_packagesFalse, )其中GCS_PYTHON、GCS_OUTPUT等常量定义在 utils.py默认值为占位符gs://INVALID BUCKET NAME/...实际使用时需要通过环境变量如APACHE_BEAM_PYTHON、APACHE_BEAM_GCS_OUTPUT注入真实 GCS 路径。DataflowRunner把 Python 管道提交到 Google Cloud Dataflow将runner切换为DataflowRunner并配合dataflow_config即可把管道提交到 Dataflow 服务见 example_python.pyfrom airflow.providers.google.cloud.operators.dataflow import DataflowConfiguration start_python_pipeline_dataflow_runner BeamRunPythonPipelineOperator( task_idstart_python_pipeline_dataflow_runner, runnerDataflowRunner, py_fileGCS_PYTHON, pipeline_options{ tempLocation: GCS_TMP, stagingLocation: GCS_STAGING, output: GCS_OUTPUT, }, py_options[], py_requirements[apache-beam[gcp]2.59.0], py_interpreterpython3, py_system_site_packagesFalse, dataflow_configDataflowConfiguration( job_name{{task.task_id}}, # 支持 Jinja 模板 project_idGCP_PROJECT_ID, locationus-central1, ), )deferrable 模式释放 Worker 插槽Beam Python Operator 支持 deferrable异步模式当任务进入等待状态时Airflow 将执行权交给 TriggerWorker 插槽被释放从而显著减少集群中空闲 Operator 长期占用的资源。使用方式只需在参数中加deferrableTrue完整示例见 example_python_async.py其中既包含 DirectRunner 的本地/GCS 文件变体也包含 DataflowRunner 变体。从源码看deferrable 的实现分为两条路径非 Dataflow 运行器Operator 直接调用self.defer()挂起BeamPythonPipelineTriggeroperators/beam.py由 Trigger 负责启动管道进程并异步轮询完成状态DataflowRunnerOperator 先同步提交作业execute_on_dataflow解析出dataflow_job_id后挂起DataflowJobStateCompleteTrigger若 Google 提供器支持或回退到DataflowJobStatusTrigger等待JOB_STATE_DONEoperators/beam.py。Dataflow 任务的额外能力XCom 推送执行期间一旦拿到作业 IDdataflow_job_id属性会立即通过 XCom 以dataflow_job_id键推送给下游operators/beam.py因此下游的DataflowJobStatusSensor可以在作业完成前就开始监听执行结果execute_on_dataflow返回{dataflow_job_id: self.dataflow_job_id}优雅取消实现on_kill任务被 kill 时调用DataflowHook.cancel_job取消对应作业operators/beam.py。Dataflow 异步提交 Sensor 组合模式如果需要提交后不等待、由独立 Sensor 轮询的完全异步模式可参考 example_python_dataflow.py在dataflow_config中设置wait_until_finishedFalse然后用DataflowJobStatusSensor通过 XCom 拉取dataflow_job_id并轮询JOB_STATE_DONEstart_python_job_dataflow_runner_async BeamRunPythonPipelineOperator( task_idstart_python_job_dataflow_runner_async, runnerDataflowRunner, py_fileGCS_PYTHON_DATAFLOW_ASYNC, pipeline_options{ tempLocation: GCS_TMP, stagingLocation: GCS_STAGING, output: GCS_OUTPUT, }, dataflow_configDataflowConfiguration( job_name{{task.task_id}}, project_idGCP_PROJECT_ID, locationus-central1, wait_until_finishedFalse, ), ) wait_for_python_job_dataflow_runner_async_done DataflowJobStatusSensor( task_idwait-for-python-job-async-done, job_id{{task_instance.xcom_pull(start_python_job_dataflow_runner_async)[dataflow_job_id]}}, expected_statuses{DataflowJobStatus.JOB_STATE_DONE}, project_idGCP_PROJECT_ID, locationus-central1, ) start_python_job_dataflow_runner_async wait_for_python_job_dataflow_runner_async_done二、运行 Java 管道BeamRunJavaPipelineOperatorBeamRunJavaPipelineOperator用于启动 Java 编写的 Beam 管道定义见 operators/beam.py。核心参数jar必填。指向自执行self-executingBeam JAR 的引用支持 GCS 路径自动下载或本地绝对路径支持模板渲染job_class要执行的 Beam 管道类名它通常不是 JAR 内配置的 main class例如示例中的org.apache.beam.examples.WordCountpipeline_options向作业传递的管道选项。DirectRunner 示例以下来自 example_beam.py先用GCSToLocalFilesystemOperator把 JAR 从 GCS 下载到本地文件名带{{ ds_nodash }}模板再交给 Beam Operatorfrom airflow.providers.apache.beam.operators.beam import BeamRunJavaPipelineOperator from airflow.providers.google.cloud.transfers.gcs_to_local import GCSToLocalFilesystemOperator jar_to_local_direct_runner GCSToLocalFilesystemOperator( task_idjar_to_local_direct_runner, bucketGCS_JAR_DIRECT_RUNNER_BUCKET_NAME, object_nameGCS_JAR_DIRECT_RUNNER_OBJECT_NAME, filename/tmp/beam_wordcount_direct_runner_{{ ds_nodash }}.jar, ) start_java_pipeline_direct_runner BeamRunJavaPipelineOperator( task_idstart_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_classorg.apache.beam.examples.WordCount, ) jar_to_local_direct_runner start_java_pipeline_direct_runnerDataflowRunner 示例example_java_dataflow.py 展示了 Dataflow 场景dataflow_config既可以传DataflowConfiguration对象也可以直接传字典start_java_pipeline_dataflow BeamRunJavaPipelineOperator( task_idstart_java_pipeline_dataflow, runnerDataflowRunner, jar/tmp/beam_wordcount_dataflow_runner_{{ ds_nodash }}.jar, pipeline_options{ tempLocation: GCS_TMP, stagingLocation: GCS_STAGING, output: GCS_OUTPUT, }, job_classorg.apache.beam.examples.WordCount, dataflow_config{job_name: {{task.task_id}}, location: us-central1}, )Java 特有的执行细节从 execute / execute_on_dataflow 源码 可见 Java Operator 与 Python 的差异jar以gs://开头时同样先下载到本地临时文件Dataflow 模式下支持dataflow_config.check_if_running去重逻辑当check_if_running WaitForRun时先查询同名作业是否已在运行若发现同名 streaming 作业处于 RUNNING 状态会停止执行并返回{dataflow_job_id: None}提示修改job_name或设置check_if_runningFalse防止重复提交pipeline_options[jobName]会在启动前被自动设置为生成的dataflow_job_name同样支持deferrableTrue、DataflowJobLink、on_kill取消与 XCom 推送。三、运行 Go 管道BeamRunGoPipelineOperatorBeamRunGoPipelineOperator用于启动 Go 编写的 Beam 管道定义与文档见 operators/beam.py。核心参数go_fileBeam 管道 Go 源文件引用如/local/path/to/main.go或gs://bucket/path/to/main.golauncher_binary为启动平台编译的 Go 管道二进制如/local/path/to/launcher-main或gs://bucket/path/to/launcher-mainworker_binary为 Worker 平台编译的二进制。当 Worker 的 OS/架构与启动平台不同时需要提供对应 Beam Go 交叉编译场景若未设置则默认等于launcher_binary的值若launcher_binary未设置提供worker_binary不生效约束go_file与launcher_binary必须且只能提供一个否则执行时抛出ValueError(Exactly one of go_file and launcher_binary must be set)operators/beam.py。Go 文件的执行方式官方指南指出本地文件系统上的 Go 源文件其执行等价于go run go_file若从 GCS 拉取则执行前会先初始化模块并安装依赖即依次执行go mod init example.com/main与go mod tidy。这在源码中对应_GoFile.should_init_go_module标志仅当文件来自 GCS 时才置为True从而触发模块初始化operators/beam.py。DirectRunner 示例example_go.py 给出本地文件与 GCS 文件两种 DirectRunner 用法from airflow.providers.apache.beam.operators.beam import BeamRunGoPipelineOperator # 本地 Go 源文件 start_go_pipeline_local_direct_runner BeamRunGoPipelineOperator( task_idstart_go_pipeline_local_direct_runner, go_filefiles/apache_beam/examples/wordcount.go, ) # GCS 上的 Go 源文件 start_go_pipeline_direct_runner BeamRunGoPipelineOperator( task_idstart_go_pipeline_direct_runner, go_fileGCS_GO, pipeline_options{output: GCS_OUTPUT}, )DataflowRunner 示例Dataflow 场景需要额外指定WorkerHarnessContainerImageGo SDK 运行在 Dataflow 上时使用apache/beam_go_sdk:latest镜像见 example_go.pyfrom airflow.providers.google.cloud.operators.dataflow import DataflowConfiguration start_go_pipeline_dataflow_runner BeamRunGoPipelineOperator( task_idstart_go_pipeline_dataflow_runner, runnerDataflowRunner, go_fileGCS_GO, pipeline_options{ tempLocation: GCS_TMP, stagingLocation: GCS_STAGING, output: GCS_OUTPUT, WorkerHarnessContainerImage: apache/beam_go_sdk:latest, }, dataflow_configDataflowConfiguration( job_name{{task.task_id}}, project_idGCP_PROJECT_ID, locationus-central1 ), )Go 特有细节Go Operator不支持impersonation若dataflow_config.impersonation_chain被设置会打印警告并在执行中跳过operators/beam.pyGCS 上的二进制launcher_binary/worker_binary会被下载到临时目录并赋予可执行权限_make_executable下载过程使用ThreadPoolExecutor并行拉取operators/beam.py若launcher_binary与worker_binary指向同一路径只下载一次两者复用同一文件Go 管道同样支持wait_until_finishedFalse异步提交 DataflowJobStatusSensor组合见 example_go_dataflow.py。四、Dataflow 配置详解DataflowConfiguration当runnerDataflowRunner时dataflow_config接受DataflowConfiguration对象或字典其关键字段结合 BeamDataflowMixin 源码 与 Google 提供器的DataflowConfiguration包括字段作用job_nameDataflow 作业名未设置时默认取task_id支持模板project_idGCP 项目 ID未设置时回退为 DataflowHook 探测到的项目location作业区域如us-central1未设置时使用 Google 提供器的默认区域wait_until_finished是否等待作业结束False即异步提交模式poll_sleep轮询间隔impersonation_chain服务账号模拟链会写入 pipeline optionimpersonateServiceAccountdrain_pipeline、cancel_timeout取消/排空作业时的行为控制check_if_runningJava Operator 特有的同名作业查重策略WaitForRun等multiple_jobsJava 作业名匹配到多个作业时是否全部等待append_job_name是否追加时间戳等后缀生成唯一作业名service_account写入 pipeline optionserviceAccountBeamDataflowMixin在提交前会统一完成以下增强operators/beam.py把project、region写入 pipeline options把serviceAccount、impersonateServiceAccount注入 pipeline options自动合并labels追加airflow-version标签如v3-0-0格式便于在 Dataflow 控制台识别作业来源通过process_line_and_extract_dataflow_job_id_callback从子进程输出中解析并持续更新dataflow_job_id一旦解析到 ID 立即 XCom 推送。五、运行前置条件与注意事项依赖安装需要安装apache-airflow-providers-apache-beam若使用 DataflowRunner 还需安装apache-airflow-providers-google源码中BeamDataflowMixin.__init__在缺少 Google 提供器时会抛出AirflowOptionalProviderFeatureException明确提示见 operators/beam.py。gcloud 工具Dataflow 场景下 Worker 需安装 Google Cloud SDK 的gcloud命令并完成 GCP 认证连接 IDgoogle_cloud_default对应环境变量与凭据配置。GCS 文件下载py_file、jar、go_file、launcher_binary、worker_binary及requirements_file均支持gs://前缀Operator 会在执行前用GCSHook下载到本地临时文件示例代码中的GCS_*常量为占位符需通过环境变量替换为真实路径。Python 2 兼容若 Airflow 运行在 Python 2需将py_interpreter指定为python2并保证py_file为 Python 2 代码强烈建议统一使用 Python 3。虚拟环境隔离py_requirements会创建临时虚拟环境安装依赖py_system_site_packagesTrue会继承 Airflow 实例的全部 Python 包除非 Dataflow 作业必需否则不建议开启。运行器能力差异不同 RunnerDirectRunner / DataflowRunner / SparkRunner / FlinkRunner / PortableRunner 等对管道能力的支持各不相同选型前应参考 Apache Beam 官方的 Runner Capability Matrix。参考与延伸阅读提供器索引与安装说明providers/apache/beam/docs/index.rst、installing-providers-from-sources.rst完整系统测试示例 DAGexample_python.py、example_python_async.py、example_python_dataflow.py、example_beam.py、example_java_dataflow.py、example_go.py、example_go_dataflow.py全部位于 providers/apache/beam/tests/system/apache/beam/Operator 单元测试tests/unit/apache/beam/operators/test_beam.pyHook 与 Trigger 实现hooks/beam.py、triggers/beam.pyApache Beam 官方文档管道模型、各 Runner 使用说明与 Dataflow 命令行/监控接口可作为进一步学习的背景资料具体以官方站点的 Documentation 与 Google Cloud Dataflow 产品文档为准。【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

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

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

免费获取报价