Airflow 必知必会(3)

举报
剑指南天 发表于 2026/09/17 20:37:49 2026/09/17
【摘要】 Airflow 是使用最广泛的开源工作流编排平台之一,理解 Airflow 的重要特性对于创建好的工作流至关重要。

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 能够成为最受欢迎的工作流编排工具的原因之一。

【版权声明】本文为华为云社区用户原创内容,未经允许不得转载,如需转载请自行联系原作者进行授权。如果您发现本社区中有涉嫌抄袭的内容,欢迎发送邮件进行举报,并提供相关证据,一经查实,本社区将立刻删除涉嫌侵权内容,举报邮箱: cloudbbs@huaweicloud.com
  • 点赞
  • 收藏
  • 关注作者

评论(0)

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

全部回复

上滑加载中

设置昵称

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

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

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