Airflow 基础

举报
剑指南天 发表于 2026/09/10 12:30:43 2026/09/10
【摘要】 Airflow是一个以编程的方式编排和监视工作流的平台,原理是将工作流编写为若干任务按照依赖顺序组成的有向无环图(DAG)。

1.概述

Airflow是一个以编程的方式编排和监视工作流的平台,原理是将工作流编写为若干任务按照依赖顺序组成的有向无环图(DAG)。

Airflow是编排调度器,不是流处理引擎,不能处理实时数据流;它也不是 ETL 引擎本身——任务里真正的重活(抽数、转换、计算)由算子调用的外部系统完成。Airflow负责的是什么时候、按什么顺序、由谁去干

2. 架构

3. 核心组件

Metadata Database:stores the state of tasks,Dags and variables——所有组件的唯一事实源。组件之间不直接互相调用,全部通过元数据库交换状态

API Server:提供 REST API ,Web UI 和 Task Execution API。

Dag processor:解析 Dag bundle 里的 Python 文件,把解析结果序列化后写进元数据库。

Scheduler:“monitors all tasks and Dags, then triggers the task instances once their dependencies are complete”——触发到期的计划工作流、把就绪任务提交给 Executor。生产环境可跑多个 Scheduler 实现高可用,一致性靠元数据库的行级锁。

Triggerer:在 asyncio 事件循环里等外部事件,让进入 deferred 状态的任务不占着 worker 槽位。

Worker:真正执行任务的进程(Celery worker 进程 / K8s Pod),只在分布式 Executor 下出现。

4. 关键概念

DAG:一个 Dag 就是一条工作流 —— A Dag is a model that encapsulates everything needed to execute a workflow,封装了执行一个工作流所需的全部内容:任务、依赖、调度规则。

Task:工作流的最小工作单元。

Operator:预定义的任务模板,官方称 “a reusable, pre-made Task template”。如 BashOperator 执行一条 bash 命令、PythonOperator 调用一个 Python 函数;把它实例化就得到一个 Task。

Sensor:专门等待外部条件的算子:等一个文件出现、等另一个 DAG 的某个任务完成。官方定义:它们 “wait for something to occur”。

DAG Run:DAG 的运行实例。官方定义:“an object representing an instantiation of the Dag in time”(DAG 在时间轴上的一次实例化)。

Task Instance:DAG Run中的 task 运行实例,官方定义 “a specific run of that task for a given Dag”。

XCom:任务之间传递小批量元数据的机制(Cross-communications 的缩写)。

Executor:决定任务以什么方式并发执行(决定任务在哪儿、以多少并发执行) —— LocalExecutor、CeleryExecutor、KubernetesExecutor。

Catchup:决定 DAG 从 start_date到“现在”之间,如果因为新建/暂停后恢复/改调度而漏跑了很多个调度点,Scheduler 是否自动把每个漏掉的区间都建 Dag Run 补跑。

Backfill :人为指定时间范围重跑。不一定从 start_date开始,可早于 DAG 当前生效范围。

Asset:数据/事件驱动。

git-sync:一个容器定期 git pull到共享目录,scheduler/dag-processor/api-server/worker 都挂同一目录。

DAG Bundle:DAG Bundle 是 Airflow 3.x 里“DAG 代码来源”的抽象:不只是 .py,还包括 DAG 依赖的 includes、配置、SQL、脚本等一整包资源。

logical date:标记的是数据区间的起点,不是实际跑起来的时刻。

5. 一个 DAG 从部署到执行的完整流程

  1. Dag processor 扫描 Dag bundle,解析 Python 文件,把 DAG 的序列化结果写进元数据库(默认每分钟收集一轮解析结果)。
  2. Scheduler 检查各 DAG 的 timetable。时间一到,Scheduler 为这个 DAG 创建新的 DAG Run。
  3. Scheduler 逐个检查这次 DAG Run 里的 Task Instance,根据依赖关系标记 Task Instance 的状态,提交状态为 scheduled 的任务实例到 Executor 。
  4. 在 Pool 与并发限制的约束下,Scheduler 把 scheduled 的任务实例提交给 Executor 的队列(状态 → queued)。这一步用数据库行级锁保护:多个 Scheduler 并存时,同一任务不会被重复提交。
  5. Worker 从队列取到任务,拉起进程执行,运行日志实时采集。
  6. 任务结束,结果回报 Scheduler,状态写入数据库:成功 → success,下游任务的依赖解除,回到第 c 步继续推进;失败但还有重试次数 → up_for_retry,间隔 retry_delay 后重新入队;失败且重试耗尽 → failed,它的下游会被标记 upstream_failed。
  7. 当所有「叶子任务」(没有下游的任务)都到达终态,整条 DAG Run 被标记为 success 或 failed。
【版权声明】本文为华为云社区用户原创内容,未经允许不得转载,如需转载请自行联系原作者进行授权。如果您发现本社区中有涉嫌抄袭的内容,欢迎发送邮件进行举报,并提供相关证据,一经查实,本社区将立刻删除涉嫌侵权内容,举报邮箱: cloudbbs@huaweicloud.com
  • 点赞
  • 收藏
  • 关注作者

评论(0

0/1000
抱歉,系统识别当前为高风险访问,暂不支持该操作

全部回复

上滑加载中

设置昵称

在此一键设置昵称,即可参与社区互动!

*长度不超过10个汉字或20个英文字符,设置后3个月内不可修改。

*长度不超过10个汉字或20个英文字符,设置后3个月内不可修改。