ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

Apache Airflow Apache Beam Provider 完全指南:安装、依赖与 Python/Java/Go 流水线编排实战

2026/9/13 23:59:38 拓冰建站 浏览量
Apache Airflow Apache Beam Provider 完全指南:安装、依赖与 Python/Java/Go 流水线编排实战 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/airflowApache Airflow 的apache-airflow-providers-apache-beam6.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.rstPIP packageVersion requiredapache-airflow2.11.0apache-airflow-providers-common-compat1.12.0apache-beam2.76.0pyarrow16.1.0; python_version 3.14pyarrow22.0.0; python_version 3.14numpy1.22.4; python_version 3.11numpy1.23.2; python_version 3.12 and python_version 3.11numpy1.26.0; python_version 3.12 and python_version 3.14numpy2.4.3; python_version 3.14可见依赖会随 Python 版本动态选择pyarrow与numpy的版本例如 Python 3.14 使用 pyarrow16.1.0Python 3.14 使用 pyarrow22.0.0numpy 则在 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。安装时指定 extrapip install apache-airflow-providers-apache-beam[google]Dependent packageExtraapache-airflow-providers-googlegoogle3.2 可选依赖googleextra 同时会引入 Beam 的 GCP 扩展ExtraDependenciesgoogleapache-beam[gcp]2.76.03.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、Twister2Runnerdefault_pipeline_options默认管道选项适合存放对 DAG 中所有 Beam Operator 通用的高层参数如 project、zonepipeline_options管道选项字典会与default_pipeline_options合并后传给 Beam。值类型决定生成的命令行参数None值该选项被显式跳过避免被当作字符串None传给 BeamTrue生成无值的--key选项False选项被跳过但对use_public_ips这类特殊标志会生成对应的取反标志--no_use_public_ips见 beam_options_to_args 中的_FLAG_THAT_SETS_FALSE_VALUE映射list为每个元素生成一个--keyvalue例如[A,B]生成--keyA --keyBdict序列化为 JSON 字符串传入常用于 labels其他类型使用 Python 文本表示。gcp_conn_id连接 GCS 使用的 Airflow 连接 ID默认google_cloud_defaultdataflow_configDataflow 专用配置DataflowConfiguration或等价 dict仅在 runner 为DataflowRunner时生效若 runner 不是 DataflowRunner 却配置了该项Operator 会打印告警日志。5.2 Dataflow 集成要点当runnerDataflowRunner时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 节点安装gcloudGoogle Cloud SDK命令行工具。六、运行 Python 流水线BeamRunPythonPipelineOperator6.1 核心参数py_file必填支持模板要执行的 Beam Python 管道文件可以是本地绝对路径也可以是 GCS 上的gs://路径Airflow 会通过GCSHook.provide_file下载为临时文件见 executepy_interpreter执行管道所用的 Python 版本默认python3。若 Airflow 实例运行在 Python 2 环境则指定python2并确保py_file为 Python 2 代码官方建议优先使用 Python 3py_options额外 Python 选项如[-m, -v]py_requirements指定后将创建临时 Python 虚拟环境并在其中安装这些依赖然后在此环境内运行管道——也可用来安装特定版本的apache-beampy_system_site_packages创建虚拟环境时是否包含 Airflow 实例的系统站点包。默认False除非 Dataflow 作业确实需要否则不建议开启deferrable是否以可延迟异步模式运行默认读取operators.default_deferrable配置默认 False。需要注意的参数组合约束来自 hooks/beam.py 的start_python_pipeline如果指定了py_requirements但列表为空、且py_system_site_packagesFalse会抛出异常——因为此时虚拟环境中没有apache-beam包作业无法执行。修复方法是要么在系统安装apache-beam并设置py_system_site_packagesTrue要么把apache-beam加入py_requirements列表。6.2 完整 DAG 示例DirectRunner以下示例取自系统测试 example_python.pyfrom airflow import models from airflow.providers.apache.beam.operators.beam import BeamRunPythonPipelineOperator with models.DAG( example_beam_native_python, start_dateSTART_DATE, scheduleNone, # 按需覆盖 catchupFalse, default_argsDEFAULT_ARGS, tags[example], ) as dag: # 本地文件 DirectRunner通过模块方式运行 wordcount 示例 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 文件 DirectRunner start_python_pipeline_direct_runner BeamRunPythonPipelineOperator( task_idstart_python_pipeline_direct_runner, py_fileGCS_PYTHON, # 形如 gs://bucket/path/to/pipeline.py py_options[], pipeline_options{output: GCS_OUTPUT}, py_requirements[apache-beam[gcp]2.59.0], py_interpreterpython3, py_system_site_packagesFalse, )6.3 DataflowRunner 示例from 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}}, project_idGCP_PROJECT_ID, locationus-central1 ), )job_name支持 Jinja 模板示例中使用{{task.task_id}}让作业名与任务 ID 一致dataflow_config也接受等价的 dict 写法。示例 DAG 中还展示了 SparkRunner、FlinkRunner 的用法只要把runner替换并传入对应 pipeline options如 Flink 的output、Spark 的endpoint即可。七、运行 Java 流水线BeamRunJavaPipelineOperator7.1 核心参数jar必填支持模板自执行的 Apache Beam JAR 包路径可以是本地绝对路径或 GCSgs://路径若在 GCS 上Operator 会先下载到本地再执行见 executejob_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.pyfrom 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_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, ) # DirectRunner 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, ) # DataflowRunnerdataflow_config 使用 dict 写法 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}, ) jar_to_local start_java_pipeline_direct_runnerJava Operator 在 Dataflow 模式下还支持check_if_running检查若同名 Dataflow 流式作业已处于 RUNNING 状态可以选择不重复提交CheckJobRunning.WaitForRun语义避免流式作业重名冲突。八、运行 Go 流水线BeamRunGoPipelineOperator8.1 核心参数go_fileBeam 管道 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时若二进制位于 GCSOperator 会用线程池并行下载并为下载的二进制添加可执行权限见 _GoBinary.download_from_gcs。8.3 完整 DAG 示例取自 example_go.pyfrom 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_idstart_go_pipeline_local_direct_runner, go_filefiles/apache_beam/examples/wordcount.go, ) # GCS Go 源码 DirectRunner start_go_pipeline_direct_runner BeamRunGoPipelineOperator( task_idstart_go_pipeline_direct_runner, go_fileGCS_GO, pipeline_options{output: GCS_OUTPUT}, ) # GCS Go 源码 DataflowRunner需指定 WorkerHarnessContainerImage 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 管道跑 Dataflow 时需要在pipeline_options中显式指定WorkerHarnessContainerImage示例使用官方镜像apache/beam_go_sdk:latest。九、可延迟Deferrable异步执行模式三个 Operator 均支持deferrableTrue的异步模式其价值在于当任务进入等待阶段时Worker 槽位被释放由 Trigger 在后台轮询状态集群不再被空闲 Operator 或 Sensor 占用显著降低资源浪费。对于非 Dataflow 的本地流水线Operator 会defer到 BeamPythonPipelineTriggerPython或BeamJavaPipelineTriggerJavaTrigger 内部使用BeamAsyncHook通过 asyncio 子进程异步执行并持续读取日志见 hooks/beam.py 中的BeamAsyncHook对于 Dataflow 作业Operator 先启动作业并拿到 job id随后委托给DataflowJobStateCompleteTriggerGoogle Provider 支持时或DataflowJobStatusTrigger期望状态JOB_STATE_DONE等待作业完成恢复执行时调用execute_complete(context, event)若事件状态为error则抛出AirflowException否则记录完成信息见 execute_complete。异步模式示例取自 example_python_async.py与同步写法唯一区别是增加deferrableTruestart_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, )十、底层执行机制子进程与参数转换无论哪种语言最终都会走到 run_beam_command使用subprocess.Popen启动命令shellFalse通过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保持兼容。所有命令最后都会统一追加--runnerrunner参数。此外start_python_pipeline在启动前会通过import apache_beam; print(apache_beam.__version__)探测 Beam 版本并记录日志且当使用impersonate_service_account选项时会校验 Beam 版本必须 2.39.0。十一、单元测试与进一步阅读单元测试覆盖 Hook、Operator、Trigger 三个层次tests/unit/apache/beamhooks/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.rstPython 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.0Dataflow 场景需通过[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),仅供参考