Article
第一章:Airflow 概述与核心概念
1.1 什么是 Apache Airflow
| 概念名称 | 说明 | 注意事项 |
|---|---|---|
| Apache Airflow | 一个开源的、由 Airbnb 开发并捐赠给 Apache 基金会的工作流自动化平台,用于编排、调度和监控复杂的数据流水线(Data Pipelines)。 | Airflow 不是数据处理引擎(如 Spark),而是工作流调度器,用于定义任务依赖与执行顺序。 |
| 工作流(Workflow) | 一组具有依赖关系的任务集合,按预定逻辑和时间执行。Airflow 使用 DAG 来建模工作流。 | 工作流应具备可重复性、可观测性和可恢复性。 |
| DAG(Directed Acyclic Graph) | 有向无环图,Airflow 中工作流的基本组织单位,定义任务及其执行顺序。 | 图中不能存在循环依赖,否则调度器将拒绝加载。 |
| 代码即配置(Code as Configuration) | 使用 Python 脚本定义工作流,使 DAG 可版本控制、可测试、可复用。 | 提高了灵活性,但也要求开发者具备一定 Python 编程能力。 |
1.2 Airflow 的四大核心组件(DAGs, Operators, Tasks, Executors)
| 组件名称 | 说明 | 注意事项 |
|---|---|---|
| DAGs(有向无环图) | 工作流的容器,定义一组任务及其依赖关系。每个 .py 文件通常包含一个或多个 DAG 对象。 | 必须设置 start_date 和 schedule_interval,否则 DAG 不会触发。 |
| Operators(操作符) | 定义具体任务行为的类,如执行 Shell 命令、调用 Python 函数等。是任务的”模板”。 | Operators 是声明式的,不立即执行,仅在调度时实例化为 TaskInstance。 |
| Tasks(任务) | DAG 中的节点,由 Operator 实例化而来,代表一个具体的执行单元。 | 任务必须有唯一 task_id,并在 DAG 内定义依赖关系。 |
| Executors(执行器) | 负责实际执行任务的组件。决定任务在何处以及如何运行。 | 不同执行器支持不同并发模型和部署方式,选择需结合生产需求。 |
1.3 Airflow 架构解析(Web Server, Scheduler, Metadata Database, Worker)
| 架构组件 | 说明 | 注意事项 |
|---|---|---|
| Web Server | 提供可视化界面(UI),用于查看 DAG 状态、日志、触发任务、管理变量和连接等。基于 Flask 实现。 | 应配置 HTTPS 和用户认证(如 RBAC)以保障安全。 |
| Scheduler | 核心调度进程,持续监控 DAG 文件夹,根据调度周期和依赖关系提交任务到 Executor。 | 需保持高可用,避免单点故障;支持多实例但需协调。 |
| Metadata Database | 存储元数据,包括 DAG 结构、任务实例状态、变量、连接、日志元信息等。常用 PostgreSQL 或 MySQL。 | 数据库性能直接影响 Airflow 整体表现,需定期维护和备份。 |
| Worker | 执行具体任务的进程或容器,由 Executor 分配任务。在 Local/Celery/Kubernetes 模式下运行。 | Worker 需具备执行任务所需的依赖环境(如 Python 包、工具命令)。 |
| DAG Processor | 解析 DAG 文件的后台进程(Airflow 2.0+ 中增强),将 DAG 定义加载到数据库。 | 多个 DAG 文件时,解析可能成为瓶颈,可调整 dag_dir_list_interval。 |
1.4 Airflow 的适用场景与优势
| 类别 | 说明 | 注意事项 |
|---|---|---|
| 适用场景:ETL 流程 | 自动化抽取、转换、加载数据,从多种源系统到数据仓库。 | 适合批处理,实时流处理建议结合 Kafka/Flink。 |
| 适用场景:数据管道编排 | 协调多个任务(如数据清洗、模型训练、报告生成)的依赖与执行。 | 可结合传感器(Sensor)等待外部事件。 |
| 适用场景:定时作业调度 | 替代 crontab,提供更强大的依赖管理、重试机制和监控能力。 | 更适合复杂依赖,简单定时任务可能”杀鸡用牛刀”。 |
| 优势:可视化 UI | 提供 DAG 图、任务日志、运行历史等,便于调试与监控。 | UI 功能强大,但大量 DAG 时可能加载较慢。 |
| 优势:可扩展性 | 支持自定义 Operator、Hook、Executor,易于集成新系统。 | 自定义组件需充分测试,避免影响调度稳定性。 |
| 优势:社区活跃 | Apache 顶级项目,文档丰富,插件生态完善。 | 注意版本兼容性,升级前需评估变更影响。 |
| 优势:代码即配置 | DAG 使用 Python 编写,支持参数化、循环、条件等编程特性。 | 需规范代码结构,避免 DAG 文件过于复杂。 |
第二章:环境搭建与基础配置
2.1 安装 Airflow(pip, conda, Docker)
| 安装方式 | 说明 | 注意事项 |
|---|---|---|
| pip 安装 | 使用 pip install apache-airflow 安装到本地 Python 环境。 | 建议使用虚拟环境(如 venv 或 conda),避免依赖冲突。 |
| conda 安装 | 通过 conda-forge 通道安装:conda install -c conda-forge apache-airflow。 | conda 更好管理复杂依赖(如数据库驱动),适合数据科学环境。 |
| Docker 安装 | 使用官方镜像 apache/airflow,适合快速搭建和生产部署。 | 需配置 volume 持久化 DAG 和日志,注意用户权限(UID/GID)。 |
| Airflow Quick Start | 使用官方 Docker Compose 文件快速启动包含 Web Server、Scheduler、DB 的完整环境。 | 仅用于本地测试,生产环境需拆分部署并配置高可用。 |
2.2 初始化数据库与 Web 服务器启动
| 操作步骤 | 说明 | 注意事项 |
|---|---|---|
airflow db init | 初始化元数据库,创建所有必要表(如 dag, task_instance, variable)。 | 首次运行前确保数据库服务已启动且连接配置正确。 |
airflow users create | 创建第一个用户(用户名、密码、角色等),用于登录 Web UI。 | 至少创建一个 Admin 用户,否则无法登录。 |
airflow webserver | 启动 Web Server,默认监听 8080 端口。 | 生产环境建议配置反向代理(如 Nginx)和 HTTPS。 |
airflow scheduler | 启动调度器,开始解析 DAG 并调度任务。 | 调度器需持续运行,建议使用进程管理工具(如 systemd、supervisord)。 |
| 日志路径 | 默认日志存储在 $AIRFLOW_HOME/logs。 | 可通过配置更改路径,生产环境建议对接集中式日志系统。 |
2.3 Airflow 配置文件 airflow.cfg 解析
| 配置项(节.参数) | 说明 | 注意事项 |
|---|---|---|
core.dags_folder | DAG 文件存放路径,默认为 $AIRFLOW_HOME/dags。 | 确保路径可读,DAG 文件需在此目录或其子目录下。 |
core.load_examples | 是否加载示例 DAG,设为 False 可减少干扰。 | 生产环境建议关闭。 |
core.executor | 执行器类型:SequentialExecutor, LocalExecutor, CeleryExecutor 等。 | 单机测试可用 Local,生产建议 Celery 或 Kubernetes。 |
webserver.web_server_port | Web Server 监听端口,默认 8080。 | 若端口被占用,可修改后重启 webserver。 |
scheduler.scheduler_zone | 调度器时区,默认为 UTC。 | 建议保持 UTC,避免本地时区带来的混乱。 |
core.default_timezone | DAG 默认时区,影响 start_date 解析。 | 可设为 Asia/Shanghai 等,但需与调度器时区协调。 |
logging.base_log_folder | 日志存储根目录。 | 确保磁盘空间充足,定期清理旧日志。 |
2.4 使用环境变量配置 Airflow
| 环境变量 | 对应配置项 | 用途说明 | 注意事项 |
|---|---|---|---|
AIRFLOW__CORE__EXECUTOR | core.executor | 设置执行器类型,如 CeleryExecutor。 | 双下划线 __ 表示节与参数的分隔,大小写不敏感。 |
AIRFLOW__DATABASE__SQL_ALCHEMY_CONN | core.sql_alchemy_conn | 数据库连接字符串,如 postgresql://user:pass@host:5432/airflow。 | 敏感信息建议通过 secrets backend 管理。 |
AIRFLOW__WEBSERVER__WEB_SERVER_PORT | webserver.web_server_port | 修改 Web Server 端口。 | 修改后需重启 webserver 生效。 |
AIRFLOW_HOME | 环境变量本身 | Airflow 主目录,默认为 ~/airflow。 | 可自定义路径,所有配置、日志、数据库默认在此目录下。 |
AIRFLOW__CORE__LOAD_EXAMPLES | core.load_examples | 控制是否加载示例 DAG。 | 设为 false 可避免测试环境干扰。 |
2.5 访问 Web UI 与基本界面介绍
| UI 区域 | 说明 | 注意事项 |
|---|---|---|
| DAGs 列表页 | 显示所有已加载的 DAG,状态(如 Running, Success, Failed)、最近运行时间、调度周期等。 | 点击 DAG 名称进入详情页,可过滤和搜索。 |
| DAG 详情页(Graph View) | 展示任务依赖图,不同颜色表示任务状态(绿色=成功,红色=失败,黄色=运行中)。 | 可右键任务节点进行手动操作(如 Clear, Run)。 |
| Tree View | 以树形结构展示多个 DagRun 及其任务执行情况,适合查看历史运行。 | 横轴为时间,纵轴为任务,便于分析执行耗时。 |
| Grid View | 表格形式展示任务实例状态,支持筛选和批量操作。 | 适合快速定位失败任务。 |
| Logs 查看 | 点击任务实例可查看详细日志输出,包括上下文和执行命令。 | 日志支持分页和搜索,是调试关键工具。 |
| Admin 菜单 | 包含 Variables、Connections、Users、Plugins 等管理功能。 | Variables 和 Connections 常用于参数化 DAG。 |
第三章:DAG 基础与生命周期
3.1 DAG 的定义与基本结构
| 方法/参数 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| DAG 类初始化 | DAG(dag_id, description=None, schedule_interval=None, start_date=None, end_date=None, catchup=True, default_args=None, params=None, tags=None) | 创建一个 DAG 实例,定义工作流的基本属性。 | from airflow import DAGfrom datetime import datetimedag = DAG( dag_id='example_dag', description='A simple tutorial DAG', schedule_interval='@daily', start_date=datetime(2025, 1, 1), catchup=False, tags=['example']) | dag_id 必须唯一且不包含空格或特殊字符(可用下划线);start_date 是必需的。 |
3.2 DAG 的导入上下文与调度周期(schedule_interval)
| 参数/值 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
schedule_interval | schedule_interval='@daily' 或 timedelta(hours=2) 或 cron 表达式 | 定义 DAG 的调度频率,控制任务何时触发。 | from datetime import timedeltadag = DAG( dag_id='my_dag', schedule_interval=timedelta(days=1), start_date=datetime(2025, 1, 1))# 或使用 cronschedule_interval='0 0 * * *' # 每天零点 | None 表示仅手动触发;@once 表示只运行一次;避免使用 @hourly 等模糊表达影响可读性。 |
| 内置字符串别名 | @once, @hourly, @daily, @weekly, @monthly, @yearly | 简化常见调度周期的定义。 | schedule_interval='@daily' | 实际对应 cron 表达式,如 @daily 等价于 0 0 * * *。 |
timedelta 调度 | timedelta(minutes=30), timedelta(hours=2) | 用于固定时间间隔调度。 | from datetime import timedeltaschedule_interval=timedelta(minutes=15) | 适用于简单周期任务,不支持复杂时间逻辑。 |
| cron 表达式 | '分 时 日 月 周' | 精确控制调度时间,支持复杂周期。 | schedule_interval='30 8 * * 1-5' # 工作日 8:30 | 注意时区问题,默认为 UTC;建议使用 time 模块或注释说明本地时间。 |
3.3 DAG 的 start_date、end_date 与 catchup 控制
| 参数 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
start_date | start_date=datetime(2025, 1, 1) | 定义 DAG 首次执行的逻辑时间点,用于计算调度实例。 | from datetime import datetimestart_date=datetime(2025, 9, 1, 0, 0) | 必须设置;若使用 datetime.now() 可能导致调度异常;建议使用 pendulum 或固定时间。 |
end_date | end_date=datetime(2025, 12, 31) | 定义 DAG 最后一次执行的时间,超过此时间不再调度。 | end_date=datetime(2025, 12, 31) | 可选;常用于临时任务或项目周期限制。 |
catchup | catchup=True 或 catchup=False | 控制是否补跑历史 DagRun(从 start_date 到当前时间)。 | catchup=False | 若设为 True 且 start_date 较早,可能导致大量任务堆积;生产环境建议设为 False。 |
catchup_by_default | 在 airflow.cfg 中设置 | 全局默认 catchup 行为。 | catchup_by_default = False | 可通过配置统一管理,避免每个 DAG 重复设置。 |
3.4 DAG 的参数化配置(params, default_args)
| 参数 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
default_args | default_args = { 'owner': 'airflow', 'retries': 1, 'retry_delay': timedelta(minutes=5) } | 定义 DAG 中所有任务的默认参数,减少重复代码。 | default_args = { 'owner': 'data_team', 'retries': 2, 'retry_delay': timedelta(minutes=10)}dag = DAG('my_dag', default_args=default_args, ...)# 单个任务可覆盖task = PythonOperator( task_id='my_task', retries=3, # 覆盖 default_args ...) | 支持 owner, retries, retry_delay, email, depends_on_past 等;任务级设置优先级更高。 |
params | params={'file_path': '/data/input.csv'} | 定义 DAG 级参数,可在任务中通过 Jinja 模板或上下文访问。 | dag = DAG( 'parametrized_dag', params={ 'region': 'cn', 'batch_size': 1000 })# 在模板中使用:{{ params.region }} | 常用于动态配置任务行为;支持 Web UI 手动触发时传参。 |
3.5 DAG 的上下文管理与执行生命周期
| 概念 | 说明 | 注意事项 |
|---|---|---|
| 执行生命周期阶段 | 1. DAG 文件解析 → 2. DagRun 创建 → 3. TaskInstance 调度 → 4. 任务执行 → 5. 状态更新 → 6. 日志记录 | 整个流程由 Scheduler 和 Executor 协同完成。 |
| 上下文变量(Context) | Airflow 在执行任务时注入的一组变量,如 {{ ds }}, {{ execution_date }}, {{ dag_run }} 等。 | 可在 Operator 的 op_args, templates_dict 等中使用 Jinja 模板访问。 |
execution_date | DAG 实例的逻辑执行时间,非实际运行时间,用于确定数据处理的时间窗口。 | execution_date 是调度时间点,通常滞后于 start_date;例如 @daily 的 execution_date 是前一天 00:00。 |
templates_dict | 允许将任意字段模板化,通过上下文变量动态赋值。 | templates_dict={'message': 'Processing data for {{ ds }}'} |
第四章:任务与 Operators
4.1 Operator 概述与分类(BaseOperator, BashOperator, PythonOperator 等)
| Operator 类 | 说明 | 注意事项 |
|---|---|---|
BaseOperator | 所有 Operator 的基类,提供通用参数如 task_id, owner, retries 等。 | 不可直接实例化,用于继承自定义 Operator。 |
BashOperator | 执行 Shell 命令或脚本。 | 依赖系统环境,确保命令在 Worker 上可用。 |
PythonOperator | 调用 Python 函数,支持传参和上下文。 | 函数必须定义在可导入位置;避免阻塞调度器。 |
EmailOperator | 发送邮件通知,常用于任务成功/失败提醒。 | 需配置 SMTP 连接(conn_id='smtp')。 |
SimpleHttpOperator | 向 HTTP 端点发送请求(GET/POST 等)。 | 支持基本认证、JSON 数据发送。 |
BranchPythonOperator | 根据 Python 函数返回值决定后续执行路径。 | 返回值应为 task_id 或 task_id 列表。 |
DummyOperator | 不执行任何操作,仅用于流程控制(如占位、合并分支)。 | 轻量级,常用于逻辑分组。 |
4.2 BashOperator:执行 Shell 命令
| 参数 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
bash_command | bash_command='echo "Hello Airflow"' | 指定要执行的 Shell 命令或脚本。 | from airflow.operators.bash import BashOperatortask = BashOperator( task_id='print_date', bash_command='date', dag=dag) | 支持多行命令(用分号或 && 连接);命令失败会触发重试。 |
env | env={'MY_VAR': 'value'} | 设置环境变量传递给 Shell 命令。 | env={'ENV': 'prod', 'PATH': '/usr/local/bin'} | 可覆盖默认环境;敏感信息建议使用 Variables。 |
output_encoding | output_encoding='utf-8' | 指定命令输出的编码格式。 | output_encoding='gbk' | 处理中文输出时可能需要设置。 |
4.3 PythonOperator:调用 Python 函数
| 参数 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
python_callable | python_callable=my_function | 指定要调用的 Python 函数。 | def greet(**context): print(f"Hello from {context['ds']}")task = PythonOperator( task_id='greet_task', python_callable=greet, dag=dag) | 函数必须已定义且可导入;不支持 lambda。 |
op_args | op_args=[arg1, arg2] | 传递位置参数给函数。 | op_args=['data.csv', 100] | 参数在函数调用时按顺序传入。 |
op_kwargs | op_kwargs={'name': 'Alice'} | 传递关键字参数给函数。 | op_kwargs={'filename': 'report.txt'} | 常用于配置化参数。 |
provide_context | 已废弃(Airflow 2.0+) | 旧版本用于自动注入上下文。 | —— | 新版本上下文自动提供,无需设置。 |
4.4 EmailOperator:发送邮件通知
| 参数 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
to | to='user@example.com' 或 ['a@x.com','b@y.com'] | 邮件接收者地址。 | to='team@company.com' | 支持单个或多个收件人。 |
subject | subject='DAG Success' | 邮件主题。 | subject='ETL Job Completed on {{ ds }}' | 支持 Jinja 模板。 |
html_content | html_content='<h1>Success</h1>' | 邮件正文(HTML 格式)。 | html_content='Data processed for {{ ds }}' | 也支持纯文本。 |
conn_id | conn_id='smtp_default' | SMTP 连接 ID,需提前在 UI 中配置。 | conn_id='custom_smtp' | 默认使用 smtp_default,需确保连接有效。 |
4.5 SimpleHttpOperator:调用 HTTP 接口
| 参数 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
http_conn_id | http_conn_id='my_api' | HTTP 连接 ID,包含 URL 基地址、认证等。 | http_conn_id='jira_api' | 需在 UI 中预先配置 Connection。 |
endpoint | endpoint='/users' | 请求的端点路径(附加到 base URL 后)。 | endpoint='/tasks/create' | 不包含协议和主机。 |
method | method='GET' 或 'POST' | HTTP 请求方法。 | method='POST' | 支持 GET, POST, PUT, DELETE 等。 |
data | data={'key': 'value'} | POST 请求的表单数据。 | data={'name': 'test'} | 也可使用 json 参数发送 JSON。 |
json | json={'payload': 1} | 发送 JSON 格式请求体。 | json={'event': 'trigger'} | 自动设置 Content-Type 为 application/json。 |
4.6 BranchPythonOperator:条件分支控制
| 参数 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
python_callable | python_callable=branch_func | 返回目标 task_id 的函数。 | def decide_branch(**context): if context['ds'] == '2025-09-30': return 'task_a' else: return 'task_b'branch_task = BranchPythonOperator( task_id='branch_task', python_callable=decide_branch, dag=dag) | 返回值必须是存在的 task_id。 |
follow_task_ids_if_false | follow_task_ids_if_false=['task_c'] | 当无分支匹配时执行的任务(可选)。 | —— | 通常不需要设置,未选中任务状态为 skipped。 |
4.7 DummyOperator:占位与流程控制
| 参数 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
task_id | task_id='start' | 唯一任务标识。 | from airflow.operators.dummy import DummyOperatorstart = DummyOperator( task_id='start', dag=dag) | 无实际执行逻辑,仅用于依赖管理。 |
| —— | —— | —— | —— | 常用于合并多个分支(join)或作为流程起点/终点。 |
第五章:任务依赖与执行流程
5.1 使用 >> 和 << 定义任务依赖
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
>> (右移操作符) | task_a >> task_b | 定义 task_a 执行成功后执行 task_b,即 task_b 依赖于 task_a。 | from airflow.operators.dummy import DummyOperatorstart = DummyOperator(task_id='start')process = DummyOperator(task_id='process')end = DummyOperator(task_id='end')start >> process >> end | 推荐写法,简洁直观;支持链式调用。 |
<< (左移操作符) | task_b << task_a | 等价于 task_a >> task_b,定义 task_b 依赖于 task_a。 | process << start # 等同于 start >> process | 语义上”task_b 从 task_a 接收”,较少使用,但功能相同。 |
5.2 set_upstream() 与 set_downstream() 方法
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
set_downstream() | task.set_downstream(downstream_task) | 设置当前任务的下游任务。 | start.set_downstream(process)process.set_downstream(end) | 可接受单个任务或任务列表;适用于动态依赖构建。 |
set_upstream() | task.set_upstream(upstream_task) | 设置当前任务的上游任务。 | end.set_upstream(process) # 等价于 process >> end | 与 set_downstream 互为反向操作;建议统一风格。 |
| 接收列表 | task.set_downstream([t1, t2]) | 一次性设置多个下游任务。 | start.set_downstream([task1, task2, task3]) | 适用于扇出(fan-out)结构。 |
5.3 多任务依赖与链式调用
| 模式 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| 扇出(Fan-out) | task >> [t1, t2, t3] | 一个任务完成后并行触发多个下游任务。 | start >> [extract_a, extract_b, extract_c] | 三个任务将并行执行,互不依赖。 |
| 扇入(Fan-in) | [t1, t2, t3] >> task | 多个任务完成后触发一个下游任务。 | [transform_a, transform_b] >> load | 所有上游任务必须成功,下游才执行。 |
| 链式调用 | t1 >> t2 >> t3 >> t4 | 定义线性任务流。 | task1 >> task2 >> task3 >> task4 | 最常见模式,清晰表达执行顺序。 |
| 混合依赖 | (t1 >> t2) << t3 | 组合多种依赖关系。 | (extract >> transform) << DummyOperator(task_id='trigger') | 使用括号控制优先级,避免歧义。 |
5.4 任务组(TaskGroup)的使用(Airflow 2.0+)
| 参数/方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| TaskGroup 上下文管理器 | with TaskGroup(group_id, tooltip=None) as tg: | 将多个任务组织为逻辑组,在 UI 中折叠显示。 | from airflow.utils.task_group import TaskGroupwith TaskGroup('data_processing') as tg: step1 = BashOperator(...) step2 = PythonOperator(...) step1 >> step2 | 提升 DAG 可读性,避免 UI 过于复杂。 |
group_id | group_id='preprocess' | 任务组的唯一标识。 | group_id='feature_engineering' | 必须唯一,不能与其他任务或组重名。 |
prefix_group_id | prefix_group_id=False | 是否将 group_id 作为组内任务 task_id 的前缀。 | prefix_group_id=True # 默认值 | 设为 False 可自定义 task_id。 |
tooltip | tooltip='Data cleaning steps' | 鼠标悬停时显示的提示信息。 | tooltip='This group handles raw data' | 增强可维护性,建议添加描述。 |
5.5 子 DAG 与 TriggerDagRunOperator
| 概念/方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| 子 DAG(SubDAG) | 独立 DAG 文件中定义,被主 DAG 调用 | 模块化复杂流程(已不推荐) | def subdag(parent_dag_name, child_dag_name, args): dag = DAG(...) # 定义子任务 return dagfrom airflow.operators.subdag import SubDagOperatorsubdag_op = SubDagOperator( task_id='subdag_task', subdag=subdag('parent', 'child', args)) | SubDagOperator 已弃用,存在性能问题和锁竞争。 |
| TriggerDagRunOperator | TriggerDagRunOperator(trigger_dag_id='target_dag') | 主动触发另一个 DAG 的运行,实现跨 DAG 编排。 | from airflow.operators.trigger_dagrun import TriggerDagRunOperatortrigger = TriggerDagRunOperator( task_id='trigger_etl', trigger_dag_id='daily_etl', conf={'region': 'cn'} # 传递参数) | 推荐替代 SubDAG 的方式;支持传递 conf 参数。 |
conf | conf={'key': 'value'} | 传递给被触发 DAG 的配置字典,可在其上下文中通过 dag_run.conf 访问。 | conf={'source': 'web', 'batch': 100} | 实现参数化调度,增强灵活性。 |
第六章:XCom 与任务间通信
6.1 XCom 机制原理与使用场景
| 概念 | 说明 | 注意事项 |
|---|---|---|
| XCom(Cross-Communication) | Airflow 中任务间传递小量数据的机制,全称 “Cross-Communication”。 | 用于传递状态、文件路径、结果标识等,非大数据传输。 |
| push | 任务将数据写入 XCom 表,供其他任务读取。 | 通常由 return 值或 xcom_push() 显式触发。 |
| pull | 任务从 XCom 表读取数据。 | 使用 {{ task_instance.xcom_pull(...) }} 或 Python 代码。 |
| 数据存储位置 | 存储在元数据库的 xcom 表中。 | 默认不加密,敏感信息需处理。 |
| 使用场景 | - 上游任务返回文件路径 - 条件分支判断依据 - 动态生成任务参数 | 不适合传输大文件(如 CSV、日志),建议使用外部存储(S3、HDFS)。 |
6.2 push 与 pull XCom 数据
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| 自动 push | return value in PythonOperator | 函数返回值自动作为 XCom 推送,key=‘return_value’。 | def get_message(): return "Hello"task = PythonOperator( python_callable=get_message, task_id='push_msg') | 默认行为,可通过 do_xcom_push=False 关闭。 |
| 手动 push | task_instance.xcom_push(key='status', value='success') | 在任务中显式推送自定义键值对。 | def manual_push(**context): context['task_instance'].xcom_push( key='process_id', value=12345 ) | 适用于推送多个值或非返回值数据。 |
| pull(Jinja) | {{ task_instance.xcom_pull(task_ids='task_a', key='return_value') }} | 在模板字段中拉取 XCom 数据。 | BashOperator( bash_command='echo {{ task_instance.xcom_pull(task_ids="get_msg") }}', ...) | 仅适用于支持模板的字段(如 bash_command)。 |
| pull(Python) | ti.xcom_pull(task_ids='task_a') | 在 Python 函数中拉取数据。 | def consume(**context): msg = context['task_instance'].xcom_pull( task_ids='push_msg' ) print(msg) | 需传入 **context 或显式获取 task_instance。 |
6.3 使用 PythonOperator 传递 XCom
| 场景 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| 返回值传递 | return 语句 | 最简单方式,返回值自动 push。 | def step1(): return '/tmp/data.csv'def step2(**context): path = context['task_instance'].xcom_pull(task_ids='step1') print(f"Processing {path}") | 适合单一结果传递。 |
| 多 key 传递 | xcom_push(key='key_name', value=val) | 推送多个命名数据项。 | def multi_push(**context): ti = context['task_instance'] ti.xcom_push(key='count', value=100) ti.xcom_push(key='status', value='done') | 避免键名冲突,建议使用有意义的 key。 |
| 拉取多个任务 | xcom_pull(task_ids=['t1','t2']) | 从多个上游任务拉取数据。 | data = ti.xcom_pull(task_ids=['extract_us', 'extract_eu']) | 返回值为列表,顺序与 task_ids 一致。 |
6.4 XCom 后端配置与性能注意事项
| 配置项/建议 | 说明 | 注意事项 |
|---|---|---|
xcom_backend | 在 airflow.cfg 中设置自定义 XCom 后端类。 | 如 airflow.models.xcom.BaseXCom 的子类,可用于压缩或外部存储。 |
| 数据大小限制 | 元数据库字段类型(如 MySQL TEXT)限制单条 XCom 大小。 | 建议不超过 48KB;大对象应存入外部存储(S3、数据库),仅传递 URI。 |
| 性能影响 | 频繁读写 XCom 会增加数据库压力。 | 避免在循环中大量 push/pull;考虑缓存或批处理。 |
| 敏感数据 | XCom 默认明文存储。 | 不要传递密码、密钥等;使用 Connection 或 Variable 加密存储。 |
| 清理策略 | XCom 数据随 DagRun 清理而删除。 | 启用 delete_child_runs 等配置确保自动清理,避免表膨胀。 |
第七章:Variables 与 Connections
7.1 Variables:全局变量管理
| 概念 | 说明 | 注意事项 |
|---|---|---|
| Variables | Airflow 中的全局键值存储,用于在 DAG 和任务之间共享配置数据(如路径、开关、参数)。 | 适合存储轻量级配置,不适用于频繁读写或大数据。 |
| 存储位置 | 默认存储在元数据库的 variable 表中,支持 JSON 序列化。 | 敏感变量可启用加密(需配置 fernet_key)。 |
| 访问方式 | 可通过 UI、CLI、Python API 或 Jinja 模板访问。 | 建议在 default_args 或任务中使用 Variable.get() 动态读取。 |
| 延迟加载 | 变量在 DAG 解析时加载,若变量不存在会抛出异常(除非设置 default_var)。 | 避免在 DAG 文件顶层直接调用 Variable.get(),可能导致解析失败;建议在任务执行时调用或使用 try/except。 |
7.2 使用 Variable.set() 与 Variable.get()
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
Variable.set() | Variable.set(key, value, description=None, serialize_json=False) | 设置或更新一个变量。 | from airflow.models import VariableVariable.set( key='api_endpoint', value='https://api.example.com/v1', description='Production API URL')# 存储字典Variable.set('config', {'batch': 1000, 'region': 'cn'}, serialize_json=True) | serialize_json=True 时,值将被 JSON 序列化;否则转为字符串。 |
Variable.get() | Variable.get(key, default_var=None, deserialize_json=False) | 获取变量值。 | # 获取字符串url = Variable.get('api_endpoint')# 获取 JSON 对象config = Variable.get('config', deserialize_json=True)# 安全获取(带默认值)timeout = Variable.get('timeout', default_var=30) | 若变量不存在且未设 default_var,会抛出 KeyError;生产环境建议始终提供默认值。 |
| Jinja 模板访问 | {{ var.value.<variable_key> }} | 在 Operator 模板字段中使用变量。 | BashOperator( task_id='print_url', bash_command='curl {{ var.value.api_endpoint }}') | 仅适用于支持模板的字段;JSON 变量需用 var.json.key 访问。 |
7.3 Connections:外部系统连接配置
| 概念 | 说明 | 注意事项 |
|---|---|---|
| Connection | 封装外部系统(数据库、API、FTP 等)的连接信息,包括主机、端口、登录凭据等。 | 避免在代码中硬编码敏感信息;通过 UI 或 CLI 配置。 |
conn_id | 连接的唯一标识符,用于在 Operator 或 Hook 中引用。 | 建议命名规范,如 postgres_prod, smtp_notification。 |
| 字段组成 | 包括:Conn Type、Host、Schema、Login、Password、Port、Extra(JSON 格式附加参数)。 | Password 自动加密存储(需配置 fernet_key);Extra 可用于传递 SSL 配置、时区等。 |
| 配置方式 | 可通过 Web UI(Admin → Connections)、CLI(airflow connections add)或环境变量设置。 | 环境变量格式:AIRFLOW_CONN_<CONN_ID>=<connection_uri>,如 AIRFLOW_CONN_POSTGRES_DB=postgres://user:pass@host:5432/db。 |
7.4 在 Operator 中使用 Connection
| Operator | 使用方式 | 代码示例 | 注意事项 |
|---|---|---|---|
| PostgresOperator | postgres_conn_id='my_postgres' | PostgresOperator( task_id='run_query', sql='SELECT * FROM users;', postgres_conn_id='prod_db') | 需确保 Connection 类型为 postgres。 |
| MySqlOperator | mysql_conn_id='my_mysql' | MySqlOperator( task_id='insert_data', sql='INSERT INTO logs VALUES (...)', mysql_conn_id='analytics_db') | 类似 PostgresOperator,适用于 MySQL。 |
| SimpleHttpOperator | http_conn_id='my_api' | SimpleHttpOperator( task_id='call_api', endpoint='/data', http_conn_id='jira_api') | Connection 类型应为 http,Host 字段填写基础 URL。 |
| EmailOperator | conn_id='smtp_default' | EmailOperator( to='admin@company.com', subject='Alert', html_content='Error occurred', conn_id='smtp_prod') | 默认使用 smtp_default,需配置 SMTP 服务器信息。 |
| BashOperator / PythonOperator | 通过 Hook 使用 | 见第八章示例 | 通用方式,适用于所有系统。 |
第八章:Hooks 与外部系统集成
8.1 Hook 概念与作用
| 概念 | 说明 | 注意事项 |
|---|---|---|
| Hook | Airflow 中封装与外部系统交互逻辑的类,提供统一接口访问数据库、API、文件系统等。 | 是 Operator 的底层支持,也可在 PythonOperator 中直接调用。 |
| 作用 | - 抽象连接管理(自动读取 Connection) - 封装常用操作(如 execute, get_records) - 支持重试、日志、上下文集成 | 减少重复代码,提高可维护性。 |
| 命名规范 | 通常为 [System]Hook,如 PostgresHook, HttpHook。 | 使用前需安装对应 provider 包,如 apache-airflow-providers-postgres。 |
| 与 Operator 关系 | Operator 调用 Hook 执行实际操作,Hook 更底层、更灵活。 | 自定义 Operator 时常用 Hook 实现核心逻辑。 |
8.2 BashHook:与 Shell 交互
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
run_command() | BashHook().run_command(command, env=None) | 执行 Shell 命令并返回结果。 | from airflow.hooks.bash import BashHookhook = BashHook()result = hook.run_command('ls -la /tmp')print(result.output) | 返回 CommandResult 对象,包含 exit_code, output, error。 |
| 直接使用 | BashHook().run_command(...) | 快速执行命令,常用于 PythonOperator。 | def run_shell(**context): hook = BashHook() hook.run_command('python /scripts/process.py') | 注意命令路径和环境变量;失败不会自动重试。 |
8.3 PythonHook:调用 Python 上下文
| 说明 | 注意事项 |
|---|---|
| PythonHook 并非 Airflow 标准 Hook。Python 操作通常通过 PythonOperator 直接调用函数,或在自定义 Hook 中使用 Python 逻辑。 | 不需要专门的 Hook 来执行 Python 代码;PythonOperator 已封装执行环境。 |
| 上下文访问 | 在 python_callable 中通过 **context 参数获取 DAG 运行信息。 |
| ⚠️ 说明:Airflow 无 PythonHook 类。Python 脚本执行由 PythonOperator 和 PythonVirtualenvOperator 处理。 |
8.4 PostgresHook:连接 PostgreSQL
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
get_conn() | hook.get_conn() | 获取 psycopg2 连接对象。 | conn = hook.get_conn()cursor = conn.cursor() | 可用于执行复杂事务或自定义操作。 |
run() | hook.run(sql, parameters=None) | 执行 SQL 语句(DDL/DML)。 | hook.run('CREATE TABLE IF NOT EXISTS tmp (id INT)') | 自动提交事务;支持参数化查询。 |
get_records() | hook.get_records(sql) | 执行查询并返回所有记录(列表套元组)。 | rows = hook.get_records('SELECT * FROM users')for row in rows: print(row) | 适合小结果集;大查询建议用 get_pandas_df。 |
get_first() | hook.get_first(sql) | 返回查询的第一行。 | first_user = hook.get_first('SELECT name FROM users LIMIT 1') | 适用于单值查询。 |
get_pandas_df() | hook.get_pandas_df(sql) | 返回 pandas DataFrame。 | df = hook.get_pandas_df('SELECT * FROM sales') | 需安装 pandas;适合数据分析任务。 |
| 使用示例 | —— | —— | from airflow.providers.postgres.hooks.postgres import PostgresHookdef query_data(**context): hook = PostgresHook(postgres_conn_id='prod_db') df = hook.get_pandas_df('SELECT * FROM orders WHERE ds = {{ ds }}') return df.shape[0] | 在 PythonOperator 中调用此函数,返回值可自动 XCom Push。 |
8.5 MySqlHook:连接 MySQL
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
get_conn() | hook.get_conn() | 获取 MySQL 连接(使用 mysqlclient 或 PyMySQL)。 | conn = hook.get_conn() | 注意驱动兼容性。 |
run() | hook.run(sql, parameters=None) | 执行 SQL 命令。 | hook.run('INSERT INTO logs VALUES (%s, %s)', parameters=[1, 'info']) | 支持参数化防止 SQL 注入。 |
get_records() | hook.get_records(sql) | 获取查询结果列表。 | results = hook.get_records('SELECT user_id FROM active_users') | 返回格式同 PostgresHook。 |
get_first() | hook.get_first(sql) | 获取第一行结果。 | user = hook.get_first('SELECT name FROM users ORDER BY id DESC') | 适用于单记录查询。 |
get_pandas_df() | hook.get_pandas_df(sql) | 返回 DataFrame。 | df = hook.get_pandas_df('SELECT * FROM transactions') | 需安装 pandas。 |
| 使用示例 | —— | —— | from airflow.providers.mysql.hooks.mysql import MySqlHookdef check_count(): hook = MySqlHook(mysql_conn_id='analytics_db') count = hook.get_first('SELECT COUNT(*) FROM events')[0] return count > 1000 | 可用于 BranchPythonOperator 判断分支。 |
8.6 HttpHook:HTTP 请求封装
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
get_conn() | hook.get_conn() | 获取 requests Session 对象。 | session = hook.get_conn() | 可复用连接,支持持久化。 |
run() | hook.run(endpoint=None, data=None, json=None, headers=None, extra_options=None) | 发送 HTTP 请求。 | response = hook.run( endpoint='/users', json={'name': 'Alice'}, headers={'Content-Type': 'application/json'}) | 方法由 Connection 的 method 字段决定,或通过 extra_options 覆盖。 |
| 使用示例 | —— | —— | from airflow.providers.http.hooks.http import HttpHookdef call_api(**context): hook = HttpHook( method='POST', http_conn_id='my_api' ) response = hook.run( endpoint='/events', json={'event': 'trigger', 'ds': context['ds']} ) return response.json() | 可在 PythonOperator 中调用,返回 API 响应用于后续处理。 |
第九章:调度与执行器
9.1 调度器(Scheduler)工作原理
| 概念 | 说明 | 注意事项 |
|---|---|---|
| Scheduler 作用 | Airflow 的核心组件,负责解析 DAG 文件、创建 DagRun、调度 TaskInstance、监控任务状态。 | 需持续运行,通常作为守护进程启动。 |
| DAG 解析 | Scheduler 周期性扫描 dags_folder,加载所有 .py 文件并实例化 DAG 对象。 | 解析时间过长会影响调度性能;避免在 DAG 文件中执行耗时操作。 |
| DagRun 创建 | 根据 schedule_interval 和 start_date 判断是否需要创建新的 DagRun 实例。 | 创建时间基于 execution_date(逻辑时间),非当前时间。 |
| TaskInstance 调度 | 对于就绪的 DagRun,Scheduler 将符合条件的 TaskInstance 状态设为 scheduled,交由 Executor 执行。 | 调度依赖任务依赖关系、资源限制、depends_on_past 等参数。 |
| 心跳机制 | Scheduler 定期更新自身状态(heartbeat),确保高可用模式下只有一个活跃实例。 | 多 Scheduler 部署时依赖元数据库锁机制。 |
| 性能优化 | 可通过 parsing_processes、min_file_process_interval 等配置优化解析性能。 | 大量 DAG 文件时建议启用并行解析并延长扫描间隔。 |
9.2 常用调度参数(schedule_interval, interval, cron 表达式)
| 参数 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
schedule_interval | schedule_interval='@daily' 或 timedelta(hours=2) 或 '0 8 * * *' | 定义 DAG 的调度周期。 | # 每天 8:00schedule_interval='0 8 * * *'# 每 2 小时schedule_interval=timedelta(hours=2) | None 表示仅手动触发;@once 表示运行一次。 |
| 内置别名 | @once, @hourly, @daily, @weekly, @monthly, @yearly | 简化常见调度周期。 | schedule_interval='@hourly' | @daily 等价于 0 0 * * *(UTC 时间)。 |
timedelta | timedelta(minutes=30), timedelta(days=1) | 用于固定间隔调度。 | from datetime import timedeltaschedule_interval=timedelta(minutes=15) | 不支持复杂时间逻辑(如工作日)。 |
| cron 表达式 | '分 时 日 月 周' | 精确控制调度时间。 | '30 9 * * 1-5' # 工作日 9:30 | 注意时区:Airflow 默认使用 UTC;可通过 timezone 参数调整。 |
| None | schedule_interval=None | 仅支持手动触发或由 TriggerDagRunOperator 触发。 | schedule_interval=None | 常用于临时任务或事件驱动流程。 |
9.3 执行器类型对比(Sequential, Local, Celery, Kubernetes)
| 执行器 | 说明 | 优点 | 缺点 | 适用场景 |
|---|---|---|---|---|
| SequentialExecutor | 顺序执行器,任务串行运行,使用单进程。 | 无需外部依赖,配置简单,适合本地测试。 | 无法并行,性能极低。 | 本地开发、学习、单元测试。 |
| LocalExecutor | 本地执行器,支持多进程并行,在单机上运行任务。 | 支持并行,无需额外服务,性能优于 Sequential。 | 仍受限于单机资源,不支持分布式。 | 单机部署、中小规模 DAG。 |
| CeleryExecutor | 分布式执行器,通过消息队列(如 RabbitMQ、Redis)协调多个 Worker。 | 支持跨多机扩展,高可用,适合大规模任务。 | 架构复杂,需维护消息队列和 Worker 集群。 | 生产环境、大规模并行任务。 |
| KubernetesExecutor | 动态执行器,每个任务在独立的 Kubernetes Pod 中运行。 | 资源隔离好,弹性伸缩,环境隔离。 | 依赖 Kubernetes 集群,冷启动较慢,成本较高。 | 云原生环境、多租户、资源波动大场景。 |
9.4 CeleryExecutor 配置与分布式部署
| 配置项 | 说明 | 示例值 | 注意事项 |
|---|---|---|---|
executor | 在 airflow.cfg 中设置执行器类型。 | executor = CeleryExecutor | 需重启 Scheduler 和 Worker 生效。 |
broker_url | 消息队列地址,用于任务分发。 | redis://localhost:6379/0 或 amqp://user:pass@rabbitmq:5672// | 推荐使用 Redis 或 RabbitMQ;确保高可用。 |
result_backend | 任务结果存储地址,用于状态同步。 | db+postgresql://user:pass@db:5432/airflow 或 redis://... | 必须与 broker_url 分开;建议使用数据库或 Redis。 |
| worker 启动 | 在各 Worker 节点运行 airflow celery worker。 | airflow celery worker -c 4 | -c 指定并发数;Worker 需能访问 DAG 文件和依赖库。 |
| Scheduler 配置 | 启动 airflow scheduler,连接同一元数据库和消息队列。 | —— | 确保 Scheduler 与 Worker 网络互通。 |
| 高可用 | 可部署多个 Worker 和 Scheduler(仅一个活跃)。 | 使用 celery multi 管理 Worker 进程。 | Scheduler 通过元数据库锁实现主备。 |
9.5 KubernetesExecutor 简介
| 概念 | 说明 | 注意事项 |
|---|---|---|
| 工作原理 | 每个 TaskInstance 由 Scheduler 提交为一个独立的 Kubernetes Pod,执行完成后自动销毁。 | 无需预启动 Worker,资源按需分配。 |
| 配置文件 | airflow.cfg 中设置 executor = KubernetesExecutor。 | 需配置 kubernetes 相关参数。 |
pod_template_file | 指定 Pod 模板 YAML 文件,定义镜像、资源限制、卷挂载等。 | pod_template_file = /path/to/pod_template.yaml |
dags_in_image | DAG 文件是否已打包在镜像中。 | 若为 True,则无需 git-sync;否则需配置 git_repo。 |
| 优势 | - 弹性伸缩 - 环境隔离 - 多语言支持 - 与 K8s 生态集成 | 适合云原生架构。 |
| 缺点 | - 冷启动延迟(Pod 创建时间) - 成本较高(按 Pod 计费) - 配置复杂 | 不适合高频短任务。 |
第十章:错误处理与重试机制
10.1 任务失败与重试参数(retries, retry_delay)
| 参数 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
retries | retries=3 | 定义任务失败后的最大重试次数。 | PythonOperator( task_id='process_data', retries=2, retry_delay=timedelta(minutes=5)) | 默认为 0;设置为 0 表示不重试。 |
retry_delay | retry_delay=timedelta(minutes=5) | 两次重试之间的等待时间。 | from datetime import timedeltaretry_delay=timedelta(seconds=30) | 支持 timedelta 或 Relativity(Airflow 2.0+)。 |
retry_exponential_backoff | retry_exponential_backoff=True | 是否启用指数退避重试(1次: delay, 2次: 2×delay, 3次: 4×delay…)。 | retry_exponential_backoff=True | 可缓解服务压力,避免雪崩。 |
max_retry_delay | max_retry_delay=timedelta(hours=1) | 指数退避时的最大延迟时间。 | max_retry_delay=3600 # 秒 | 防止延迟过长。 |
10.2 on_failure_callback 与全局回调
| 回调类型 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
on_failure_callback | on_failure_callback=my_func | 任务失败时执行的回调函数。 | def notify_failure(context): print(f"Task {context['task_instance'].task_id} failed!")PythonOperator( on_failure_callback=notify_failure, ...) | 函数必须接受 context 参数;常用于发邮件、告警。 |
| 全局配置 | 在 default_args 中设置 | 为 DAG 中所有任务统一设置回调。 | default_args = { 'on_failure_callback': global_alert, 'retries': 1} | 任务级回调优先级更高。 |
| DAG 级回调 | on_failure_callback in DAG | DAG 整体失败时触发(Airflow 2.2+ 支持 on_failure_callback for DagRun)。 | dag = DAG( 'my_dag', on_failure_callback=dag_failed, ...) | 需注意版本支持;可用于清理资源或通知。 |
10.3 使用 Trigger Rules 控制失败任务流
| Trigger Rule | 说明 | 适用场景 | 注意事项 |
|---|---|---|---|
all_success | 默认规则,所有上游任务成功才执行。 | 普通线性流程。 | 任一上游失败则跳过。 |
all_failed | 所有上游任务失败才执行。 | 故障恢复、清理任务。 | 常用于”全部失败后报警”或”重试前检查”。 |
all_done | 所有上游任务完成(无论成功或失败)即执行。 | 汇总、归档、通知任务。 | 保证后续任务总能执行。 |
one_success | 至少一个上游任务成功即执行。 | 分支并行处理,任一成功即继续。 | 适用于”或”逻辑。 |
one_failed | 至少一个上游任务失败即执行。 | 监控、告警任务。 | 可快速响应异常。 |
none_failed | 上游任务中没有失败的(即全部成功或跳过)即执行。 | 安全后续操作。 | 比 all_success 更宽松(允许 skipped)。 |
none_skipped | 上游任务中没有被跳过的即执行。 | 确保所有路径都被处理。 | 常用于数据一致性检查。 |
10.4 TaskInstance 与 DagRun 的失败状态处理
| 概念 | 说明 | 注意事项 |
|---|---|---|
| TaskInstance 状态 | 包括 success, failed, up_for_retry, skipped, up_for_reschedule 等。 | 失败后根据 retries 决定是否重试;达到上限后标记为 failed。 |
| DagRun 状态 | success, failed, running, queued 等。 | 当所有任务完成且无失败时为 success;任一任务最终失败则为 failed。 |
| 手动重试 | 在 Web UI 中点击”Clear”可重置任务状态,触发重试。 | 会清除历史尝试记录;可用于修复后重新运行。 |
| 跳过机制 | 由 BranchPythonOperator 或 ShortCircuitOperator 导致任务状态为 skipped。 | skipped 不影响 all_success 规则,但影响 none_skipped。 |
| 失败处理策略 | - 自动重试 - 告警通知(on_failure_callback) - 进入备用流程(Trigger Rule) - 手动干预 | 建议结合监控系统(如 Prometheus + Alertmanager)实现自动化告警。 |
| 日志查看 | 失败任务的日志可在 Web UI 的”Logs”标签中查看。 | 日志包含错误堆栈,是排查问题的主要依据。 |
第十一章:高级特性与最佳实践
11.1 DAG 的动态生成
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| 循环生成 | for item in list: | 根据配置列表批量创建相似 DAG 或任务。 | configs = ['site_a', 'site_b']for site in configs: dag = DAG( f'crawl_{site}', schedule_interval='@daily', default_args=default_args ) # 定义任务... globals()[dag.dag_id] = dag | 使用 globals() 注册 DAG 到模块级作用域,Airflow 才能识别。 |
| 函数生成 | def create_dag(dag_id, schedule, ...) | 封装 DAG 创建逻辑,提高复用性。 | def create_etl_dag(dag_id, source): dag = DAG(dag_id, ...) extract = PythonOperator(..., dag=dag) return dagdag_a = create_dag('etl_user', 'users')globals()['etl_user'] = dag_a | 避免重复代码;适用于多数据源、多租户场景。 |
| 从外部加载 | json.load(), YAML, Database | 从配置文件或数据库读取元数据生成 DAG。 | import jsonwith open('dags_config.json') as f: config = json.load(f)for cfg in config['dags']: dag = DAG(cfg['id'], ...) globals()[cfg['id']] = dag | 实现配置驱动的调度系统,便于集中管理。 |
11.2 使用 Jinja 模板引擎
| 特性 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| 变量插值 | {{ var }} | 在模板字段中插入变量值。 | BashOperator( bash_command='echo "Date: {{ ds }}"') | 仅支持 Operator 中标记为 template_fields 的字段。 |
| 过滤器 | {{ var | upper }} | 对变量进行格式化处理。 | 'File_{{ ds_nodash | upper }}.csv' → File_20250930.CSV | 常用过滤器:upper, lower, replace, jsonify。 |
| 控制结构 | {% if condition %}...{% endif %} | 条件渲染模板内容。 | bash_command='{% if tomorrow_ds %}process {{ tomorrow_ds }}{% else %}exit 0{% endif %}' | 支持 if, for 等 Jinja 控制语句。 |
| 宏调用 | {% call macro() %} | 调用预定义宏。 | {% call macros.alert() %}Error occurred{% endcall %} | 需在 Airflow 配置中启用宏。 |
11.3 宏(Macros)与上下文变量
| 类别 | 变量名 | 说明 | 示例值 | 注意事项 |
|---|---|---|---|---|
| 常用上下文变量 | {{ ds }} | 执行日期(YYYY-MM-DD) | 2025-09-30 | 基于 execution_date,非当前时间。 |
{{ ds_nodash }} | 无连字符的执行日期 | 20250930 | 常用于文件命名。 | |
{{ ts }} | 执行时间戳(ISO 格式) | 2025-09-30T00:00:00+00:00 | 包含时区。 | |
{{ execution_date }} | 执行日期时间对象 | datetime.datetime | 可用于日期计算。 | |
{{ tomorrow_ds }}, {{ prev_ds }} | 明天/前一天日期 | 2025-10-01 | 基于 schedule_interval 推算。 | |
{{ dag_run.conf }} | 外部传入的配置 | {‘batch’: 100} | 由 TriggerDagRunOperator 或 API 传入。 | |
| 内置宏 | {{ macros.ds_add(ds, 7) }} | 日期加减(天) | 2025-10-07 | ds_add('2025-09-30', -1) 返回前一天。 |
{{ macros.uuid() }} | 生成 UUID | a1b2c3d4-… | 用于唯一标识。 | |
{{ macros.jsonify(dict) }} | 字典转 JSON 字符串 | ’{“a”:1}‘ | 用于 API 请求体。 |
11.4 SubDAG 与 TaskGroup 对比
| 特性 | SubDAG | TaskGroup | 推荐使用 |
|---|---|---|---|
| 架构 | 独立 DAG,通过 SubDagOperator 引入。 | 逻辑分组,非独立 DAG,仅 UI 折叠。 | ✅ TaskGroup |
| 执行方式 | 由 Worker 执行,占用一个任务槽。 | 内部任务由 Scheduler 直接调度,无额外开销。 | ✅ TaskGroup |
| 并行性 | 内部任务并行受限于 SubDAG 的并发配置。 | 与主 DAG 任务统一调度,并行性更好。 | ✅ TaskGroup |
| UI 展示 | 可点击进入查看内部任务。 | 在主 DAG 中折叠/展开显示。 | 各有优势 |
| 性能 | 有调度延迟,存在锁竞争(已弃用)。 | 高效,无额外调度开销。 | ✅ TaskGroup |
| 维护性 | 需维护多个 DAG 文件,耦合度高。 | 单文件内组织,易于维护。 | ✅ TaskGroup |
| 适用场景 | 已不推荐使用 | 模块化复杂流程、提升可读性 | ✅ TaskGroup |
| ⚠️ 结论:Airflow 官方已弃用 SubDAG,强烈推荐使用 TaskGroup 替代。 |
11.5 DAG 编写最佳实践(idempotency, logging, testing)
| 实践 | 说明 | 示例/建议 | 注意事项 |
|---|---|---|---|
| 幂等性(Idempotency) | 同一 execution_date 多次运行结果一致,避免重复写入。 | - 写数据库前先删除当日数据 - 使用 UPSERT 而非 INSERT - 文件写入包含 ds 路径 | 确保重试或手动重跑不会导致数据重复。 |
| 日志记录(Logging) | 在任务中添加详细日志,便于排查问题。 | def my_task(**context): logging.info(f"Processing for {context['ds']}") # ... | 使用标准 logging 模块;避免 print()。 |
| 错误处理 | 捕获异常并提供有意义的错误信息。 | try: result = api_call()except Exception as e: logging.error(f"API failed: {e}") raise | 确保异常被 raise,否则任务不会标记为失败。 |
| 单元测试 | 为 Python 函数编写测试,确保逻辑正确。 | 使用 pytest 测试 python_callable 函数。 | DAG 文件本身不需测试,但核心逻辑应可独立测试。 |
| 配置外化 | 使用 Variables 和 Connections 管理配置。 | 不要硬编码 URL、路径、密码。 | 提高可移植性和安全性。 |
| 依赖管理 | 明确声明 Python 依赖(requirements.txt 或 Docker)。 | 生产环境使用虚拟环境或容器。 | 避免版本冲突。 |
第十二章:监控、安全与生产部署
12.1 Web UI 监控 DagRun 与 TaskInstance
| 功能 | 说明 | 注意事项 |
|---|---|---|
| DAG 列表页 | 查看所有 DAG 状态、最近运行、调度延迟。 | 红色表示失败,绿色表示成功,灰色表示未运行。 |
| DAG 详情页 | 查看任务依赖图(Graph View)、时间线(Gantt)、日志、任务实例列表。 | Graph View 可直观查看执行流程和瓶颈。 |
| Gantt 图 | 展示任务执行的时间线,识别长任务和并行情况。 | 横条长度表示执行时长,间隙表示等待时间。 |
| Task Instance 日志 | 点击任务查看详细日志输出。 | 日志包含执行命令、返回码、错误堆栈。 |
| 清除(Clear) | 重置任务状态,触发重试。 | 会清除历史尝试记录;谨慎用于生产环境。 |
| 设置为成功/失败 | 手动修改任务状态,用于跳过故障任务。 | 仅用于紧急修复,避免数据不一致。 |
12.2 日志查看与远程日志存储(S3, GCS)
| 配置项 | 说明 | 示例值 | 注意事项 |
|---|---|---|---|
remote_logging | 是否启用远程日志存储。 | remote_logging = True | 默认 False,日志存本地。 |
remote_base_log_folder | 远程日志存储路径。 | s3://airflow-logs/{{ ti.dag_id }}/{{ ti.task_id }} | 支持 S3、GCS、Azure Blob 等。 |
remote_log_conn_id | 访问远程存储的 Connection ID。 | remote_log_conn_id = aws_logs | 需配置对应云存储的认证信息。 |
encrypt_s3_logs | 是否启用 S3 服务器端加密。 | encrypt_s3_logs = True | 增强安全性。 |
| 优势 | - 集中管理 - 持久化存储 - 与云架构集成 | —— | 避免本地磁盘满导致问题。 |
| 注意事项 | 确保 Airflow 服务有权限写入远程存储;配置生命周期策略自动清理旧日志。 | —— | 日志延迟可能影响实时排查。 |
12.3 用户权限管理(RBAC)
| 角色 | 权限说明 | 适用对象 |
|---|---|---|
| Admin | 全部权限,可管理用户、DAG、变量、连接等。 | 系统管理员 |
| User | 可查看 DAG、触发运行、查看日志,但不能修改。 | 普通开发/运维 |
| Op (Operator) | 在 User 基础上可清除任务、设置状态。 | 运维人员 |
| Viewer | 仅可查看,不能触发或修改。 | 审计、监控人员 |
| Public | 匿名用户权限(通常禁用)。 | —— |
| 自定义角色 | 可通过 UI 或 CLI 创建,精细控制权限。 | 满足特定安全需求 |
| 🔐 建议:生产环境启用 RBAC,遵循最小权限原则。 |
12.4 生产环境部署建议(高可用、备份、升级)
| 方面 | 建议 | 说明 |
|---|---|---|
| 高可用 | - 多节点部署 Scheduler(主备) - 多 Worker 节点(Celery/K8s) - 元数据库主从复制 | 避免单点故障。 |
| 备份 | - 定期备份元数据库(PostgreSQL/MySQL) - 备份 DAG 文件和配置 | 恢复时需同时恢复 DB 和文件。 |
| 升级 | - 先在测试环境验证 - 使用 airflow db upgrade 升级数据库 schema- 逐步滚动升级组件 | 遵循官方升级指南,避免跳版本升级。 |
| 监控 | - 监控 Scheduler 延迟(dag_processing.last_duration) - 监控任务失败率、队列长度 - 集成 Prometheus + Grafana | 及时发现性能瓶颈。 |
| 安全 | - 启用 HTTPS - 配置 fernet_key 加密敏感数据 - 限制网络访问 | 保护元数据和凭据。 |
| 资源隔离 | - 使用容器或虚拟机隔离组件 - 限制 Worker 资源使用 | 防止任务失控影响系统。 |
12.5 性能调优与元数据库维护
| 优化项 | 配置/操作 | 说明 | 注意事项 |
|---|---|---|---|
| DAG 解析 | parsing_processes = 4min_file_process_interval = 300 | 提高解析并发,减少扫描频率。 | 避免频繁解析大量 DAG 文件。 |
| Scheduler 性能 | scheduler_heartbeat_sec = 5max_tis_per_query = 512 | 调整心跳和查询批次。 | 根据负载调整,避免数据库压力过大。 |
| 数据库清理 | airflow db clean | 定期清理历史 DagRun、TaskInstance、Log 记录。 | 可设置 --clean-before-timestamp 删除旧数据;建议结合脚本定期执行。 |
| 索引优化 | 为 dag_run, task_instance, xcom 表添加索引 | 加速查询。 | 重点关注 execution_date, dag_id, state 字段。 |
| XCom 清理 | 避免存储大对象;设置 max_xcom_size 限制。 | 减少元数据库膨胀。 | 大数据应存外部存储,仅传递 URI。 |
| 连接池 | 配置 sql_alchemy_pool_size, sql_alchemy_max_overflow | 优化数据库连接复用。 | 避免连接耗尽。 |