
Apache Airflow Apache Beam Provider 实战指南Python / Java / Go 管道调度与 Dataflow 集成【免费下载链接】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 数据管道编排进 Airflow DAG 的一等公民能力。本指南围绕该 Provider 的安装配置、三种语言Python / Java / Go的管道运行运算符、DirectRunner 与 DataflowRunner 两种执行模式以及 deferrable可延迟异步执行展开读者读完即可在自己的 Airflow 集群中接入 Beam 管道并深入理解其底层 Hook、Trigger 与 Dataflow 配置的协作机制。Provider 概览Apache Beam 是一个用于定义批处理和流式数据并行处理管道的统一开源模型支持使用 Beam SDK 编写管道程序再交由 Apache Flink、Apache Spark、Google Cloud Dataflow 等分布式处理后端执行。本 Provider 将这一能力封装为 Airflow 运算符使得 Beam 管道的启动、监控、日志追踪可以完全融入 Airflow 的调度与重试体系。在仓库中该 Provider 的元数据定义于 provider.yaml其状态为ready、生命周期为production当前版本为6.2.4。Provider 对外暴露的组件包括运算符airflow.providers.apache.beam.operators.beam模块中的BeamRunPythonPipelineOperator、BeamRunJavaPipelineOperator、BeamRunGoPipelineOperatorHookairflow.providers.apache.beam.hooks.beam模块中的BeamHook负责实际执行管道命令Triggerairflow.providers.apache.beam.triggers.beam模块中的BeamPythonPipelineTrigger、BeamJavaPipelineTrigger支撑 deferrable 异步执行。安装与依赖要求基础安装在已有 Airflow 安装最低支持版本见下之上通过 pip 直接安装pip install apache-airflow-providers-apache-beam该包支持的 Python 版本为 3.10、3.11、3.12、3.13、3.14。版本依赖根据仓库内 README.rst 的 Requirements 表核心依赖及其版本约束如下PIP 包版本要求apache-airflow2.11.0apache-airflow-providers-common-compat1.12.0apache-beam2.76.0pyarrow16.1.0; python_version 3.1422.0.0; python_version 3.14numpy依 Python 版本而定3.11需1.22.43.11 3.12需1.23.23.12 3.14需1.26.03.14需2.4.3注意pyarrow与numpy的版本要求随 Python 版本分化安装时 pip 会根据解释器版本自动选择合适约束。可选跨 Provider 依赖部分功能需要额外安装 Google Provider例如将管道跑在 Dataflow 服务上。通过 extra 安装pip install apache-airflow-providers-apache-beam[google]该 extra 对应apache-airflow-providers-google与apache-beam[gcp]2.76.0。从源码实现看beam.py运算符在导入时会探测 Google Provider 是否安装未安装时使用 Dataflow 相关功能将抛出AirflowOptionalProviderFeatureException提示先安装对应版本的 google provider。通用执行模型Runner 与 Pipeline Options所有 Beam 管道运算符都继承自抽象基类BeamBasePipelineOperatorbeam.py其核心参数在两个层面控制管道执行runner管道运行的后端默认为DirectRunner可选项还包括DataflowRunner、SparkRunner、FlinkRunner、PortableRunnerdefault_pipeline_options与pipeline_options两组合并后作为 Beam 管道执行参数。default_pipeline_options用于存放对整个 DAG 中所有 Beam 任务通用的高层选项如项目与区域而pipeline_options存放单个任务特有选项。值的类型决定了命令行参数的产生方式见 beam.py值为None或True生成单个--key无值选项值为False跳过该选项值为列表如[A, B]生成--keyA --keyB多条选项其他类型替换为 Python 文本表示。管道选项中的 key 会被统一转换为 snake_case 后传入执行层_init_pipeline_options(format_pipeline_optionsTrue)。此外运算符会自动为管道注入airflow-version标签便于在 Dataflow 控制台识别任务来源。运行 Python 管道BeamRunPythonPipelineOperatorBeamRunPythonPipelineOperator用于启动 Python 编写的 Beam 管道。其必填参数py_file指定包含管道逻辑的 Python 文件可以位于本地文件系统提供绝对路径或 GCSAirflow 会自动下载到本地临时文件再执行。关键参数说明py_file管道文件路径支持模板渲染templatedpy_interpreter执行管道的 Python 解释器默认python3。若 Airflow 实例运行于 Python 2需指定python2并确保py_file使用 Python 2 语法——不过官方推荐尽量使用 Python 3py_options附加的 Python 选项例如[-m, -v]传入后等价于python -m ...的调用方式py_requirements若指定将创建一个临时 Python 虚拟环境并安装所列依赖管道在其中运行。典型用法是把apache-beam[gcp]装进虚拟环境避免与 Airflow 宿主环境互相污染py_system_site_packages是否让虚拟环境共享 Airflow 实例的全部 Python 包。仅在设置py_requirements时生效官方建议避免开启除非 Dataflow 任务确实需要gcp_conn_id当py_file位于 GCS 时用于连接 Google Cloud Storage 的连接 ID默认google_cloud_defaultdataflow_config当 runner 为DataflowRunner时的 Dataflow 配置类型为DataflowConfiguration或字典deferrable是否启用可延迟模式默认读取 Airflow 配置[operators] default_deferrable。DirectRunner 本地执行以下示例运行本地 Python 文件通过模块方式调用 wordcount 示例并基于py_requirements创建虚拟环境见 example_python.pystart_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, )当py_file位于 GCS 时使用pipeline_options传入输出参数见 example_python.pystart_python_pipeline_direct_runner BeamRunPythonPipelineOperator( task_idstart_python_pipeline_direct_runner, py_fileGCS_PYTHON, # gs://... 路径Airflow 自动下载 py_options[], pipeline_options{output: GCS_OUTPUT}, py_requirements[apache-beam[gcp]2.59.0], py_interpreterpython3, py_system_site_packagesFalse, )DataflowRunner 云端执行切到DataflowRunner时需通过dataflow_config指定任务名、项目与区域并配合tempLocation、stagingLocation等选项见 example_python.pystart_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}}, project_idGCP_PROJECT_ID, locationus-central1 ), )需要特别注意的是当 Beam 管道运行在 Dataflow 服务上时Airflow worker 必须安装gcloud命令行工具Google Cloud SDK见 operators.rst 中的说明。从源码看Dataflow 分支的执行流程为beam.pyBeamDataflowMixin._set_dataflow()构建DataflowHook从管道输出行中提取 Dataflow job idprocess_line_and_extract_dataflow_job_id_callback把 job id 通过 XCom 以dataflow_job_id键推送随后调用DataflowJobLink.persist()记录监控链接再wait_for_done()阻塞等待任务完成。因此下游任务可以通过ti.xcom_pull(task_ids..., keydataflow_job_id)拿到 job id配合 Dataflow 监控接口进行后续追踪。运行 Java 管道BeamRunJavaPipelineOperatorBeamRunJavaPipelineOperator用于启动 Java 编写的 Beam 管道必填参数jar指定包含管道逻辑的 JAR 包同样支持本地绝对路径或 GCS 路径。DirectRunner 执行以下 DAG 先用GCSToLocalFilesystemOperator把 JAR 从 GCS 下载到本地再运行管道见 example_beam.py。注意jar中使用了{{ ds_nodash }}模板变量说明该参数支持模板渲染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 执行Java 管道上云时dataflow_config可以直接传字典见 example_java_dataflow.py。从源码的_cast_dataflow_config()可以看出字典形式会被内部转换为DataflowConfiguration对象未显式指定job_name时默认使用task_idstart_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}, )运行 Go 管道BeamRunGoPipelineOperatorBeamRunGoPipelineOperator用于启动 Go 编写的 Beam 管道必填参数go_file指定 Go 源文件支持本地绝对路径或 GCS 路径。从文档说明operators.rst可知其执行语义本地文件等价于go run go_fileGCS 文件下载后先执行go run init example.com/main初始化模块再执行go mod tidy安装依赖最后运行管道。DirectRunner 执行# 本地文件见 example_go.py 对应段落 start_go_pipeline_local_direct_runner BeamRunGoPipelineOperator( task_idstart_go_pipeline_local_direct_runner, go_filefiles/apache_beam/examples/wordcount.go, ) # GCS 文件 start_go_pipeline_direct_runner BeamRunGoPipelineOperator( task_idstart_go_pipeline_direct_runner, go_fileGCS_GO, pipeline_options{output: GCS_OUTPUT}, )DataflowRunner 执行Go 管道跑 Dataflow 时需要额外指定WorkerHarnessContainerImage指向 Beam Go SDK 镜像见 example_go.pystart_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 管道同样支持 SparkRunner、FlinkRunner例如pipeline_options{endpoint: /your/spark/endpoint}配置 Spark 集群端点。可延迟模式Deferrable与异步执行Beam 管道往往运行数分钟到数小时传统阻塞式执行会长时间占用 worker 槽位。三个运算符均支持deferrableTrue在需要等待时把任务挂起deferred将恢复职责移交给 Trigger从而释放 worker 资源减少集群在空闲 Operator / Sensor 上的浪费。以 Python 管道为例见 example_python_async.py只需增加一个参数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, deferrableTrue, )从源码实现看deferrable 的挂起逻辑分为两种情况beam.py非 Dataflow Runner任务直接通过self.defer(triggerBeamPythonPipelineTrigger(...))挂起由 Beam Python 管道 Trigger 异步监控本地/远程进程的完成状态Dataflow Runner任务先同步启动管道并拿到dataflow_job_id再挂起。触发器的选择取决于 Google Provider 的版本优先使用DataflowJobStateCompleteTrigger配合wait_until_finished不可用时回退到DataflowJobStatusTrigger等待JOB_STATE_DONE。挂起期间任务不占用 worker 槽位Trigger 触发后任务通过execute_complete()恢复事件状态为error时抛出AirflowException使任务失败beam.py。在 DAG 中组合编排多个 Beam 任务可以在同一 DAG 中自由组合例如将 DirectRunner 与 DataflowRunner 串成依赖链见 example_python.py( [ start_python_pipeline_local_direct_runner, start_python_pipeline_direct_runner, ] start_python_pipeline_dataflow_runner start_python_pipeline_local_flink_runner start_python_pipeline_local_spark_runner )仓库中的系统测试示例 DAG 均以get_test_run(dag)结尾见各示例文件末尾这使示例既能作为独立 DAG 运行也能通过 pytest 纳入系统测试体系。若要在本地跑通这些示例需按 utils.py 的约定配置 GCS 路径、项目 ID 等环境变量并保证 Airflow 与 GCP 之间的认证gcloud 凭证、连接google_cloud_default等就绪。参考与进一步阅读完整运算符指南见 operators.rst其中包含全部示例 DAG 的引用片段变更记录见 changelog.rst可用于追溯各版本行为变化安全说明与配置项见 security.rst相关实现源码运算符 beam.py、Hook hooks/beam.py、Trigger triggers/beam.py单元测试覆盖了 Hook、运算符与 Trigger 的关键行为见 tests/unit/apache/beam 目录。适用前提与限制本文描述基于仓库中当前 Provider 版本6.2.4运行 Dataflow 场景需要同时安装 Google Provider[google]extra并在 worker 上部署 gcloud SDK不同 runner 对 Beam SDK 的版本要求与能力矩阵以 Beam 官方文档为准。【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考