Airflow 必知必会(3)
1.概述
Airflow 是使用最广泛的开源工作流编排平台之一,理解 Airflow 的重要特性对于创建好的工作流至关重要。

2. DAG(准确说是 DagRun)和 Task(TaskInstance)各有独立的状态机
2.1 DagRun 状态
|
状态 |
含义 |
何时出现 |
|---|---|---|
|
queued |
排队等待调度器分配 |
3.x 新增,手动触发或 scheduler 创建后、还没开始跑 |
|
running |
正在运行 |
至少有一个 task 在跑或 pending |
|
success |
全部 task 成功 |
所有 task 都 success(或被 skip 但 trigger_rule 允许) |
|
failed |
运行失败 |
至少一个 task 失败,且没有其他路径可走 |
|
skipped |
整体被跳过 |
3.x 新增,所有 entry task 都被 skipped(如 branch 全部未选中) |
2.2 TaskInstance 状态
|
状态 |
含义 |
何时出现 |
|---|---|---|
|
none |
还没被调度 |
DAG 刚解析,task 还没进入调度队列 |
|
scheduled |
已调度,等 executor 取走 |
scheduler 判定可以跑了,放入 executor 队列 |
|
queued |
executor 已取走,等 worker |
已发给 worker 进程/线程,还没开始执行 |
|
running |
正在执行 |
worker 里代码在跑 |
|
success |
执行成功 |
task 正常结束,exit code 0 |
|
failed |
执行失败 |
抛异常、exit code 非 0、或被外部 kill |
|
skipped |
被跳过 |
branch 未选中 / ShortCircuit 返回 False / trigger_rule 不满足 |
|
up_for_retry |
等待重试 |
task 失败且 retries > 0,在 retry_delay倒计时中 |
|
up_for_reschedule |
Sensor reschedule 等待 |
Sensor mode="reschedule"时,poke 返回 False 后释放 worker |
|
deferred |
异步挂起 |
deferrable task 调 self.defer()后,释放 worker 等 triggerer |
|
removed |
代码已删除 |
DAG 文件改了,这个 task_id 不再存在,旧 run 里标 removed |
3. 装饰器 @asset 和 Asset
|
|
@asset 装饰器 |
Asset() 对象 |
|---|---|---|
|
是什么 |
装饰器,把函数体绑定为"生产逻辑" |
构造函数,创建一个纯 URI 标识符 |
|
返回值 |
Asset 实例(函数被包装了) |
Asset 实例 |
|
有没有生产逻辑 |
有,函数体就是 |
没有,只是一个名字 |
|
谁用得多 |
定义"自己产出的数据" |
定义"别人产出的、我只消费" |
3.1 Asset 是"数据可用性事件"的抽象
a. Asset是一个 Python 对象(不是装饰器)
b. 它代表一个逻辑数据位置(URI 字符串)
d. 不存储数据,只记录"这个数据什么时候被更新了"
e. 通过 outlets/ inlets声明生产/消费关系
f. 下游 DAG 通过 schedule=[Asset(...)]被事件触发
g. 数据本身通过 S3/DB/文件等外部系统共享
3.2 @asset
@asset是把函数定义变成一个 Asset 对象,同时把函数体绑定为该 Asset 的"生产逻辑"。
核心规则:函数体在 Asset 被"物化"时执行。
a. schedule自动触发
b. 被 DAG 中的 task 调用

c. 触发下游DAG

producer按 @daily物化 → 触发 consumerDAG 执行。但 producer的函数体是在物化时执行的,不是在 consumer里。
4. Docker 启动 Airflow
参考 https://airflow.apache.org/docs/apache-airflow/3.3.1/howto/docker-compose/index.html#docker-compose-env-variables
4.1 准备工作
使用 root 用户,创建 airflow 用户,记录用户 id(AIRFLOW_UID),将用户加入 docker 组
useradd airflow;id -u airflow;usermod -aG docker airflow;su airflow
进入 airflow 用户 home 目录,创建自定义 airflow 文件夹
cd /home/airflow
mkdir airflow331;cd airflow331
拉取 docker-compose.yaml 文件
curl -LfO 'https://airflow.apache.org/docs/apache-airflow/3.3.1/docker-compose.yaml'
创建 airflow 服务自有文件夹
mkdir -p ./dags ./logs ./plugins ./config
编辑 .env 文件

使用下面命令生成 FERNET_KEY 的值
python3 -c 'import base64, os; print(base64.urlsafe_b64encode(os.urandom(32)).decode())'
查看 SELinux 状态,如果开启的话修改 docker-compose.yaml

4.2 初始化数据库
docker compose up airflow-init
结束显示:airflow-init-1 exited with code 0
4.3 启动 airflow
docker compose up -d

4.4 验证挂载目录是否成功

5. 总结
Airflow 提供的丰富组件和功能是 Airflow 能够成为最受欢迎的工作流编排工具的原因之一。
- 点赞
- 收藏
- 关注作者


评论(0)