Airflow 必知必会(2)
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 |
重试时发邮件 |
|
|
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)])
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 定期清理元数据库
- 点赞
- 收藏
- 关注作者
评论(0)