十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Apache Airflow TaskFlow API 实战:用纯 Python 函数编写 ETL 流水线(含 XCom、传感器与隔离环境详解)

Apache Airflow TaskFlow API 实战:用纯 Python 函数编写 ETL 流水线(含 XCom、传感器与隔离环境详解) Apache Airflow TaskFlow API 实战用纯 Python 函数编写 ETL 流水线含 XCom、传感器与隔离环境详解【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow本文基于 Airflow 官方教程 Pythonic Dags with the TaskFlow API 与配套示例 DAG 源码系统讲解如何用 TaskFlow API 以纯 Python 函数的方式编写 Airflow DAG包括dag/task装饰器用法、函数返回值经 XCom 自动传递数据、.override()参数化复用、虚拟环境/Docker/Kubernetes 隔离执行、task.sensor传感器以及模板化上下文变量与条件执行等高级模式。读完后你将能够写出比传统 Operator 风格更简洁、可维护的 Airflow 工作流并理解其背后 XCom 与依赖图的实现机制。一、总览一条完整的 TaskFlow ETL 流水线TaskFlow API 在 Airflow 2.0 中引入核心设计思路是你写普通的 Python 函数加上装饰器Airflow 负责其余一切——包括创建任务、建立依赖关系、在任务之间传递数据。官方教程以一条经典的 ETL 流水线Extract → Transform → Load为例对应的示例 DAG 源码位于 tutorial_taskflow_api.py。完整代码如下import json import pendulum from airflow.sdk import dag, task dag( scheduleNone, start_datependulum.datetime(2021, 1, 1, tzUTC), catchupFalse, tags[example], ) def tutorial_taskflow_api(): ### TaskFlow API Tutorial Documentation This is a simple data pipeline example which demonstrates the use of the TaskFlow API using three simple tasks for Extract, Transform, and Load. task() def extract(): #### Extract task A simple Extract task to get data ready for the rest of the data pipeline. In this case, getting data is simulated by reading from a hardcoded JSON string. data_string {1001: 301.27, 1002: 433.21, 1003: 502.22} order_data_dict json.loads(data_string) return order_data_dict task(multiple_outputsTrue) def transform(order_data_dict: dict): #### Transform task A simple Transform task which takes in the collection of order data and computes the total order value. total_order_value 0 for value in order_data_dict.values(): total_order_value value return {total_order_value: total_order_value} task() def load(total_order_value: float): #### Load task A simple Load task which takes in the result of the Transform task and instead of saving it to end user review, just prints it out. print(fTotal order value is: {total_order_value:.2f}) order_data extract() order_summary transform(order_data) load(order_summary[total_order_value]) tutorial_taskflow_api()上面这段代码就是整条流水线三个函数、三行调用Airflow 便能自动完成调度与编排。下面分步骤拆解其构成。二、Step 1用dag装饰器定义 DAGDAG 本质上仍是 Airflow 加载并解析的 Python 脚本但这里使用dag装饰器来定义。官方示例中 DAG 的定义部分对应源码 tutorial_taskflow_api.py 第32-37行dag( scheduleNone, start_datependulum.datetime(2021, 1, 1, tzUTC), catchupFalse, tags[example], ) def tutorial_taskflow_api(): ...为了让 Airflow 发现这个 DAG只需在模块级别调用被dag装饰的函数tutorial_taskflow_api()需要注意一个版本演进细节自 Airflow 2.4 起使用dag装饰器或以with块定义 DAG 时不再需要把 DAG 赋给全局变量Airflow 会自动发现它。DAG 加载后可以在 Airflow UI 的 Graph View 中直观查看任务之间的连接方式。三、Step 2用task编写任务在 TaskFlow 中每个任务就是一个普通 Python 函数加上task装饰器后Airflow 就能调度和执行它。以extract任务为例对应源码 tutorial_taskflow_api.py 第50-61行task() def extract(): #### Extract task A simple Extract task to get data ready for the rest of the data pipeline. In this case, getting data is simulated by reading from a hardcoded JSON string. data_string {1001: 301.27, 1002: 433.21, 1003: 502.22} order_data_dict json.loads(data_string) return order_data_dicttransform与load任务采用同样的模式。这里有几个关键点函数返回值会自动传给下游任务——不需要手动使用 XCom。TaskFlow 底层依然使用 XCom 管理数据传递只是把手动管理 XCom 的复杂性抽象掉了。multiple_outputsTrue的行为transform使用了task(multiple_outputsTrue)这告诉 Airflow 函数返回的是一个字典应将其拆分为独立的 XCom。字典中的每个 key 各自成为一个 XCom 条目下游任务可以直接引用特定值如order_summary[total_order_value]。如果省略multiple_outputsTrue整个字典会作为单个 XCom 存储只能整体访问。四、Step 3通过函数调用构建流程任务定义完成后像调用普通 Python 函数一样调用它们即可构建流水线对应源码 tutorial_taskflow_api.py 第96-98行order_data extract() order_summary transform(order_data) load(order_summary[total_order_value])Airflow 利用这种函数式调用设置任务依赖并管理数据传递。仅此三行代码Airflow 就知道如何调度和编排整条流水线。这里有一个容易误解的点在 DAG 定义阶段extract()的调用并不会真正执行函数体而是返回一个代表结果 XCom 的对象从 TaskFlow 概念文档 描述看即XComArg。你可以把XComArg作为下游任务或传统 Operator 的输入Airflow 会据此自动声明依赖方向——即compose_email在get_ip的下游。装饰器实现可参考 task-sdk 中的 decorator 基类其中包含XComArg的引入与 expand/mapping 相关逻辑。五、运行你的 DAG启用并触发 DAG 的标准步骤打开 Airflow UI在 DAG 列表中找到该 DAG点击开关将其启用点击 “Trigger Dag” 按钮手动触发或等待其按 schedule 自动运行。六、幕后机制与传统 Operator 写法的对比如果你用过 Airflow 1.xTaskFlow 用起来可能像“魔法”。教程中给出了同一 DAG 在传统写法下的形态——用PythonOperator 手动 XComimport json import pendulum from airflow.sdk import DAG from airflow.providers.standard.operators.python import PythonOperator def extract(): # Old way: simulate extracting data from a JSON string data_string {1001: 301.27, 1002: 433.21, 1003: 502.22} return json.loads(data_string) def transform(ti): # Old way: manually pull from XCom order_data_dict ti.xcom_pull(task_idsextract) total_order_value sum(order_data_dict.values()) return {total_order_value: total_order_value} def load(ti): # Old way: manually pull from XCom total ti.xcom_pull(task_idstransform)[total_order_value] print(fTotal order value is: {total:.2f}) with DAG( dag_idlegacy_etl_pipeline, scheduleNone, start_datependulum.datetime(2021, 1, 1, tzUTC), catchupFalse, tags[example], ) as dag: extract_task PythonOperator(task_idextract, python_callableextract) transform_task PythonOperator(task_idtransform, python_callabletransform) load_task PythonOperator(task_idload, python_callableload) extract_task transform_task load_task两种写法产生完全相同的结果但传统方式要求显式管理 XCom 和任务依赖ti.xcom_pull(task_ids...)与链。TaskFlow 写法中XCom 的存取与依赖图的构建全部自动化你可以专注于业务逻辑。XCom 是如何工作的TaskFlow 函数的返回值会被自动存为 XCom。这些值可以在 UI 的 “XCom” 标签页中检查。对于传统 Operator仍然可以手动调用xcom_pull()。七、错误处理与重试通过装饰器参数即可为任务配置重试task(retries3) def my_task(): ...这有助于确保瞬时故障不会直接导致任务失败。八、任务参数化与复用.override()装饰过的任务可以在多个 DAG 中复用并通过.override()覆盖task_id、retries等参数start add_task.override(task_idstart)(1, 2)你甚至可以从共享模块导入已装饰的任务函数。官方示例 example_python_decorator.py 展示了这一模式——循环生成 5 个睡眠任务每个任务用不同的task_idtask def my_sleeping_function(random_base): This is a function that will run within the DAG execution time.sleep(random_base) for i in range(5): sleeping_task my_sleeping_function.override(task_idfsleep_for_{i})(random_basei / 10) run_this log_the_sql sleeping_task该示例还展示了 Airflow 3.2 对异步 callable 的原生支持task直接装饰async def函数任务体内可以await。九、高级模式隔离执行环境当某些任务需要与 DAG 其余部分不同的 Python 依赖专用库或系统级包时TaskFlow 支持多种执行环境来隔离依赖。9.1 动态创建的虚拟环境task.virtualenv在任务运行时创建一个临时 virtualenv适合实验性或动态任务但可能有冷启动开销对应 example_python_decorator.py 第96-119行task.virtualenv( task_idvirtualenv_python, requirements[colorama0.4.0], system_site_packagesFalse ) def callable_virtualenv(): Example function that will be performed in a virtual environment. Importing at the module level ensures that it will not attempt to import the library before it is installed. from time import sleep from colorama import Back, Fore, Style print(Fore.RED some red text) print(Back.GREEN and with a green background) print(Style.DIM and in dim text) print(Style.RESET_ALL) for _ in range(4): print(Style.DIM Please wait..., flushTrue) sleep(1) print(Finished)注意示例中的细节库的 import 放在函数体内确保在安装该库之前不会尝试导入。9.2 外部 Python 环境task.external_python使用预装好的 Python 解释器执行任务适合环境一致或共享 virtualenv 的场景对应 example_python_decorator.py 第125-143行PATH_TO_PYTHON_BINARY sys.executable task.external_python(task_idexternal_python, pythonPATH_TO_PYTHON_BINARY) def callable_external_python(): import sys from time import sleep print(fRunning task via {sys.executable}) print(Sleeping) for _ in range(4): print(Please wait..., flushTrue) sleep(1) print(Finished)9.3 Docker 环境task.docker在 Docker 容器中运行任务适合把任务所需的一切打包进镜像前提是你的 worker 上可用 Docker。官方系统测试示例 example_taskflow_api_docker_virtualenv.py 中transform任务跑在 Docker 里而extract用 virtualenvtask.docker(imagepython:3.9-slim-bookworm, multiple_outputsTrue) def transform(order_data_dict: dict): #### Transform task A simple Transform task which takes in the collection of order data and computes the total order value. total_order_value 0 for value in order_data_dict.values(): total_order_value value return {total_order_value: total_order_value}注意Docker 装饰器要求 Airflow 2.2 且安装 Docker provider。该示例同时演示了task.virtualenv的serializerdill参数用于序列化函数体。9.4 KubernetesPodOperatortask.kubernetes在 Kubernetes Pod 中运行任务与主 Airflow 环境完全隔离适合大任务或需要自定义运行时的任务对应 example_kubernetes_decorator.pytask.kubernetes( imagepython:3.9-slim-buster, namek8s_test, namespacedefault, in_clusterFalse, config_file/path/to/.kube/config, ) def execute_in_k8s_pod(): import time print(Hello from k8s pod) time.sleep(2)注意Kubernetes 装饰器要求 Airflow 2.4 且安装 Kubernetes provider。十、高级模式传感器task.sensor允许用 Python 函数构建轻量、可复用的传感器同时支持 poke 与 reschedule 两种模式。官方示例 example_sensor_decorator.py 演示了 reschedule 模式import pendulum from airflow.sdk import PokeReturnValue, dag, task dag( scheduleNone, start_datependulum.datetime(2021, 1, 1, tzUTC), catchupFalse, tags[example], ) def example_sensor_decorator(): # Using a sensor operator to wait for the upstream data to be ready. task.sensor(poke_interval60, timeout3600, modereschedule) def wait_for_upstream() - PokeReturnValue: return PokeReturnValue(is_doneTrue, xcom_valuexcom_value) task def dummy_operator() - None: pass wait_for_upstream() dummy_operator() tutorial_etl_dag example_sensor_decorator()要点poke_interval60表示每 60 秒检查一次timeout3600表示最长等待 1 小时modereschedule表示检查不通过时释放 worker 资源重新调度。传感器函数返回PokeReturnValue对象其中is_done指示是否完成xcom_value会作为该任务的 XCom 输出传给下游。十一、与传统任务混用装饰任务可以与经典 Operator 自由组合这在对接社区 provider 或渐进式迁移到 TaskFlow 时特别有用。两种衔接方式用把 TaskFlow 任务与传统任务链接起来通过.output属性把 TaskFlow 任务的返回值传给传统 Operator 的参数。例如 TaskFlow 概念文档 中的示例get_ip()与compose_email()是 TaskFlow 任务EmailOperator是传统 Operator但它直接消费email_info[subject]/email_info[body]来自compose_email返回字典的 XComArg 索引Airflow 会自动判定其下游依赖from airflow.sdk import task from airflow.providers.smtp.operators.smtp import EmailOperator task def get_ip(): return my_ip_service.get_main_ip() task(multiple_outputsTrue) def compose_email(external_ip): return { subject: fServer connected from {external_ip}, body: fYour server executing Airflow is connected from the external IP {external_ip}br } email_info compose_email(get_ip()) EmailOperator( task_idsend_email_notification, toexampleexample.com, subjectemail_info[subject], html_contentemail_info[body], )十二、TaskFlow 中的模板化与上下文变量与任务装饰器一样TaskFlow 函数的参数自动支持模板化——包括从文件加载内容或使用运行时参数。12.1 显式接收上下文变量执行 callable 时Airflow 会传入一组关键字参数与 Jinja 模板中可用的上下文完全一致。要接收某个上下文变量只需把它作为函数签名的关键字参数task def my_python_callable(*, ti, next_ds): pass上面的 callable 将收到ti与next_ds两个上下文变量的值。12.2 用**kwargs接收完整上下文也可以选择接收整个上下文task def my_python_callable(**kwargs): ti kwargs[ti] next_ds kwargs[next_ds]但需注意这会带来轻微的性能损耗——Airflow 需要展开整个上下文而其中可能包含大量你用不到的内容。因此官方推荐优先使用显式参数。example_python_decorator.py 中print_context任务即展示了这一用法task(task_idprint_the_context) def print_context(dsNone, **kwargs): Print the Airflow context and ds variable from the context. pprint(kwargs) print(ds) return Whatever you return gets printed in the logs12.3 深层调用中获取上下文get_current_context有时你想在调用栈深处访问上下文又不想把上下文变量从任务 callable 一路传下去。此时可以用get_current_context方法from airflow.sdk import get_current_context def some_function_in_your_library(): context get_current_context() ti context[ti]12.4 模板化文件templates_exts与templates_dict传入装饰函数的参数会自动模板化你也可以用templates_exts模板化文件扩展名task(templates_exts[.sql]) def read_sql(sql): ...官方示例中还演示了用templates_dict从文件加载并渲染 SQLtask(task_idlog_sql_query, templates_dict{query: sql/sample.sql}, templates_exts[.sql]) def log_sql(**kwargs): log.info(Python task decorator query: %s, str(kwargs[templates_dict][query]))十三、条件执行用task.run_if()或task.skip_if()根据运行时的动态条件控制任务是否执行无需修改 DAG 结构task.run_if(lambda ctx: ctx[task_instance].task_id run) task.bash() def echo(): return echo run十四、关于可传递对象类型由于 TaskFlow 依赖 XCom 在任务间传递变量用作参数的变量必须可序列化。Airflow 开箱支持所有内建类型int、str 等也支持dataclass或attr.define装饰的对象。若需自定义序列化可为类添加serialize()方法与静态方法deserialize(data, version)并用__version__: ClassVar[int]做对象版本管理详见 TaskFlow 概念文档 的 Passing Arbitrary Objects As Arguments 一节。一个实用特性若用Assetattr.define装饰作为输入参数它会自动注册为 inlet若任务返回值是Asset或list[Asset]会自动注册为 outlet——这让 TaskFlow DAG 直接获得资产感知调度能力。十五、后续探索方向完成第一条 TaskFlow 流水线后推荐沿以下方向深入与 教程原文 的 “What to Explore Next” 一致给 DAG 添加新任务比如 filter 或 validation 步骤修改返回值并传递多个输出multiple_outputs探索重试与.override(task_id...)覆盖打开 Airflow UI检查数据如何在任务间流动包括任务日志与依赖关系学习资产感知工作流/authoring-and-scheduling/asset-scheduling与调度选项继续阅读 TaskFlow 核心概念文档 获取上下文变量、日志、任意对象参数与自定义对象版本化的完整说明进入下一篇教程 pipeline 学习 TaskGroup 等结构化模式。附本文引用的仓库文件文件作用airflow-core/docs/tutorial/taskflow.rst本文对应的官方教程原文airflow-core/src/airflow/example_dags/tutorial_taskflow_api.pyTaskFlow ETL 示例 DAG教程主体代码providers/standard/src/airflow/providers/standard/example_dags/example_python_decorator.pyPython 装饰器示例虚拟环境、外部解释器、override、异步任务providers/standard/src/airflow/providers/standard/example_dags/example_sensor_decorator.pytask.sensor传感器示例providers/docker/tests/system/docker/example_taskflow_api_docker_virtualenv.pytask.docker/task.virtualenv组合示例providers/cncf/kubernetes/tests/system/cncf/kubernetes/example_kubernetes_decorator.pytask.kubernetes示例airflow-core/docs/core-concepts/taskflow.rstTaskFlow 核心概念XComArg、上下文、对象序列化task-sdk/src/airflow/sdk/bases/decorator.py装饰器底层实现XComArg、任务展开与校验逻辑版本适用说明本仓库示例代码中的dag、task、get_current_context等导入自airflow.sdkAirflow 3.x 的新 SDK 入口Airflow 2.x 环境下等价导入为from airflow.decorators import task, dag。Docker 装饰器需 Airflow 2.2 与 Docker providerKubernetes 装饰器需 Airflow 2.4 与 Kubernetes provider。【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表