Airflow 必知必会(2)

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

1.概述

Airflow 是使用最广泛的开源工作流编排平台之一,理解 Airflow 的重要特性对于创建好的工作流至关重要。

2. DAG 对象详解

2.1 DAG 的定义

2.2 身份与元数据相关的参数

参数

类型

说明

默认值

dag_id

str

全局唯一标识符,UI/CLI/API 都用它。只允许字母数字、中划线、下划线

必填

description

str

DAG 描述,UI 上显示。具体在 Dag 的 Details 

""

tags

list[str]

标签,UI 里按标签筛选/分组

[]

owner

str

负责人(其实放在 default_args 里更常见)

"airflow"


doc_md


str

3.0+ 新增,UI 显示名(可含中文/空格)。具体在 Dag 的 Dag Docs 按钮

""


dag_display_name


str

3.0+ 新增,UI 显示名(可含中文/空格)

""

2.3 调度核心

参数

类型

说明

默认值

start_date

datetime

最早可以跑的时间点,也是数据区间的起点

必填

end_date

datetime

DAG 生命周期结束时间,到了就不再调度

None

schedule

str / timedelta / Asset / list

调度表达式(cron、预设、时间间隔、数据资产)

None(手动触发)

catchup

bool

是否补跑 start_date到当前的所有历史区间

True

max_active_runs

int

同时运行的最大 DagRun 数

None(无限制)

max_active_tasks

int

同时运行的最大 task 实例数(3.x 替代旧 concurrency)

None(用 airflow.cfg 的 max_active_tasks_per_dag)

2.4 依赖与串行控制

参数

类型

说明

默认值

depends_on_past

bool

同一个 DAG 内,当前 run 必须等上一次 run 成功才跑

False

wait_for_downstream

bool

配合 depends_on_past,等上一次 run 的所有下游 task 完成​

False

2.5 超时与失败处理

参数

类型

说明

默认值

dagrun_timeout

timedelta

单个 DagRun 超过这个时间就 fail

None

fail_stop

bool

3.x 新增:任意 task 失败立即停掉整个 DagRun

False

on_failure_callback

callable

DagRun 失败时的回调

None

on_success_callback

callable

DagRun 成功时的回调

None

on_execute_callback

callable

DagRun 开始执行时的回调

None

2.6 模板与Jinja

参数

类型

说明

默认值

template_searchpath

list[str]

模板文件搜索路径(你之前问过的)

None

template_undefined

jinja2.Undefined

模板变量未定义时的行为

StrictUndefined(严格模式)

user_defined_macros

dict

DAG 级自定义宏(你之前问过的)

{}

user_defined_filters

dict

DAG 级自定义 Jinja filter

{}

jinja_environment_kwargs

dict

传给 Jinja Environment 的额外参数

None

render_template_as_native_obj

bool

3.x:渲染结果自动转 Python 原生类型

False

2.7 default_args(批量传参给所有 task)

2.8 初始化与权限

参数

类型

说明

默认值

is_paused_upon_creation

bool

首次加载 DAG 时是否暂停

None(用 airflow.cfg 的 dags_are_paused_at_creation)

access_control

dict

角色权限控制(3.x 需要 apache-airflow-providers-fab)

None

params

dict

DAG 级参数,UI 手动触发时可覆盖,task 可以从上下文中的‘ "params" 关键词调用

{}

3. @task 的参数

3.1 基本参数

参数

类型

默认值

说明

task_id

str

函数名

覆盖默认任务 ID

multiple_outputs

bool

False

返回 dict时是否展开为多个 XCom

task_display_name

str

None

UI 显示名(Airflow 3.x)

3.2 Asset 相关

参数

类型

说明

outlets

list[Asset]

声明本 task 产出了哪些 Asset

inlets

list[Asset]

声明本 task 消费了哪些 Asset

3.3 重试与容错

参数

类型

默认值

说明

retries

int

0

失败重试次数

retry_delay

timedelta

timedelta(seconds=300)

重试间隔

retry_exponential_backoff

bool

False

指数退避

max_retry_delay

timedelta

timedelta(hours=1)

最大重试间隔

execution_timeout

timedelta

None

单次执行超时

timeout

timedelta

None

任务总超时(含重试)

3.4 触发规则与依赖

参数

类型

默认值

说明

trigger_rule

str

"all_success"

触发条件

depends_on_past

bool

False

是否依赖上一次 DagRun 成功

wait_for_downstream

bool

False

等待下游完成

ignore_first_depends_on_past

bool

False

第一次运行忽略 depends_on_past

3.5 资源与队列

参数

