Airflow 进阶
概述
掌握 Airflow 基础后,进阶能力决定生产可用性:动态 DAG(按配置生成)、TaskFlow API(简洁依赖)、重试/超时/告警(可靠性)、Pool(资源控制)。本文逐一深入。
一、动态 DAG 生成
1.1 为什么动态
多张表 ETL 逻辑相同:
手写 N 个 DAG → 重复
按配置生成 → 一处维护1.2 实现
python
# 配置驱动生成
tables = ['orders', 'users', 'payments']
for table in tables:
dag_id = f'etl_{table}'
with DAG(dag_id=dag_id, schedule_interval='0 2 * * *', ...) as dag:
extract = PythonOperator(task_id='extract',
python_callable=extract_fn,
op_kwargs={'table': table})
load = PythonOperator(task_id='load',
python_callable=load_fn,
op_kwargs={'table': table})
extract >> load
globals()[dag_id] = dag| 方法 | 说明 |
|---|---|
| 循环生成 | 按列表生成 DAG |
| 配置文件 | YAML 驱动 |
| 工厂函数 | 封装 DAG 创建 |
1.3 注意
| 注意 | 说明 |
|---|---|
| 扫描耗时 | 动态文件别太多 |
| 幂等 | 生成确定性 |
| 版本管理 | 配置入 Git |
二、TaskFlow API
2.1 对比传统
传统:Operator + XCom(繁琐)
TaskFlow:Python 函数装饰器 + 自动依赖2.2 用法
python
from airflow.decorators import dag, task
@dag(schedule_interval='0 2 * * *', start_date=datetime(2026, 8, 1))
def etl_flow():
@task
def extract():
return {'total': 1000} # 自动 XCom
@task
def transform(data: dict):
data['total'] *= 2
return data
@task
def load(data: dict):
print(data)
load(transform(extract())) # 依赖自动推断
etl_flow()| 优势 | 说明 |
|---|---|
| 简洁 | 函数即任务 |
| 自动依赖 | 返回值传递 |
| 类型 | 函数签名 |
TaskFlow 内部仍用 XCom:
小数据适用
大数据走外部存储三、上下游依赖
3.1 跨 DAG 依赖
场景:DAG A 完成后跑 DAG B| 方式 | 说明 |
|---|---|
| ExternalTaskSensor | 等待上游任务 |
| TriggerDagRunOperator | 触发下游 DAG |
python
from airflow.sensors.external_task import ExternalTaskSensor
wait_upstream = ExternalTaskSensor(
task_id='wait_landing',
external_dag_id='landing_job',
external_task_id='load_done',
timeout=3600,
dag=dag,
)3.2 同 DAG 依赖
>> 与 << 操作符 + set_upstream/set_downstream依赖设计原则:
最小依赖(降耦合)
幂等可重跑
失败影响可控四、重试与超时
4.1 重试
python
default_args = {
'retries': 3, # 重试次数
'retry_delay': timedelta(minutes=5), # 间隔
'retry_exponential_backoff': True, # 指数退避
'max_retry_delay': timedelta(hours=1),
}| 参数 | 说明 |
|---|---|
| retries | 重试次数 |
| retry_delay | 间隔 |
| 指数退避 | 避免风暴 |
| max_retry | 上限 |
4.2 超时
python
PythonOperator(
task_id='t',
python_callable=fn,
execution_timeout=timedelta(hours=2), # 执行超时
dag=dag,
)| 参数 | 说明 |
|---|---|
| execution_timeout | 任务超时 |
| dagrun_timeout | DAG 超时 |
| 超时处理 | 失败/告警 |
超时设计:
防止任务挂死
明确 SLA五、告警
5.1 告警方式
| 方式 | 说明 |
|---|---|
| 邮件 | 默认 |
| 钉钉/企微 | 回调 |
| Webhook | 通用 |
| Slack | 团队 |
5.2 配置
python
default_args = {
'email': ['data@example.com'],
'email_on_failure': True,
'email_on_retry': False,
}python
# 回调函数
def on_failure(context):
ti = context['task_instance']
send_alert(f"Task {ti.task_id} failed: {context.get('exception')}")
PythonOperator(..., on_failure_callback=on_failure)5.3 告警内容
| 内容 | 说明 |
|---|---|
| 任务名 | 谁失败 |
| DAG | 哪个工作流 |
| 时间 | 执行时间 |
| 错误 | 异常信息 |
| 日志 | 链接 |
告警最佳实践:
失败必告警
重试可静默
分级(失败/延迟)六、池 Pool
6.1 作用
资源控制:
限制并发任务数
保护下游/系统python
# 定义 Pool(airflow pools 命令)
PythonOperator(
task_id='heavy',
python_callable=fn,
pool='heavy_pool', # 池名
priority_weight=10, # 优先级
pool_slots=1, # 占用槽位
dag=dag,
)| 概念 | 说明 |
|---|---|
| Pool | 资源池 |
| Slots | 并发槽位 |
| priority_weight | 优先级 |
| queue | 队列(Celery) |
6.2 使用场景
| 场景 | Pool |
|---|---|
| 数据库任务 | db_pool(限并发) |
| 重任务 | heavy_pool |
| 对外 API | api_pool |
Pool 意义:
防止并发打爆下游
保证关键任务优先七、其他进阶
7.1 Branch
分支任务(按条件走不同路径):
BranchPythonOperator7.2 触发规则
trigger_rule:
all_success(默认)
all_failed
one_success
none_failed7.3 回填与补数
catchup=True 补跑历史
手动 Trigger DAG(指定日期)八、生产实践
| 实践 | 说明 |
|---|---|
| DAG 版本化 | 代码入 Git |
| 配置化 | YAML 驱动 |
| 幂等任务 | 可重跑 |
| 监控告警 | 全链路 |
| 资源池 | 防风暴 |
| 日志 | 集中管理 |
| 常见坑 | 处理 |
|---|---|
| 动态 DAG 慢 | 控制规模 |
| 大 XCom | 外部存储 |
| 任务挂死 | 超时 |
| 并发风暴 | Pool |
| 上游失败 | 依赖+告警 |
常见问题速查
| 问题 | 要点 |
|---|---|
| 动态 DAG 怎么做 | 配置循环生成 |
| TaskFlow 优势 | 简洁自动依赖 |
| 跨 DAG 依赖 | Sensor/Trigger |
| 重试策略 | 指数退避 |
| 告警怎么配 | 回调 + 通知 |
| Pool 作用 | 并发控制 |