Article

任务调度 Airflow

更新于:2026-07-13

第一章: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_dateschedule_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-airflowconda 更好管理复杂依赖(如数据库驱动),适合数据科学环境。
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_folderDAG 文件存放路径,默认为 $AIRFLOW_HOME/dags确保路径可读,DAG 文件需在此目录或其子目录下。
core.load_examples是否加载示例 DAG,设为 False 可减少干扰。生产环境建议关闭。
core.executor执行器类型:SequentialExecutor, LocalExecutor, CeleryExecutor 等。单机测试可用 Local,生产建议 Celery 或 Kubernetes。
webserver.web_server_portWeb Server 监听端口,默认 8080。若端口被占用,可修改后重启 webserver。
scheduler.scheduler_zone调度器时区,默认为 UTC。建议保持 UTC,避免本地时区带来的混乱。
core.default_timezoneDAG 默认时区,影响 start_date 解析。可设为 Asia/Shanghai 等,但需与调度器时区协调。
logging.base_log_folder日志存储根目录。确保磁盘空间充足,定期清理旧日志。

2.4 使用环境变量配置 Airflow

环境变量对应配置项用途说明注意事项
AIRFLOW__CORE__EXECUTORcore.executor设置执行器类型,如 CeleryExecutor双下划线 __ 表示节与参数的分隔,大小写不敏感。
AIRFLOW__DATABASE__SQL_ALCHEMY_CONNcore.sql_alchemy_conn数据库连接字符串,如 postgresql://user:pass@host:5432/airflow敏感信息建议通过 secrets backend 管理。
AIRFLOW__WEBSERVER__WEB_SERVER_PORTwebserver.web_server_port修改 Web Server 端口。修改后需重启 webserver 生效。
AIRFLOW_HOME环境变量本身Airflow 主目录,默认为 ~/airflow可自定义路径,所有配置、日志、数据库默认在此目录下。
AIRFLOW__CORE__LOAD_EXAMPLEScore.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 DAG
from datetime import datetime

dag = 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_intervalschedule_interval='@daily'timedelta(hours=2) 或 cron 表达式定义 DAG 的调度频率,控制任务何时触发。from datetime import timedelta

dag = DAG(
dag_id='my_dag',
schedule_interval=timedelta(days=1),
start_date=datetime(2025, 1, 1)
)

# 或使用 cron
schedule_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 timedelta
schedule_interval=timedelta(minutes=15)
适用于简单周期任务,不支持复杂时间逻辑。
cron 表达式'分 时 日 月 周'精确控制调度时间,支持复杂周期。schedule_interval='30 8 * * 1-5' # 工作日 8:30注意时区问题,默认为 UTC;建议使用 time 模块或注释说明本地时间。

3.3 DAG 的 start_date、end_date 与 catchup 控制

参数语法用途代码示例注意事项
start_datestart_date=datetime(2025, 1, 1)定义 DAG 首次执行的逻辑时间点,用于计算调度实例。from datetime import datetime
start_date=datetime(2025, 9, 1, 0, 0)
必须设置;若使用 datetime.now() 可能导致调度异常;建议使用 pendulum 或固定时间。
end_dateend_date=datetime(2025, 12, 31)定义 DAG 最后一次执行的时间,超过此时间不再调度。end_date=datetime(2025, 12, 31)可选;常用于临时任务或项目周期限制。
catchupcatchup=Truecatchup=False控制是否补跑历史 DagRun(从 start_date 到当前时间)。catchup=False若设为 Truestart_date 较早,可能导致大量任务堆积;生产环境建议设为 False
catchup_by_defaultairflow.cfg 中设置全局默认 catchup 行为。catchup_by_default = False可通过配置统一管理,避免每个 DAG 重复设置。

3.4 DAG 的参数化配置(params, default_args)

