如何用Apache Airflow TaskFlow API构建无返回值任务流?
用TaskFlow API实现无返回值、无数据传递的任务流
TaskFlow API完全支持这种仅依赖执行顺序、无返回值的任务流场景,写法比传统方式更简洁——核心是用@task装饰普通函数,再通过>>运算符定义任务依赖即可,和你熟悉的旧方式逻辑完全一致。
完整示例代码
from airflow import DAG from airflow.decorators import task from datetime import datetime @dag( schedule_interval="@daily", start_date=datetime(2024, 1, 1), catchup=False, tags=["example", "taskflow"] ) def no_return_value_dag(): @task def extract(): # 这里写抽取逻辑:比如读取数据源、写入临时存储等 print("执行数据抽取操作") # 不需要返回任何值 @task def transform(): # 这里写转换逻辑:比如清洗数据、处理格式等 print("执行数据转换操作") # 不需要返回任何值 @task def load(): # 这里写加载逻辑:比如写入数据库、数据仓库等 print("执行数据加载操作") # 不需要返回任何值 # 定义任务依赖顺序,和旧方式逻辑一致 extract() >> transform() >> load() # 实例化DAG dag = no_return_value_dag()
关键说明
- 每个任务用
@task装饰后,Airflow会自动将其封装成Task实例,即使函数没有返回值也不影响任务调度 - 依赖关系通过
>>运算符串联,逻辑和传统的extract >> transform >> load完全相同 - 任务内部可以执行任意业务逻辑:读写文件、调用外部API、执行SQL语句等,不需要为了传递数据而强制返回值
内容的提问来源于stack exchange,提问作者TeoK
相关产品推荐
相关产品推荐

