Apache Airflow 101:使用 Airflow SDK 编写你的第一个 DAG 工作流(附完整源码解析与测试指南)
【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow
本教程以 Apache Airflow 官方入门文档 fundamentals.rst 为主线,结合仓库内完整示例 tutorial.py,系统讲解 DAG 的概念、DAG 定义文件的编写方式、Operator 与 Task 的关系、Jinja 模板渲染、任务依赖编排以及命令行测试方法。读完本文,你将具备独立编写、校验、本地测试并提交一个可被 Scheduler 调度执行的 Airflow 工作流的能力。
什么是 DAG?
DAG(Directed Acyclic Graph,有向无环图)是 Airflow 中工作流的核心抽象。简单来说,DAG 是一组任务的集合,这些任务按照它们之间的关系与依赖被组织起来——它就像一张工作流的路线图,清楚地展示了每个任务如何与其他任务连接。
"有向"意味着任务之间的依赖有方向(谁先谁后);"无环"意味着依赖关系不能形成闭环,否则调度器无法确定执行起点。Airflow 会在解析 DAG 时检测环路,一旦发现循环依赖或同一依赖被重复引用,就会抛出错误(见下文"设置任务依赖"一节)。
一个完整的 Pipeline 定义示例
仓库中的 tutorial.py 是官方入门示例,虽然初看有些内容,但每一行都值得拆解。我们先给出完整代码,随后逐段解释:
# [START tutorial] import textwrap from datetime import datetime, timedelta # Operators; we need this to operate! from airflow.providers.standard.operators.bash import BashOperator # The DAG object; we'll need this to instantiate a DAG from airflow.sdk import DAG with DAG( "tutorial", default_args={ "depends_on_past": False, "retries": 1, "retry_delay": timedelta(minutes=5), }, description="A simple tutorial DAG", schedule=timedelta(days=1), start_date=datetime(2021, 1, 1), catchup=False, tags=["example"], ) as dag: t1 = BashOperator( task_id="print_date", bash_command="date", ) t2 = BashOperator( task_id="sleep", depends_on_past=False, bash_command="sleep 5", retries=3, ) t1.doc_md = textwrap.dedent( """\ #### Print the current date This task runs `date` by using the `bash_command` argument on `BashOperator`. ... """ ) dag.doc_md = __doc__ templated_command = textwrap.dedent( """ {% for i in range(5) %} echo "{{ ds }}" echo "{{ macros.ds_add(ds, 7)}}" {% endfor %} """ ) t3 = BashOperator( task_id="templated", depends_on_past=False, bash_command=templated_command, ) t1 >> [t2, t3] # [END tutorial]理解 DAG 定义文件
把 Airflow 的 Python 脚本想象成一个用代码描述 DAG 结构的配置文件——这一点非常重要:你在其中定义的 task 实际运行在另一个环境(Scheduler 派发、Worker 执行)中,因此这个脚本本身不是用来做数据处理的。
它的主要职责是定义DAG对象,并且必须能够被快速求值。原因在于:Airflow 的 Dag File Processor(DAG 文件处理器)会定期检查 DAG 文件夹中的每个文件,一旦发现变更就重新解析。如果脚本在导入阶段执行了重量级操作(例如连接数据库、发起网络请求),会拖慢整个解析过程,甚至导致 DAG 无法被识别。因此最佳实践是:DAG 文件中只做"定义"工作,把真正的计算留给任务执行阶段。
导入模块
与其他 Python 脚本一样,第一步是导入所需库。tutorial 示例导入了三部分内容:
import textwrap from datetime import datetime, timedelta # Operators; we need this to operate! from airflow.providers.standard.operators.bash import BashOperator # The DAG object; we'll need this to instantiate a DAG from airflow.sdk import DAGtextwrap:标准库,用于textwrap.dedent去除多行字符串的公共缩进,让文档字符串与模板字符串书写更整洁;datetime/timedelta:标准库,用于设置start_date、retry_delay等时间参数;BashOperator:来自airflow.providers.standard.operators.bash,是本次使用的 Operator;DAG:来自airflow.sdk,是本仓库中 DAG 对象的官方入口(注意新版 Airflow 已将其收敛到airflow.sdk命名空间)。
关于 Python 与 Airflow 的模块管理机制(例如哪些目录会被自动扫描、模块名如何解析),可参考仓库文档 modules_management.rst。
设置默认参数(default_args)
创建 DAG 及其任务时,你可以把参数直接传给每个 task,也可以把一组公共参数放在字典中统一定义。后者通常更高效、更整洁,因为所有 task 会默认继承这些参数,同时仍允许在单个 task 上覆盖。
tutorial 示例中的default_args:
default_args={ "depends_on_past": False, "retries": 1, "retry_delay": timedelta(minutes=5), # 'queue': 'bash_queue', # 'pool': 'backfill', # 'priority_weight': 10, # 'end_date': datetime(2016, 1, 1), # 'wait_for_downstream': False, # 'execution_timeout': timedelta(seconds=300), # 'on_failure_callback': some_function, # or list of functions # 'on_success_callback': some_other_function, # or list of functions # 'on_retry_callback': another_function, # or list of functions # 'sla_miss_callback': yet_another_function, # or list of functions # 'on_skipped_callback': another_function, #or list of functions # 'trigger_rule': 'all_success' },被注释掉的参数展示了default_args的常见能力边界,它们都可以按需启用:
| 参数 | 含义 | 示例值 |
|---|---|---|
depends_on_past | 当前任务实例是否依赖上一个调度周期的任务实例成功 | False |
retries | 失败后自动重试的次数 | 1 |
retry_delay | 两次重试之间的等待时长 | timedelta(minutes=5) |
queue | 任务被派发到的队列 | 'bash_queue' |
pool | 任务使用的资源池 | 'backfill' |
priority_weight | 任务优先级权重 | 10 |
end_date | 任务不再被调度的时间点 | datetime(2016, 1, 1) |
wait_for_downstream | 是否等待下游任务的上一个实例成功 | False |
execution_timeout | 任务实例运行的最长时限,超时则失败 | timedelta(seconds=300) |
on_failure_callback/on_success_callback/on_retry_callback/on_skipped_callback | 失败/成功/重试/跳过时的回调函数(或函数列表) | some_function |
sla_miss_callback | SLA 未达成的回调 | yet_another_function |
trigger_rule | 任务触发条件规则 | 'all_success' |
如果想深入了解BaseOperator的全部参数,可查阅 airflow.sdk.BaseOperator 文档 及仓库中 BaseOperator 源码 对应的BaseOperator类定义。
创建 DAG
接下来实例化一个DAG对象来承载任务。tutorial 示例:
with DAG( "tutorial", default_args={...}, description="A simple tutorial DAG", schedule=timedelta(days=1), start_date=datetime(2021, 1, 1), catchup=False, tags=["example"], ) as dag:各参数含义:
"tutorial"(位置参数):dag_id,即 DAG 的唯一标识符。它在你的 Airflow 实例中必须全局唯一,UI、CLI、API 都通过它引用该 DAG;default_args:上一步定义的默认参数字典,会传递给每个任务;description:DAG 的简要描述,展示在 UI 的 DAG 列表中;schedule=timedelta(days=1):调度周期,这里表示每天运行一次。除了timedelta,还可以使用 cron 表达式字符串(如"0 0 * * *")或@daily等预设值;start_date:DAG 开始生效的时间。注意 Airflow 的调度是"结束时间对齐"的:第一个 DAG Run 的 logical date(逻辑日期)通常等于start_date,但真正触发时间在其之后的一个调度周期;catchup=False:关闭回填。若为True,当 DAG 从start_date到当前时间之间有大量错过的调度周期时,Scheduler 会一次性补跑所有错过的 DAG Run;设为False则只调度最新周期,避免刚上线时批量触发;tags=["example"]:为 DAG 打标签,便于在 UI 中按标签筛选;with ... as dag::上下文管理器语法,块内实例化的任务会自动注册到该 DAG 上,这是新版 Airflow 推荐、也是本仓库示例采用的写法。
理解 Operator
Operator 是 Airflow 中的"工作单元",是构建工作流的积木,决定了任务将要执行什么动作。所有 Operator 都继承自BaseOperator,因此共享运行任务所需的核心参数(如retries、depends_on_past、execution_timeout等)。
社区与官方提供大量 Operator,常见的有:
PythonOperator:执行 Python 可调用对象;BashOperator:执行 Bash 命令或脚本(本教程主角);KubernetesPodOperator:在 Kubernetes 集群中拉起 Pod 执行任务;- 以及各类 Provider 提供的专有 Operator(可浏览仓库 providers 目录)。
除了直接实例化 Operator,Airflow 还提供更"Pythonic"的 TaskFlow API,可以用装饰器把普通 Python 函数变成任务,本文暂不展开。
定义任务(Task)
要使用 Operator,必须先把它实例化为任务(Task)。任务决定了 Operator 在 DAG 上下文中如何执行工作。task_id是每个任务的唯一标识符。
tutorial 示例实例化了两次BashOperator:
t1 = BashOperator( task_id="print_date", bash_command="date", ) t2 = BashOperator( task_id="sleep", depends_on_past=False, bash_command="sleep 5", retries=3, )注意这里把Operator 专有参数(bash_command)与BaseOperator继承来的通用参数(retries、depends_on_past)混合使用,这让代码更简洁。t2还把retries覆盖为3,演示了按任务覆盖默认值的能力。
任务参数的优先级如下:
- 显式传入的参数(如
t2的retries=3); default_args字典中的值(如retries=1);- Operator 自身的默认值(如果存在)。
注意:每个任务必须包含或继承
task_id和owner两个参数,否则 Airflow 会报错。幸运的是,全新安装的 Airflow 默认将owner设为airflow,因此你通常只需确保设置task_id即可。
使用 Jinja 模板渲染
Airflow 内置了Jinja 模板引擎,让你能访问内置变量(如{{ ds }})与宏(macros)来动态生成命令内容。{{ ds }}是最常用的模板变量,代表逻辑日期(logical date)的日期戳,格式为YYYY-MM-DD。
tutorial 示例中的模板任务:
templated_command = textwrap.dedent( """ {% for i in range(5) %} echo "{{ ds }}" echo "{{ macros.ds_add(ds, 7)}}" {% endfor %} """ ) t3 = BashOperator( task_id="templated", depends_on_past=False, bash_command=templated_command, )这段模板包含:
{% for i in range(5) %}...{% endfor %}:Jinja 控制流语句块,循环 5 次;{{ ds }}:逻辑日期的日期戳变量,例如2021-01-01;{{ macros.ds_add(ds, 7) }}:宏调用,ds_add会在给定日期上加上指定天数,这里得到2021-01-08。
实际渲染后,templated任务执行的命令大致是:
echo "2021-01-01" echo "2021-01-08" echo "2021-01-01" echo "2021-01-08" ...(共 5 组)关于模板,还有几点实用技巧:
- 传入脚本文件:
bash_command可以直接传文件名,例如bash_command='templated_command.sh',把命令逻辑拆到独立文件中,便于组织与维护; - 自定义宏与过滤器:可以在 DAG 上定义
user_defined_macros和user_defined_filters,创建自己的模板变量与过滤器; - 完整变量/宏清单:所有可在模板中引用的变量与宏,见仓库 templates-ref.rst。
为 DAG 与任务添加文档
Airflow 允许为 DAG 或单个任务附加文档,直接在 UI 中渲染查看:
- DAG 文档:以Markdown渲染在 DAG 详情页;
- 任务文档:支持纯文本、Markdown、reStructuredText、JSON、YAML 等多种格式。
当任务文档使用doc_md时,Airflow 渲染常见的 Markdown 特性,包括:行内代码、围栏代码块、以math围栏包裹的公式(用KaTeX渲染)以及mermaid围栏中的Mermaid 流程图。
在 tutorial DAG 中,print_date任务(t1)通过doc_md展示了这些能力:
t1.doc_md = textwrap.dedent( """\ #### Print the current date This task runs `date` by using the `bash_command` argument on `BashOperator`. In the Task Instance Details page, Airflow renders this documentation from the task's `doc_md` field. After this task succeeds, Airflow can run both downstream tasks: `sleep` and `templated`. ```bash dateMath fences are rendered with KaTeX. This tutorial starts one task and then branches into two downstream tasks:
1\ \text{upstream task} + 2\ \text{downstream tasks} = 3\ \text{tasks}The same dependency is shown as a Mermaid diagram:
""" )
与此同时,DAG 级文档可以用模块 docstring 或直接赋值: ```python dag.doc_md = __doc__ # 使用文件开头的 docstring # 或者直接写字符串 dag.doc_md = """ This is a documentation placed anywhere """下图展示了任务doc_md在 UI 的 Task Instance Details 页面中的渲染效果(Markdown、KaTeX 公式与 Mermaid 图并存):
实践建议:把文档紧挨着它所描述的任务编写(如t1.doc_md = ...),保持文档与代码同步演进。
设置任务依赖
Airflow 中任务之间可以互相依赖。假设有任务t1、t2、t3,可以用多种方式表达依赖:
t1.set_downstream(t2) # 这表示 t2 需要等 t1 成功运行后才能运行 # 等价于: t2.set_upstream(t1) # 也可以使用位移运算符(bit shift)链式表达: t1 >> t2 # 反向的上游依赖: t2 << t1 # 链式多依赖,位移运算符更简洁: t1 >> t2 >> t3 # 列表形式设置依赖,以下写法效果相同: t1.set_downstream([t2, t3]) t1 >> [t2, t3] [t2, t3] << t1tutorial 示例使用的正是列表形式的位移运算符写法:
t1 >> [t2, t3]它表达"t1成功之后,t2与t3并行运行"的扇形结构。
需要警惕的是,Airflow 会检测 DAG 中的环(cycle),也会检测同一依赖被重复引用的情况,一旦发现就会抛出错误。因此在编写复杂 DAG 时,应避免出现t1 >> t2 >> t1之类的循环依赖。
处理时区
创建一个时区感知(time zone aware)的 DAG很简单:使用 pendulum 库提供的时区感知日期时间即可,例如pendulum.datetime(2021, 1, 1, tz="Asia/Shanghai")。
务必避免使用标准库datetime.timezone对象,因为它们在 Airflow 场景下存在已知限制(无法携带 IANA 时区名称、转换行为有缺陷等)。Airflow 内部的时间处理统一基于 pendulum,仓库 shared/timezones 目录下的源码封装了相关逻辑,供深入研究者参考。
回顾:完整代码
完成上述步骤后,你的代码应当与仓库中的 tutorial.py 一致。整体结构如下:
- 导入模块(
textwrap、datetime、BashOperator、DAG); - 定义
default_args; - 用
with DAG(...)创建 DAG 对象; - 实例化
BashOperator得到任务t1、t2、t3; - 为任务与 DAG 添加
doc_md文档; - 定义 Jinja 模板命令;
- 用
t1 >> [t2, t3]声明依赖。
测试你的 Pipeline
写完之后就该测试了。第一步:确认脚本能通过解析。把代码保存为tutorial.py,放到airflow.cfg中dags_folder指定的 DAG 目录(默认如~/airflow/dags),然后运行:
python ~/airflow/dags/tutorial.py如果脚本无错误地运行结束,说明你的 DAG 结构定义正确。(注意:直接执行时with DAG(...)块内的任务实例化与依赖声明都会正常完成,但不会真的执行任务。)
命令行元数据校验
进一步用 CLI 命令验证元数据:
# 初始化数据库表 airflow db migrate # 打印所有已激活的 DAG 列表 airflow dags list # 打印 "tutorial" DAG 中的任务列表 airflow tasks list tutorial # 打印 "tutorial" DAG 的 graphviz 可视化表示 airflow dags show tutorial其中airflow db migrate会创建/更新元数据库 schema(首次使用 Airflow 时必须执行);airflow dags show tutorial需要安装 graphviz 支持,输出 DAG 结构的可视化描述。
测试任务实例与 DAG Run
你可以针对指定的**逻辑日期(logical date)**测试某个任务实例,这模拟了 Scheduler 在某个日期时间点上运行你的任务。
关于逻辑日期:请注意,Scheduler 是为某个具体日期时间运行你的任务,而不一定是在那个日期时间运行。逻辑日期(logical date)是 DAG Run 被命名的那个时间戳,它通常对应工作流所处理时间周期的结束时刻——或者是手动触发 DAG Run 的时刻。Airflow 用逻辑日期来组织和跟踪每次运行,你在 UI、日志和代码中都是通过它引用某次具体执行的。当通过 UI 或 API 触发 DAG 时,你也可以自行提供逻辑日期,从而"按某个时间点"运行工作流。
命令格式为:
# 命令布局: command subcommand [dag_id] [task_id] [(可选) 日期] # 测试 print_date 任务 airflow tasks test tutorial print_date 2015-06-01 # 测试 sleep 任务 airflow tasks test tutorial sleep 2015-06-01还可以查看模板是如何渲染的:
# 测试 templated 任务 airflow tasks test tutorial templated 2015-06-01这条命令会输出详细日志,并实际执行你的 bash 命令——你会看到模板循环展开后的真实命令内容。
需要记住:
airflow tasks test在本地运行任务实例,日志输出到 stdout,不在数据库中记录状态。它是调试单个任务实例的便捷工具;airflow dags test在本地运行整个 DAG Run,适合测试完整 DAG。与tasks test不同,它会创建真实的 DAG Run 并在元数据库中记录任务状态,因此需要已初始化的数据库,且 DAG 能被 Airflow 从你的 DAG 文件夹序列化。更多细节可参考 DAG 调试与 dag.test() 相关章节。
接下来做什么?
到这里,你已经成功编写并测试了第一个 Airflow 工作流。下一步,把代码合并到运行着Scheduler的代码仓库中,Scheduler 会接管你的 DAG,按schedule每天自动触发执行。
进阶方向:
- 继续学习 TaskFlow API 教程:用更 Pythonic 的方式定义工作流;
- 浏览 核心概念:深入理解 DAG、Task、Operator、Scheduler 等底层机制;
- 阅读 templates-ref.rst:掌握全部模板变量与宏;
- 参考 模块管理文档:了解 Airflow 如何加载 Python 模块。
掌握了本教程,你就拥有了 Airflow 世界中最重要的一块基石:从"会写脚本"到"能写出被调度器可靠执行的工作流"。
【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考