Apache Airflow 深入
概述
Apache Airflow 是 Python 生态最流行的工作流调度平台:用 DAG 定义任务依赖,Operator 封装执行逻辑,Executor 决定调度方式,XCom 传递任务间数据。本文讲透 DAG 定义、Operator/Sensor/Hook、Celery/Kubernetes 执行器与 XCom。
一、Airflow 定位
| 特性 | 说明 |
|---|---|
| 类型 | 工作流调度 |
| 语言 | Python |
| 核心 | DAG(有向无环图) |
| 调度 | Cron 表达式 |
| 场景 | 数据管道、ETL |
适用:
ETL 定时任务
复杂依赖编排
数据管道编排二、DAG 定义
2.1 概念
DAG = 任务集合 + 依赖关系(有向无环)
Task:一个执行单元
DAG:任务编排python
from airflow import DAG
from airflow.operators.python import PythonOperator
from datetime import datetime, timedelta
default_args = {
'owner': 'data',
'retries': 1,
'retry_delay': timedelta(minutes=5),
}
with DAG(
dag_id='etl_pipeline',
schedule_interval='0 2 * * *', # 每天 2 点
start_date=datetime(2026, 8, 1),
default_args=default_args,
catchup=False,
) as dag:
extract = PythonOperator(task_id='extract', python_callable=extract_fn)
transform = PythonOperator(task_id='transform', python_callable=transform_fn)
load = PythonOperator(task_id='load', python_callable=load_fn)
extract >> transform >> load # 依赖| 参数 | 说明 |
|---|---|
| dag_id | DAG 名 |
| schedule_interval | 调度间隔 |
| start_date | 起始时间 |
| catchup | 是否补跑历史 |
| 依赖 | >> 操作符 |
2.2 依赖关系
python
a >> b >> c # 串行
a >> [b, c] # 并行
b << a # 反向
[a, b] >> c # 汇合DAG 规则:
无环
依赖明确
可并行三、Operator
3.1 概念
Operator = 任务执行逻辑封装
一个 Task = 一个 Operator 实例3.2 常用 Operator
| Operator | 说明 |
|---|---|
| PythonOperator | 执行 Python 函数 |
| BashOperator | 执行命令 |
| SparkSubmitOperator | 提交 Spark |
| HiveOperator | Hive 任务 |
| SQLOperator | SQL 执行 |
| DummyOperator | 空任务 |
python
from airflow.operators.bash import BashOperator
run_spark = BashOperator(
task_id='run_spark',
bash_command='spark-submit /scripts/etl.py --date {{ ds }}',
dag=dag,
)| 设计 | 说明 |
|---|---|
| 单一职责 | 一个任务一件事 |
| 幂等 | 可重跑 |
| 模板 | 等变量 |
四、Sensor
4.1 作用
等待条件满足:
文件到达
数据就绪
时间到达| Sensor | 说明 |
|---|---|
| FileSensor | 等待文件 |
| ExternalTaskSensor | 等待其他 DAG 任务 |
| TimeSensor | 等待时间 |
| SqlSensor | 等待 SQL 结果 |
python
from airflow.sensors.filesystem import FileSensor
wait_file = FileSensor(
task_id='wait_file',
filepath='/data/input/{{ ds }}/ready.txt',
poke_interval=60,
timeout=3600,
dag=dag,
)| 参数 | 说明 |
|---|---|
| poke_interval | 轮询间隔 |
| timeout | 超时 |
Sensor 场景:
文件依赖、上游完成、外部系统就绪五、Hook
5.1 作用
与外部系统交互的统一封装:
数据库、云服务、API| Hook | 说明 |
|---|---|
| BaseHook | 基类 |
| SshHook | SSH 执行 |
| HttpHook | HTTP 请求 |
| HiveHook | Hive 交互 |
Hook 使用:
Operator 内部调用
连接配置管理(Connections)python
from airflow.providers.mysql.hooks.mysql import MySqlHook
hook = MySqlHook(mysql_conn_id='mysql_default')
result = hook.get_records('SELECT COUNT(*) FROM orders')| 优势 | 说明 |
|---|---|
| 复用 | 连接管理 |
| 安全 | 密码加密存储 |
| 标准 | 统一接口 |
六、Executor 执行器
6.1 执行器类型
| 执行器 | 说明 |
|---|---|
| SequentialExecutor | 单机串行 |
| LocalExecutor | 单机多进程 |
| CeleryExecutor | 分布式 |
| KubernetesExecutor | K8s 动态 |
| CeleryKubernetes | 混合 |
生产推荐:
中型 → Celery
云原生 → Kubernetes6.2 CeleryExecutor
Scheduler 分发任务 → Celery Worker:
Worker 集群执行
队列管理架构:
Airflow Scheduler
Celery Broker(Redis/RabbitMQ)
Worker(多节点)
Result Backend| 优点 | 说明 |
|---|---|
| 分布式 | 多 Worker |
| 队列 | 任务隔离 |
| 成熟 | 稳定 |
6.3 KubernetesExecutor
每个任务一个 Pod:
动态创建
资源隔离| 优点 | 说明 |
|---|---|
| 弹性 | 按任务建 Pod |
| 隔离 | 独立环境 |
| 云原生 | K8s 生态 |
| 对比 | Celery | Kubernetes |
|---|---|---|
| 资源 | 常驻 Worker | 动态 Pod |
| 启动 | 快 | 慢(拉镜像) |
| 运维 | 管理 Worker | 管理 K8s |
七、XCom
7.1 作用
任务间传递数据:
小数据(结果/参数)python
# 推送
def extract_fn(**context):
context['ti'].xcom_push(key='date', value='2026-08-05')
# 拉取
def transform_fn(**context):
date = context['ti'].xcom_pull(task_ids='extract', key='date')7.2 注意
| 注意 | 说明 |
|---|---|
| 小数据 | 适合参数/状态 |
| 大数据 | 存外部(HDFS/DB) |
| 生命周期 | 随 DAG run |
XCom 局限:
不适合大对象
用文件/存储替代八、调度机制
8.1 调度流程
DAG 扫描 → 判断是否该跑 → 生成 DagRun
→ 创建 TaskInstance → 执行器调度
→ 执行 → 记录状态| 概念 | 说明 |
|---|---|
| DagRun | DAG 一次运行 |
| TaskInstance | 任务一次实例 |
| 状态 | success/failed/upstream_failed |
| 重试 | 自动重试 |
8.2 变量与模板
变量:
{{ ds }} 执行日期
{{ execution_date }}
{{ prev_ds }}九、运维要点
| 要点 | 说明 |
|---|---|
| 监控 | 任务状态/失败 |
| 告警 | 邮件/钉钉 |
| 日志 | 任务日志 |
| 回填 | catchup/手动触发 |
| 重跑 | 清空重跑 |
常见问题:
任务失败 → 看日志
调度不触发 → 检查时间/状态
并行不够 → 扩 Worker常见问题速查
| 问题 | 要点 |
|---|---|
| DAG 是什么 | 任务依赖图 |
| Operator 类型 | Python/Bash/Spark... |
| Sensor 用途 | 等待条件 |
| 执行器选型 | Celery/K8s |
| XCom 传什么 | 小数据 |
| 依赖怎么写 | >> 操作符 |