news 2026/9/11 13:15:01

Apache Airflow 101:使用 Airflow SDK 编写你的第一个 DAG 工作流(附完整源码解析与测试指南)

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Apache Airflow 101:使用 Airflow SDK 编写你的第一个 DAG 工作流(附完整源码解析与测试指南)

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 DAG
  • textwrap:标准库,用于textwrap.dedent去除多行字符串的公共缩进,让文档字符串与模板字符串书写更整洁;
  • datetime/timedelta:标准库,用于设置start_dateretry_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_callbackSLA 未达成的回调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,因此共享运行任务所需的核心参数(如retriesdepends_on_pastexecution_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继承来的通用参数retriesdepends_on_past)混合使用,这让代码更简洁。t2还把retries覆盖为3,演示了按任务覆盖默认值的能力。

任务参数的优先级如下:

  1. 显式传入的参数(如t2retries=3);
  2. default_args字典中的值(如retries=1);
  3. Operator 自身的默认值(如果存在)。

注意:每个任务必须包含或继承task_idowner两个参数,否则 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_macrosuser_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 date

Math 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 中任务之间可以互相依赖。假设有任务t1t2t3,可以用多种方式表达依赖:

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] << t1

tutorial 示例使用的正是列表形式的位移运算符写法:

t1 >> [t2, t3]

它表达"t1成功之后,t2t3并行运行"的扇形结构。

需要警惕的是,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 一致。整体结构如下:

  1. 导入模块(textwrapdatetimeBashOperatorDAG);
  2. 定义default_args
  3. with DAG(...)创建 DAG 对象;
  4. 实例化BashOperator得到任务t1t2t3
  5. 为任务与 DAG 添加doc_md文档;
  6. 定义 Jinja 模板命令;
  7. t1 >> [t2, t3]声明依赖。

测试你的 Pipeline

写完之后就该测试了。第一步:确认脚本能通过解析。把代码保存为tutorial.py,放到airflow.cfgdags_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),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/11 13:14:41

基于LSTM的三分类中文情感分析完整实现

简介&#xff1a;Python基于LSTM三分类的文本情感分析项目&#xff0c;面向计算机相关专业学生和需要实战练习的开发者&#xff0c;适用于课程设计、期末大作业或毕业设计参考。这是一份大三学生的期末项目&#xff0c;经导师指导并认可&#xff0c;评审分99分&#xff0c;代码…

作者头像 李华
网站建设 2026/9/11 13:13:05

PCSX2模拟器上手:BIOS怎么配、渲染器怎么选、掉帧怎么排查

PCSX2模拟器上手&#xff1a;BIOS怎么配、渲染器怎么选、掉帧怎么排查 【免费下载链接】pcsx2 PCSX2 - The Playstation 2 Emulator 项目地址: https://gitcode.com/GitHub_Trending/pc/pcsx2 PCSX2是一款开源的PS2模拟器&#xff0c;在Windows、Linux和macOS上用软件模…

作者头像 李华
网站建设 2026/9/11 13:11:40

WPF数据可视化实战:高性能动态图表与仪表盘开发

1. WPF数据可视化项目概述在工业控制、物联网监控和业务分析系统中&#xff0c;数据可视化始终是核心需求。最近我完成了一个基于WPF的实时数据监控项目&#xff0c;主要实现了动态折线图和仪表盘两大核心组件。这个方案完美替代了传统WinForm图表控件&#xff0c;在医疗监护设…

作者头像 李华
网站建设 2026/9/11 13:09:31

如何用 Vosk 三步搞定离线语音识别:完整指南

如何用 Vosk 三步搞定离线语音识别&#xff1a;完整指南 【免费下载链接】vosk-api Offline speech recognition API for Android, iOS, Raspberry Pi and servers with Python, Java, C# and Node 项目地址: https://gitcode.com/GitHub_Trending/vo/vosk-api Vosk 是一…

作者头像 李华
网站建设 2026/9/11 13:08:40

车载蓝牙六大协议协同开发实战指南

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华