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( schedule=None, start_date=pendulum.datetime(2021, 1, 1, tz="UTC"), catchup=False, 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_outputs=True) 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(f"Total 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装饰器定义 DAG
DAG 本质上仍是 Airflow 加载并解析的 Python 脚本,但这里使用@dag装饰器来定义。官方示例中 DAG 的定义部分(对应源码 tutorial_taskflow_api.py 第32-37行):
@dag( schedule=None, start_date=pendulum.datetime(2021, 1, 1, tz="UTC"), catchup=False, 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_outputs=True的行为:transform使用了@task(multiple_outputs=True),这告诉 Airflow 函数返回的是一个字典,应将其拆分为独立的 XCom。字典中的每个 key 各自成为一个 XCom 条目,下游任务可以直接引用特定值(如order_summary["total_order_value"])。如果省略multiple_outputs=True,整个字典会作为单个 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.x,TaskFlow 用起来可能像“魔法”。教程中给出了同一 DAG 在传统写法下的形态——用PythonOperator+ 手动 XCom:
import 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_ids="extract") 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_ids="transform")["total_order_value"] print(f"Total order value is: {total:.2f}") with DAG( dag_id="legacy_etl_pipeline", schedule=None, start_date=pendulum.datetime(2021, 1, 1, tz="UTC"), catchup=False, tags=["example"], ) as dag: extract_task = PythonOperator(task_id="extract", python_callable=extract) transform_task = PythonOperator(task_id="transform", python_callable=transform) load_task = PythonOperator(task_id="load", python_callable=load) extract_task >> transform_task >> load_task两种写法产生完全相同的结果,但传统方式要求显式管理 XCom 和任务依赖(ti.xcom_pull(task_ids=...)与>>链)。TaskFlow 写法中,XCom 的存取与依赖图的构建全部自动化,你可以专注于业务逻辑。
XCom 是如何工作的
TaskFlow 函数的返回值会被自动存为 XCom。这些值可以在 UI 的 “XCom” 标签页中检查。对于传统 Operator,仍然可以手动调用xcom_pull()。
七、错误处理与重试
通过装饰器参数即可为任务配置重试:
@task(retries=3) def my_task(): ...这有助于确保瞬时故障不会直接导致任务失败。
八、任务参数化与复用(.override())
装饰过的任务可以在多个 DAG 中复用,并通过.override()覆盖task_id、retries等参数:
start = add_task.override(task_id="start")(1, 2)你甚至可以从共享模块导入已装饰的任务函数。官方示例 example_python_decorator.py 展示了这一模式——循环生成 5 个睡眠任务,每个任务用不同的task_id:
@task 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_id=f"sleep_for_{i}")(random_base=i / 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_id="virtualenv_python", requirements=["colorama==0.4.0"], system_site_packages=False ) 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...", flush=True) 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_id="external_python", python=PATH_TO_PYTHON_BINARY) def callable_external_python(): import sys from time import sleep print(f"Running task via {sys.executable}") print("Sleeping") for _ in range(4): print("Please wait...", flush=True) sleep(1) print("Finished")9.3 Docker 环境
@task.docker在 Docker 容器中运行任务,适合把任务所需的一切打包进镜像;前提是你的 worker 上可用 Docker。官方系统测试示例 example_taskflow_api_docker_virtualenv.py 中,transform任务跑在 Docker 里,而extract用 virtualenv:
@task.docker(image="python:3.9-slim-bookworm", multiple_outputs=True) 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的serializer="dill"参数,用于序列化函数体。
9.4 KubernetesPodOperator
@task.kubernetes在 Kubernetes Pod 中运行任务,与主 Airflow 环境完全隔离,适合大任务或需要自定义运行时的任务(对应 example_kubernetes_decorator.py):
@task.kubernetes( image="python:3.9-slim-buster", name="k8s_test", namespace="default", in_cluster=False, 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( schedule=None, start_date=pendulum.datetime(2021, 1, 1, tz="UTC"), catchup=False, tags=["example"], ) def example_sensor_decorator(): # Using a sensor operator to wait for the upstream data to be ready. @task.sensor(poke_interval=60, timeout=3600, mode="reschedule") def wait_for_upstream() -> PokeReturnValue: return PokeReturnValue(is_done=True, xcom_value="xcom_value") @task def dummy_operator() -> None: pass wait_for_upstream() >> dummy_operator() tutorial_etl_dag = example_sensor_decorator()要点:poke_interval=60表示每 60 秒检查一次,timeout=3600表示最长等待 1 小时,mode="reschedule"表示检查不通过时释放 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_outputs=True) def compose_email(external_ip): return { 'subject': f'Server connected from {external_ip}', 'body': f'Your server executing Airflow is connected from the external IP {external_ip}<br>' } email_info = compose_email(get_ip()) EmailOperator( task_id='send_email_notification', to='example@example.com', subject=email_info['subject'], html_content=email_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_id="print_the_context") def print_context(ds=None, **kwargs): """Print the Airflow context and ds variable from the context.""" pprint(kwargs) print(ds) return "Whatever you return gets printed in the logs"12.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从文件加载并渲染 SQL:
@task(task_id="log_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" 一节)。一个实用特性:若用Asset(@attr.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.py | TaskFlow ETL 示例 DAG(教程主体代码) |
| providers/standard/src/airflow/providers/standard/example_dags/example_python_decorator.py | Python 装饰器示例:虚拟环境、外部解释器、override、异步任务 |
| providers/standard/src/airflow/providers/standard/example_dags/example_sensor_decorator.py | @task.sensor传感器示例 |
| providers/docker/tests/system/docker/example_taskflow_api_docker_virtualenv.py | @task.docker/@task.virtualenv组合示例 |
| providers/cncf/kubernetes/tests/system/cncf/kubernetes/example_kubernetes_decorator.py | @task.kubernetes示例 |
| airflow-core/docs/core-concepts/taskflow.rst | TaskFlow 核心概念:XComArg、上下文、对象序列化 |
| task-sdk/src/airflow/sdk/bases/decorator.py | 装饰器底层实现(XComArg、任务展开与校验逻辑) |
版本适用说明:本仓库示例代码中的
dag、task、get_current_context等导入自airflow.sdk(Airflow 3.x 的新 SDK 入口);Airflow 2.x 环境下等价导入为from airflow.decorators import task, dag。Docker 装饰器需 Airflow 2.2+ 与 Docker provider,Kubernetes 装饰器需 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),仅供参考