类型

说明

queue

str

指定 Celery 队列

pool

str

连接池名称

priority_weight

int

优先级权重(默认 1)

weight_rule

str

权重计算方式(downstream/upstream/absolute)

executor_config

dict

执行器配置(如 K8s executor 的 pod 配置)

3.6 回调与通知

参数

类型

说明

on_failure_callback

Callable

失败回调

on_success_callback

Callable

成功回调

on_retry_callback

Callable

重试回调

on_execute_callback

Callable

执行前回调

email_on_failure

bool

失败时发邮件

email_on_retry

bool

重试时发邮件

email

list[str]

收件人列表

3.7 显式与文档

参数

类型

说明

doc

str

任务文档

doc_md

str

Markdown 文档(UI 中渲染)

ui_color

str

UI 节点背景色

ui_fgcolor

str

UI 节点前景色

3.8 动态任务映射(Dynamic Task Mapping)

参数

类型

说明

map_index_template

str

映射任务的索引模板(Jinja)

3.9 其他常用参数

参数

说明

owner

任务负责人

sla

SLA(Service Level Agreement)

subdag

是否作为 SubDag(已弃用,不推荐)

task_group

所属 TaskGroup

params

参数定义

templates_dict

模板字典

show_return_value_in_logs

是否在日志中显示返回值(默认 True)

templates_exts

显式声明启用哪种后缀的文件模板。读文件内容,并渲染

4. Dynamic Task Mapping(动态任务映射)

动态任务映射(Dynamic Task Mapping)是 Airflow 2.3+ 引入、在 3.x 中成为运行时动态流水线首选方案的核心特性。它让你用一个任务模板在 DAG 运行时刻按数据展开成 N 个并行实例,彻底告别解析时 for循环写死 Task 节点的旧模式。

4.1 expand() 单参数展开

4.2 partial()+ expand()— 固定参数与变化参数分离

4.3 expand_kwargs() — 多参数逐行定义

4.4 多参数笛卡尔积

5. Airflow 三大调度范式

范式

核心机制

触发源

典型场景

时间驱动

Scheduler 按时间表建 DagRun

时钟(cron / Timetable)

每日 ETL、小时级报表

事件驱动

外部事件 → AssetEvent → DagRun

文件到达、消息队列、API 调用

数据就绪即处理、实时响应

编排驱动

上游 DAG 主动触发下游 DAG

TriggerDagRunOperator / 手动

强顺序流水线、CI/CD

5.1 时间驱动调度(Time-based)

5.1.1 cron/预设表达式

标准 cron 由 5 个字段​ 组成(Airflow 也支持 6 字段含秒,但一般用 5 字段)

字符

含义

示例

*

每(所有值)

*在分钟位 = 每分钟

,

列表(多个值)

1,15,30= 第 1、15、30 分钟

-

范围

9-17= 9 点到 17 点

/

步长

*/10= 每 10 个单位


预设

等价于 cron

含义

@once

只运行一次

@hourly

0 * * * *

每小时第 0 分

@daily

0 0 * * *

每天 00:00

@weekly

0 0 * * 0

每周日 00:00

@monthly

0 0 1 * *

每月 1 号 00:00

@yearly

0 0 1 1 *

每年 1 月 1 日 00:00

@none/ None

不调度(手动触发)

5.2 事件驱动调度(Event-driven)

外部事件源 → BaseEventTrigger (Triggerer) → AssetWatcher → AssetEvent → Scheduler → DagRun

a. 自定义 trigger 继承 BaseEventTrigger

b. asset = Asset("example_asset", watchers=[AssetWatcher(name="test_asset_watcher", trigger=trigger)]) 

c. 将 asset 传给 DAG 的参数 schedule

5.3 编排驱动调度(Orchestration-based)

trigger = TriggerDagRunOperator(
        task_id="trigger_downstream",
        trigger_dag_id="child_report",       # 要触发的 DAG ID
        wait_for_completion=False,            # 是否等下游 DAG 跑完
        deferrable=False,                     # True 则用 async 等待
        conf={"source_date": "{{ ds }}"},     # 传给下游 DAG 的 conf
        reset_dag_run=True,                   # 如果目标 DAG Run 已存在,重置
    )

6. 个人心得

6.1 把airflow各个组件分开,尤其是 Scheduler 和 Worker 要分开,尽量使用分布式架构和容器,有利于系统的稳定性

6.2 任务的幂等性

6.3 多实现自己的 Operator 或 Hook

6.4 定期清理元数据库

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

评论(0

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

全部回复

上滑加载中

设置昵称

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

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

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