在第一个教程中,我们使用 BashOperator 等传统 Operator 构建了 Airflow DAG。接下来介绍一种更现代、更符合 Python 风格的工作流写法:Airflow 2.0 引入的 TaskFlow API。
TaskFlow API 旨在让代码更简单、清晰且易于维护。只需编写普通 Python 函数并添加装饰器,Airflow 就会负责创建任务、建立依赖,以及在任务之间传递数据。本教程将用它构建一个简单的 ETL 流水线:提取、转换、加载。
先看完整的 TaskFlow 流水线
下面是完整示例,源码文件为 airflow/example_dags/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.
Documentation that goes along with the Airflow TaskFlow API tutorial is
located
[here](https://airflow.apache.org/docs/apache-airflow/stable/tutorial_taskflow_api.html)
"""
@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()
第一步:定义 DAG
与以前一样,DAG 是由 Airflow 加载和解析的 Python 脚本。这次使用 @dag 装饰器进行定义:
@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.
Documentation that goes along with the Airflow TaskFlow API tutorial is
located
[here](https://airflow.apache.org/docs/apache-airflow/stable/tutorial_taskflow_api.html)
"""
为了让 Airflow 发现这个 DAG,需要调用被 @dag 装饰的 Python 函数:
tutorial_taskflow_api()
从 Airflow 2.4 起,如果使用 @dag 装饰器,或在 with 代码块中定义 DAG,就不再需要把它赋给全局变量,Airflow 会自动发现它。
DAG 加载后,可以在 Airflow 界面的 Graph View 中查看任务之间的连接关系。
第二步:用 @task 编写任务
在 TaskFlow 中,每个任务都是普通 Python 函数。加上 @task 装饰器,就能把它变为 Airflow 可调度、可执行的任务。下面是 extract 任务:
@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
函数返回值会传给下游任务,无需手动管理 XCom。TaskFlow 底层仍使用 XCom 传递数据,只是将之前需要显式处理的细节封装起来。transform 和 load 任务也使用同样的定义方式。
上面的 @task(multiple_outputs=True) 告诉 Airflow,函数返回的字典应拆分成多个独立 XCom。字典中的每个键都会成为单独的 XCom 条目,便于下游任务引用具体值。如果不启用 multiple_outputs=True,整个字典会作为单个 XCom 保存,需要整体读取。
第三步:建立流程
定义好任务后,只需像调用 Python 函数一样调用它们,即可构建流水线。Airflow 会根据这些调用关系建立任务依赖并管理数据传递:
order_data = extract()
order_summary = transform(order_data)
load(order_summary["total_order_value"])
仅凭这段代码,Airflow 就知道如何调度和编排整个流水线。
运行 DAG
启用并触发 DAG 的步骤如下:
- 打开 Airflow 界面。
- 在列表中找到该 DAG,点击开关启用。
- 点击“Trigger Dag”手动触发,或等待它按计划运行。
底层发生了什么
如果使用过 Airflow 1.x,可以通过对比传统方式理解 TaskFlow 的作用。
传统方式:手动连接任务与操作 XCom
TaskFlow API 出现以前,需要使用 PythonOperator 等 Operator,并通过 XCom 手动在任务之间传递数据。同一个 DAG 用传统方式编写,大致如下:
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
这个版本与 TaskFlow 示例得到相同结果,但需要显式管理 XCom 和任务依赖。
TaskFlow 方式
使用 TaskFlow 时,这些细节由框架自动处理:
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.
Documentation that goes along with the Airflow TaskFlow API tutorial is
located
[here](https://airflow.apache.org/docs/apache-airflow/stable/tutorial_taskflow_api.html)
"""
@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 仍然使用 XCom,也仍然构建依赖图,只是将相关操作封装起来,让代码更集中于业务逻辑。
XCom 如何工作
TaskFlow 的返回值会自动保存为 XCom,可在界面的“XCom”选项卡中查看。传统 Operator 仍然可以手动调用 xcom_pull() 读取数据。
错误处理与重试
可以直接在任务装饰器中配置重试,例如设置最大重试次数:
@task(retries=3)
def my_task(): ...
这有助于避免短暂故障直接导致任务最终失败。
任务参数化
被装饰的任务可以在多个 DAG 中复用,并覆盖 task_id、retries 等参数:
start = add_task.override(task_id="start")(1, 2)
也可以把这些任务放到共享模块,再从不同 DAG 中导入。
接下来可以尝试什么
完成第一个 TaskFlow 流水线后,可以进一步练习:
- 为 DAG 增加过滤或校验任务。
- 修改返回值,并传递多个输出。
- 通过
.override(task_id="...")探索参数覆盖与重试配置。 - 在 Airflow 界面中查看任务之间的数据流、任务日志及依赖关系。
可以继续阅读构建简单数据流水线、TaskFlow API与核心概念,也可以继续阅读下方的进阶模式。
TaskFlow 进阶模式
掌握基础之后,可以使用以下方式处理更复杂的工作流。
复用被装饰的任务
任务可以跨多个 DAG 或 DAG 运行复用,特别适合公用工具逻辑和共享业务规则。使用 .override() 定制 task_id、retries 等元数据:
start = add_task.override(task_id="start")(1, 2)
被装饰的任务同样可以从共享模块导入。
处理互相冲突的依赖
某些任务需要与 DAG 其他部分不同的 Python 依赖,例如专用库或系统级软件包。TaskFlow 支持多种执行环境,用来隔离这些依赖。
动态创建虚拟环境
任务运行时创建临时 virtualenv,适合实验性或动态任务,但会带来冷启动开销。示例来自标准 provider 的 example_python_decorator.py:
@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")
virtualenv_task = callable_virtualenv()
外部 Python 环境
使用预先安装的 Python 解释器执行任务,适合需要稳定环境或共享虚拟环境的情况。示例同样来自 example_python_decorator.py:
@task.external_python(task_id="external_python", python=PATH_TO_PYTHON_BINARY)
def callable_external_python():
"""
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.
"""
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")
external_python_task = callable_external_python()
Docker 环境
在 Docker 容器内执行任务,适合将任务所需的一切一起打包。worker 必须能够使用 Docker。下面的示例来自 Docker provider 的 example_taskflow_api_docker_virtualenv.py:
@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}
注意:这需要 Airflow 2.2 及 Docker provider。
KubernetesPodOperator
在 Kubernetes Pod 内执行任务,使其与 Airflow 主环境完全隔离。这种方式适合大型任务或需要自定义运行时的任务。示例来自 Kubernetes provider 的 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)
@task.kubernetes(image="python:3.9-slim-buster", namespace="default", in_cluster=False)
def print_pattern():
n = 5
for i in range(n):
# inner loop to handle number of columns
# values changing acc. to outer loop
for _ in range(i + 1):
# printing stars
print("* ", end="")
# ending line after each row
print("\r")
execute_in_k8s_pod_instance = execute_in_k8s_pod()
print_pattern_instance = print_pattern()
execute_in_k8s_pod_instance >> print_pattern_instance
注意:这需要 Airflow 2.4 及 Kubernetes provider。
使用 Sensor
使用 @task.sensor 可以把 Python 函数变成轻量、可复用的 Sensor。它同时支持 poke 与 reschedule 两种模式。示例来自标准 provider 的 example_sensor_decorator.py:
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()
与传统任务混合使用
被装饰的任务可以与传统 Operator 组合。这在使用社区 provider,或逐步迁移到 TaskFlow 时尤其有用。
可以使用 >> 串联 TaskFlow 任务与传统任务,也可以通过 .output 属性传递数据。
TaskFlow 模板处理
与传统任务一样,TaskFlow 装饰的函数支持模板参数,包括从文件加载内容和使用运行时参数。
执行函数时,Airflow 会提供一组关键字参数,它们与 Jinja 模板中可用的上下文变量完全对应。要接收其中某些值,可将相应上下文键写为函数的关键字参数。例如:
@task
def my_python_callable(*, ti, next_ds):
pass
上面的函数会获得 ti 和 next_ds 两个上下文变量。
也可以用 **kwargs 接收整个上下文。但这会带来少量性能开销,因为 Airflow 需要展开全部上下文,其中可能包含大量用不到的数据。因此更推荐像前例那样明确声明所需参数。
@task
def my_python_callable(**kwargs):
ti = kwargs["ti"]
next_ds = kwargs["next_ds"]
如果需要在调用栈较深的位置访问上下文,却不希望从任务函数逐层传递变量,可以调用 get_current_context:
from airflow.sdk import get_current_context
def some_function_in_your_library():
context = get_current_context()
ti = context["ti"]
传入被装饰函数的参数会自动进行模板处理。还可以通过 templates_exts 对文件内容应用模板:
@task(templates_exts=[".sql"])
def read_sql(sql): ...
条件执行
使用 @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 API 构建清晰、可维护的 DAG。可以继续学习资产感知调度、调度选项,或进入下一篇简单数据流水线教程。
原文来源:Pythonic Dags with the TaskFlow API。本文依据留存原文译为中文,代码示例按原文保留。
© The Apache Software Foundation. Apache Airflow、Apache、Airflow 及相关标志是 Apache Software Foundation 的注册商标或商标,其他产品及品牌归各自所有者所有。原文许可:License;Apache 许可证。











暂无评论内容