Apache Airflow 101:使用 Airflow SDK 编写你的第一个 DAG 工作流(附完整源码解析与测试指南) 📅 发布时间:2026/9/11 13:15:14 👁 浏览次数: 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 工作流的能力。什么是 DAGDAGDirected 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; well 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(minutes5), }, descriptionA simple tutorial DAG, scheduletimedelta(days1), start_datedatetime(2021, 1, 1), catchupFalse, tags[example], ) as dag: t1 BashOperator( task_idprint_date, bash_commanddate, ) t2 BashOperator( task_idsleep, depends_on_pastFalse, bash_commandsleep 5, retries3, ) 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_idtemplated, depends_on_pastFalse, bash_commandtemplated_command, ) t1 [t2, t3] # [END tutorial]理解 DAG 定义文件把 Airflow 的 Python 脚本想象成一个用代码描述 DAG 结构的配置文件——这一点非常重要你在其中定义的 task 实际运行在另一个环境Scheduler 派发、Worker 执行中因此这个脚本本身不是用来做数据处理的。它的主要职责是定义DAG对象并且必须能够被快速求值。原因在于Airflow 的 Dag File ProcessorDAG 文件处理器会定期检查 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; well 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是本次使用的 OperatorDAG来自airflow.sdk是本仓库中 DAG 对象的官方入口注意新版 Airflow 已将其收敛到airflow.sdk命名空间。关于 Python 与 Airflow 的模块管理机制例如哪些目录会被自动扫描、模块名如何解析可参考仓库文档 modules_management.rst。设置默认参数default_args创建 DAG 及其任务时你可以把参数直接传给每个 task也可以把一组公共参数放在字典中统一定义。后者通常更高效、更整洁因为所有 task 会默认继承这些参数同时仍允许在单个 task 上覆盖。tutorial 示例中的default_argsdefault_args{ depends_on_past: False, retries: 1, retry_delay: timedelta(minutes5), # queue: bash_queue, # pool: backfill, # priority_weight: 10, # end_date: datetime(2016, 1, 1), # wait_for_downstream: False, # execution_timeout: timedelta(seconds300), # 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当前任务实例是否依赖上一个调度周期的任务实例成功Falseretries失败后自动重试的次数1retry_delay两次重试之间的等待时长timedelta(minutes5)queue任务被派发到的队列bash_queuepool任务使用的资源池backfillpriority_weight任务优先级权重10end_date任务不再被调度的时间点datetime(2016, 1, 1)wait_for_downstream是否等待下游任务的上一个实例成功Falseexecution_timeout任务实例运行的最长时限超时则失败timedelta(seconds300)on_failure_callback/on_success_callback/on_retry_callback/on_skipped_callback失败/成功/重试/跳过时的回调函数或函数列表some_functionsla_miss_callbackSLA 未达成的回调yet_another_functiontrigger_rule任务触发条件规则all_success如果想深入了解BaseOperator的全部参数可查阅 airflow.sdk.BaseOperator 文档 及仓库中 BaseOperator 源码 对应的BaseOperator类定义。创建 DAG接下来实例化一个DAG对象来承载任务。tutorial 示例with DAG( tutorial, default_args{...}, descriptionA simple tutorial DAG, scheduletimedelta(days1), start_datedatetime(2021, 1, 1), catchupFalse, tags[example], ) as dag:各参数含义tutorial位置参数dag_id即 DAG 的唯一标识符。它在你的 Airflow 实例中必须全局唯一UI、CLI、API 都通过它引用该 DAGdefault_args上一步定义的默认参数字典会传递给每个任务descriptionDAG 的简要描述展示在 UI 的 DAG 列表中scheduletimedelta(days1)调度周期这里表示每天运行一次。除了timedelta还可以使用 cron 表达式字符串如0 0 * * *或daily等预设值start_dateDAG 开始生效的时间。注意 Airflow 的调度是结束时间对齐的第一个 DAG Run 的 logical date逻辑日期通常等于start_date但真正触发时间在其之后的一个调度周期catchupFalse关闭回填。若为True当 DAG 从start_date到当前时间之间有大量错过的调度周期时Scheduler 会一次性补跑所有错过的 DAG Run设为False则只调度最新周期避免刚上线时批量触发tags[example]为 DAG 打标签便于在 UI 中按标签筛选with ... as dag:上下文管理器语法块内实例化的任务会自动注册到该 DAG 上这是新版 Airflow 推荐、也是本仓库示例采用的写法。理解 OperatorOperator 是 Airflow 中的工作单元是构建工作流的积木决定了任务将要执行什么动作。所有 Operator 都继承自BaseOperator因此共享运行任务所需的核心参数如retries、depends_on_past、execution_timeout等。社区与官方提供大量 Operator常见的有PythonOperator执行 Python 可调用对象BashOperator执行 Bash 命令或脚本本教程主角KubernetesPodOperator在 Kubernetes 集群中拉起 Pod 执行任务以及各类 Provider 提供的专有 Operator可浏览仓库 providers 目录。除了直接实例化 OperatorAirflow 还提供更Pythonic的 TaskFlow API可以用装饰器把普通 Python 函数变成任务本文暂不展开。定义任务Task要使用 Operator必须先把它实例化为任务Task。任务决定了 Operator 在 DAG 上下文中如何执行工作。task_id是每个任务的唯一标识符。tutorial 示例实例化了两次BashOperatort1 BashOperator( task_idprint_date, bash_commanddate, ) t2 BashOperator( task_idsleep, depends_on_pastFalse, bash_commandsleep 5, retries3, )注意这里把Operator 专有参数bash_command与BaseOperator继承来的通用参数retries、depends_on_past混合使用这让代码更简洁。t2还把retries覆盖为3演示了按任务覆盖默认值的能力。任务参数的优先级如下显式传入的参数如t2的retries3default_args字典中的值如retries1Operator 自身的默认值如果存在。注意每个任务必须包含或继承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_idtemplated, depends_on_pastFalse, bash_commandtemplated_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_commandtemplated_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 tasks 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, tzAsia/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),仅供参考