参数语法用途代码示例注意事项
default_argsdefault_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 等;任务级设置优先级更高。
paramsparams={'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_dateDAG 实例的逻辑执行时间,非实际运行时间,用于确定数据处理的时间窗口。execution_date 是调度时间点,通常滞后于 start_date;例如 @dailyexecution_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_idtask_id 列表。
DummyOperator不执行任何操作,仅用于流程控制(如占位、合并分支)。轻量级,常用于逻辑分组。

4.2 BashOperator:执行 Shell 命令

参数语法用途代码示例注意事项
bash_commandbash_command='echo "Hello Airflow"'指定要执行的 Shell 命令或脚本。from airflow.operators.bash import BashOperator

task = BashOperator(
task_id='print_date',
bash_command='date',
dag=dag
)
支持多行命令(用分号或 && 连接);命令失败会触发重试。
envenv={'MY_VAR': 'value'}设置环境变量传递给 Shell 命令。env={'ENV': 'prod', 'PATH': '/usr/local/bin'}可覆盖默认环境;敏感信息建议使用 Variables。
output_encodingoutput_encoding='utf-8'指定命令输出的编码格式。output_encoding='gbk'处理中文输出时可能需要设置。

4.3 PythonOperator:调用 Python 函数

参数语法用途代码示例注意事项
python_callablepython_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_argsop_args=[arg1, arg2]传递位置参数给函数。op_args=['data.csv', 100]参数在函数调用时按顺序传入。
op_kwargsop_kwargs={'name': 'Alice'}传递关键字参数给函数。op_kwargs={'filename': 'report.txt'}常用于配置化参数。
provide_context已废弃(Airflow 2.0+)旧版本用于自动注入上下文。——新版本上下文自动提供,无需设置。

4.4 EmailOperator:发送邮件通知

参数语法用途代码示例注意事项
toto='user@example.com'['a@x.com','b@y.com']邮件接收者地址。to='team@company.com'支持单个或多个收件人。
subjectsubject='DAG Success'邮件主题。subject='ETL Job Completed on {{ ds }}'支持 Jinja 模板。
html_contenthtml_content='<h1>Success</h1>'邮件正文(HTML 格式)。html_content='Data processed for {{ ds }}'也支持纯文本。
conn_idconn_id='smtp_default'SMTP 连接 ID,需提前在 UI 中配置。conn_id='custom_smtp'默认使用 smtp_default,需确保连接有效。

4.5 SimpleHttpOperator:调用 HTTP 接口

参数语法用途代码示例注意事项
http_conn_idhttp_conn_id='my_api'HTTP 连接 ID,包含 URL 基地址、认证等。http_conn_id='jira_api'需在 UI 中预先配置 Connection。
endpointendpoint='/users'请求的端点路径(附加到 base URL 后)。endpoint='/tasks/create'不包含协议和主机。
methodmethod='GET''POST'HTTP 请求方法。method='POST'支持 GET, POST, PUT, DELETE 等。
datadata={'key': 'value'}POST 请求的表单数据。data={'name': 'test'}也可使用 json 参数发送 JSON。
jsonjson={'payload': 1}发送 JSON 格式请求体。json={'event': 'trigger'}自动设置 Content-Type 为 application/json。

4.6 BranchPythonOperator:条件分支控制

参数语法用途代码示例注意事项
python_callablepython_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_falsefollow_task_ids_if_false=['task_c']当无分支匹配时执行的任务(可选)。——通常不需要设置,未选中任务状态为 skipped。

4.7 DummyOperator:占位与流程控制

参数语法用途代码示例注意事项
task_idtask_id='start'唯一任务标识。from airflow.operators.dummy import DummyOperator

start = 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 DummyOperator

start = 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 >> endset_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 TaskGroup

with TaskGroup('data_processing') as tg:
step1 = BashOperator(...)
step2 = PythonOperator(...)
step1 >> step2
提升 DAG 可读性,避免 UI 过于复杂。
group_idgroup_id='preprocess'任务组的唯一标识。group_id='feature_engineering'必须唯一,不能与其他任务或组重名。
prefix_group_idprefix_group_id=False是否将 group_id 作为组内任务 task_id 的前缀。prefix_group_id=True # 默认值设为 False 可自定义 task_id。
tooltiptooltip='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 dag

from airflow.operators.subdag import SubDagOperator

subdag_op = SubDagOperator(
task_id='subdag_task',
subdag=subdag('parent', 'child', args)
)
SubDagOperator 已弃用,存在性能问题和锁竞争。
TriggerDagRunOperatorTriggerDagRunOperator(trigger_dag_id='target_dag')主动触发另一个 DAG 的运行,实现跨 DAG 编排。from airflow.operators.trigger_dagrun import TriggerDagRunOperator

trigger = TriggerDagRunOperator(
task_id='trigger_etl',
trigger_dag_id='daily_etl',
conf={'region': 'cn'} # 传递参数
)
推荐替代 SubDAG 的方式;支持传递 conf 参数。
confconf={'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 数据

方法语法用途代码示例注意事项
自动 pushreturn 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 关闭。
手动 pushtask_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_backendairflow.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:全局变量管理

概念说明注意事项
VariablesAirflow 中的全局键值存储,用于在 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 Variable

Variable.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使用方式代码示例注意事项
PostgresOperatorpostgres_conn_id='my_postgres'PostgresOperator(
task_id='run_query',
sql='SELECT * FROM users;',
postgres_conn_id='prod_db'
)
需确保 Connection 类型为 postgres。
MySqlOperatormysql_conn_id='my_mysql'MySqlOperator(
task_id='insert_data',
sql='INSERT INTO logs VALUES (...)',
mysql_conn_id='analytics_db'
)
类似 PostgresOperator,适用于 MySQL。
SimpleHttpOperatorhttp_conn_id='my_api'SimpleHttpOperator(
task_id='call_api',
endpoint='/data',
http_conn_id='jira_api'
)
Connection 类型应为 http,Host 字段填写基础 URL。
EmailOperatorconn_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 概念与作用

概念说明注意事项
HookAirflow 中封装与外部系统交互逻辑的类,提供统一接口访问数据库、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 BashHook

hook = 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 PostgresHook

def 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 MySqlHook

def 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 HttpHook

def 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_intervalstart_date 判断是否需要创建新的 DagRun 实例。创建时间基于 execution_date(逻辑时间),非当前时间。
TaskInstance 调度对于就绪的 DagRun,Scheduler 将符合条件的 TaskInstance 状态设为 scheduled,交由 Executor 执行。调度依赖任务依赖关系、资源限制、depends_on_past 等参数。
心跳机制Scheduler 定期更新自身状态(heartbeat),确保高可用模式下只有一个活跃实例。多 Scheduler 部署时依赖元数据库锁机制。
性能优化可通过 parsing_processesmin_file_process_interval 等配置优化解析性能。大量 DAG 文件时建议启用并行解析并延长扫描间隔。

9.2 常用调度参数(schedule_interval, interval, cron 表达式)

参数语法用途代码示例注意事项
schedule_intervalschedule_interval='@daily'timedelta(hours=2)'0 8 * * *'定义 DAG 的调度周期。# 每天 8:00
schedule_interval='0 8 * * *'

# 每 2 小时
schedule_interval=timedelta(hours=2)
None 表示仅手动触发;@once 表示运行一次。
内置别名@once, @hourly, @daily, @weekly, @monthly, @yearly简化常见调度周期。schedule_interval='@hourly'@daily 等价于 0 0 * * *(UTC 时间)。
timedeltatimedelta(minutes=30), timedelta(days=1)用于固定间隔调度。from datetime import timedelta
schedule_interval=timedelta(minutes=15)
不支持复杂时间逻辑(如工作日)。
cron 表达式'分 时 日 月 周'精确控制调度时间。'30 9 * * 1-5' # 工作日 9:30注意时区:Airflow 默认使用 UTC;可通过 timezone 参数调整。
Noneschedule_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 配置与分布式部署

配置项说明示例值注意事项
executorairflow.cfg 中设置执行器类型。executor = CeleryExecutor需重启 Scheduler 和 Worker 生效。
broker_url消息队列地址,用于任务分发。redis://localhost:6379/0amqp://user:pass@rabbitmq:5672//推荐使用 Redis 或 RabbitMQ;确保高可用。
result_backend任务结果存储地址,用于状态同步。db+postgresql://user:pass@db:5432/airflowredis://...必须与 broker_url 分开;建议使用数据库或 Redis。
worker 启动在各 Worker 节点运行 airflow celery workerairflow 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_imageDAG 文件是否已打包在镜像中。若为 True,则无需 git-sync;否则需配置 git_repo。
优势- 弹性伸缩
- 环境隔离
- 多语言支持
- 与 K8s 生态集成
适合云原生架构。
缺点- 冷启动延迟(Pod 创建时间)
- 成本较高(按 Pod 计费)
- 配置复杂
不适合高频短任务。

第十章:错误处理与重试机制

10.1 任务失败与重试参数(retries, retry_delay)

参数语法用途代码示例注意事项
retriesretries=3定义任务失败后的最大重试次数。PythonOperator(
task_id='process_data',
retries=2,
retry_delay=timedelta(minutes=5)
)
默认为 0;设置为 0 表示不重试。
retry_delayretry_delay=timedelta(minutes=5)两次重试之间的等待时间。from datetime import timedelta
retry_delay=timedelta(seconds=30)
支持 timedelta 或 Relativity(Airflow 2.0+)。
retry_exponential_backoffretry_exponential_backoff=True是否启用指数退避重试(1次: delay, 2次: 2×delay, 3次: 4×delay…)。retry_exponential_backoff=True可缓解服务压力,避免雪崩。
max_retry_delaymax_retry_delay=timedelta(hours=1)指数退避时的最大延迟时间。max_retry_delay=3600 # 秒防止延迟过长。

10.2 on_failure_callback 与全局回调

回调类型语法用途代码示例注意事项
on_failure_callbackon_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 DAGDAG 整体失败时触发(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 dag

dag_a = create_dag('etl_user', 'users')
globals()['etl_user'] = dag_a
避免重复代码;适用于多数据源、多租户场景。
从外部加载json.load(), YAML, Database从配置文件或数据库读取元数据生成 DAG。import json
with 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-07ds_add('2025-09-30', -1) 返回前一天。
{{ macros.uuid() }}生成 UUIDa1b2c3d4-…用于唯一标识。
{{ macros.jsonify(dict) }}字典转 JSON 字符串’{“a”:1}‘用于 API 请求体。

11.4 SubDAG 与 TaskGroup 对比

特性SubDAGTaskGroup推荐使用
架构独立 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 = 4
min_file_process_interval = 300
提高解析并发,减少扫描频率。避免频繁解析大量 DAG 文件。
Scheduler 性能scheduler_heartbeat_sec = 5
max_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优化数据库连接复用。避免连接耗尽。