ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

Apache Airflow核心原理与DAG实战:从入门到生产部署

Apache Airflow核心原理与DAG实战:从入门到生产部署 当业务里需要定时跑数据任务、串联复杂依赖、或者统一管理多个脚本调度时Airflow 是一个非常典型的解决方案。它自称 “A production scheduler, built with Colors”这句话并不夸张Airflow 不仅是一个能承载生产级调度压力的平台它的 Web UI 还通过颜色来直观表达任务状态绿色表示成功、红色表示失败、蓝色表示运行中打开界面一眼就能看清整个数据管线的健康度。本文会围绕 Apache Airflow 的核心原理、安装部署、DAG 编写、调度机制和生产环境落地展开。无论你是刚接触调度平台的新手还是准备把 Airflow 接入业务系统的后端工程师都可以跟着本文从头搭建一套可用的 Airflow 环境并理解它的设计思路。1. Airflow 是什么与核心概念1.1 调度器到底解决了什么问题在没有调度平台之前很多团队的定时任务是用 Linux 自带的 crontab 来做的。一旦脚本数量变多会出现几个问题任务之间的依赖关系难以表达例如“B 任务必须等 A 任务成功后执行”。看不到任务运行的历史记录失败原因靠猜。任务重试逻辑分散在业务代码里维护成本高。监控告警需要另外开发不同团队各写各的。任务运行资源无法统一管理高峰期容易互相影响。Airflow 这类工作流调度平台本质是把“任务编排”和“定时触发”从业务代码中抽离出来。开发者只需要用 Python 定义一张 DAG有向无环图描述任务之间的先后关系和依赖条件Airflow 负责按时触发、调度执行、记录日志、失败重试和展示状态。需要注意的是Airflow 是“调度平台”不是“计算引擎”。它本身不负责海量数据的计算而是负责把计算任务发送到合适的执行器上。数据计算仍然由 Spark、Flink、Hive 或者普通 Python 脚本完成Airflow 只关心流程编排和运行时机。1.2 Airflow 的核心组件一个完整的 Airflow 集群通常由以下部分组成组件作用Scheduler核心调度器负责解析 DAG根据调度周期生成 DAG Run并向执行器提交任务实例WebserverWeb UI 服务用来查看 DAG 状态、任务日志、触发手动运行、配置告警等Worker真正执行任务的进程在 CeleryExecutor 模式下由 Worker 拉起任务实例Metadata Database元数据库存储 DAG 定义、任务实例状态、运行历史、用户信息等通常是 PostgreSQL 或 MySQLExecutor任务执行器决定任务以什么方式运行常见的有 SequentialExecutor、LocalExecutor、CeleryExecutor、KubernetesExecutor在开发和测试环境最简单的方式是使用 LocalExecutorScheduler 在本地直接执行任务。在生产环境通常使用 CeleryExecutor 或 KubernetesExecutor 来横向扩展执行能力。1.3 DAG、Task、Operator 的关系初学者最容易混淆的是 DAG、Task、Operator 这三个概念。DAG有向无环图定义整个工作流的结构。它描述任务之间谁先谁后、谁能并行但 DAG 本身不执行任何业务逻辑。Task任务节点是 DAG 中的一个环节对应一次要执行的操作。例如一个“读取接口数据”的步骤是一个 Task一个“清洗数据”的步骤是另一个 Task。Operator任务的操作类型。Task 是节点Operator 是节点里具体要做什么事。Airflow 提供了很多内置 Operator例如 PythonOperator 执行 Python 函数、BashOperator 执行 Shell 命令、DummyOperator 表示空节点。从代码层面看一个 DAG 文件就是一份 Python 脚本。Airflow 的 Scheduler 会定时扫描 DAG 目录把文件解析成 DAG 对象然后按照调度规则来执行。2. 环境准备与安装部署2.1 安装环境说明Airflow 是基于 Python 的项目安装方式以 pip 为主。目前 Apache Airflow 的主线版本是 2.x 系列相比 1.x 在 API 和任务并发模型上有较大改进本文以 Airflow 2.x 为例。不同版本对 Python 版本有不同要求一般建议使用 Python 3.8 到 3.11 之间的版本。实际安装时需要根据你服务器的 Python 版本选择对应的约束文件。如果安装时遇到依赖冲突优先排查 Python 版本和依赖约束是否匹配。操作系统方面Linux 和 macOS 是常用环境Windows 下也可以运行但部分执行器和依赖在 Windows 上兼容性较差生产环境不建议使用 Windows 作为调度服务器。本文示例环境如下操作系统LinuxCentOS 7 或 Ubuntu 20.04 均可Python3.8Airflow2.x数据库默认使用 SQLite 初始化生产推荐 PostgreSQL2.2 安装 Apache Airflow安装前建议先创建虚拟环境避免 Airflow 的依赖和系统 Python 环境互相干扰。# 创建虚拟环境 python3 -m venv airflow_env # 激活虚拟环境 source airflow_env/bin/activate # 安装 apache-airflow这里不指定版本时安装最新版 pip install apache-airflow[celery,redis]2.8.0 \ --constraint https://raw.githubusercontent.com/apache/airflow/constraints-2.8.0/constraints-3.8.txt上面命令中的[celery,redis]是可选依赖表示同时安装 CeleryExecutor 所需的 Celery 和 Redis 客户端。如果只做单机测试可以不装这些扩展pip install apache-airflow安装完成后可以查看版本确认安装结果。airflow version2.3 初始化元数据库Airflow 运行时需要把 DAG 信息、任务运行记录、用户信息等写入元数据库。默认配置使用 SQLite适合初始化体验但生产环境需要切换到 PostgreSQL 或 MySQL。先设置 AIRFLOW_HOME 环境变量这是 Airflow 存放配置文件和日志的目录。export AIRFLOW_HOME~/airflow然后执行数据库初始化airflow db init这个命令会生成airflow.cfg配置文件。创建元数据库表结构。创建默认的 DAG 目录默认是~/airflow/dags。初始化完成后使用 SQLite 时可以在~/airflow目录下看到一个airflow.db文件。2.4 创建管理员用户如果要登录 Web UI需要先创建一个用户。Airflow 2.x 使用 Flask-AppBuilder 的用户体系可以通过命令行创建管理员账号。airflow users create \ --username admin \ --firstname Admin \ --lastname User \ --role Admin \ --email adminexample.com \ --password admin123创建成功后Web UI 登录时就可以使用这个账号。密码建议设置一个足够强度的值不要在测试环境之外使用弱密码。2.5 启动 Web 服务与调度器Airflow 需要同时启动两个进程Webserver 和 Scheduler。# 启动 Web UI默认端口 8080 airflow webserver --port 8080 # 新开一个终端启动调度器 airflow scheduler本地测试时如果配置的是 SequentialExecutor还需要手动启动一个airflow scheduler进程如果使用 LocalExecutorScheduler 会直接在本机执行任务不需要额外启动 Worker。访问http://localhost:8080用刚创建的用户登录就可以看到 Airflow 的 Web UI。启动后界面大多数区域是空的因为还没有编写任何 DAG。3. 编写第一个 DAG3.1 DAG 文件放置目录默认情况下Airflow 会扫描$AIRFLOW_HOME/dags目录下的 Python 文件。可以通过airflow.cfg修改这个目录。dags_folder /home/your_user/airflow/dags新建一个 Python 文件例如first_dag.py放在 DAG 目录下。Airflow 的 Scheduler 会定期扫描这个目录频率由dag_dir_list_interval配置项控制默认是 300 秒也就是说新 DAG 文件可能会延迟几分钟才会出现在 Web UI 中。3.2 写一个最简单的 DAG下面这个 DAG 实现了两个任务先执行一个 Python 函数再执行一条 Shell 命令。# 文件路径~/airflow/dags/first_dag.py from datetime import datetime, timedelta from airflow import DAG from airflow.operators.bash import BashOperator from airflow.operators.python import PythonOperator def say_hello(): print(Hello from Airflow DAG!) return done default_args { owner: admin, depends_on_past: False, start_date: datetime(2024, 1, 1), retries: 1, retry_delay: timedelta(minutes1), } with DAG( dag_idfirst_dag, default_argsdefault_args, descriptionA simple tutorial DAG, schedule_intervaldaily, catchupFalse, tags[example], ) as dag: hello_task PythonOperator( task_idsay_hello_task, python_callablesay_hello, ) echo_task BashOperator( task_idecho_time_task, bash_commandecho Executed at $(date), ) hello_task echo_task关键点说明dag_id是 DAG 在 Airflow 中的唯一标识不能重复。schedule_interval指定调度周期可以用 cron 表达式也可以使用daily、hourly等预设值。start_date是 DAG 的起始时间它影响调度器从何时开始生成运行实例。catchupFalse表示不补齐历史任务。如果设置为 True在 DAG 首次部署时会尝试把过去错过的调度周期全部执行一遍容易造成任务堆积。hello_task echo_task表示 hello_task 执行完成后再执行 echo_task。3.3 使用 PythonOperator 传递参数实际项目中Python 任务往往需要传入参数例如配置文件的日期、数据入库的表名等。PythonOperator 使用op_kwargs向 Python 函数传参。# 文件路径~/airflow/dags/param_dag.py from datetime import datetime from airflow import DAG from airflow.operators.python import PythonOperator def process_data(dt, table_name): print(fProcessing date: {dt}, table: {table_name}) # 这里可以调用业务处理逻辑 return f{table_name}_{dt} with DAG( dag_idparam_dag, start_datedatetime(2024, 1, 1), schedule_interval0 2 * * *, catchupFalse, ) as dag: process_task PythonOperator( task_idprocess_task, python_callableprocess_data, op_kwargs{ dt: {{ ds }}, table_name: user_log }, )这里的{{ ds }}是 Airflow 的模板变量会渲染成调度日期格式为 YYYY-MM-DD。Airflow 支持多种模板变量常用的包括模板变量含义{{ ds }}调度日期格式yyyy-mm-dd{{ ds_nodash }}调度日期格式yyyymmdd{{ execution_date }}执行时间在 2.x 中通常等同于调度周期开始时间{{ ts }}带时间和时区的时间戳3.4 任务依赖的多种表达方式Airflow 支持链式依赖、并行依赖、交叉依赖代码表达也很简洁。# 串行A - B - C task_a task_b task_c # 并行A 和 B 都完成后再执行 C [task_a, task_b] task_c # 条件分支A 执行完后B 和 C 可以并行执行 task_a [task_b, task_c]这种表达比写一大段判断逻辑更清晰也是 Airflow 的核心优势之一。4. 调度原理与任务状态4.1 DAG 调度周期逻辑刚接触 Airflow 时最容易困惑的是start_date、schedule_interval和execution_date三者的关系。Airflow 的调度行为是“一个周期结束后才生成对这个周期的执行任务”。例如schedule_interval0 2 * * *表示每天凌晨 2 点调度。start_date设置为2024-01-01那么第一个 DAG Run 的execution_date是2024-01-01 00:00:00但真正的执行时间却是2024-01-02 02:00:00。这种设计逻辑是在 2 点运行的任务通常处理的是前一天的数据。所以 Airflow 的任务实例里“调度日期”和“实际运行时间”是不一样的这在按天跑数、按小时跑数的离线数据场景中很常见。为了避免补数据时产生意料之外的运行堆叠推荐在新环境里把catchupFalse设为默认只有明确需要补历史数据时才开启。4.2 任务状态与 UI 颜色含义Airflow Web UI 最具识别度的特点就是用颜色来表达 DAG 和 Task 的运行状态。这也是标题里 “built with Colors” 最直观的体现。在 DAG 列表页和图视图中每个任务节点周围会显示不同颜色含义如下状态颜色说明success深绿色任务运行成功failed红色任务运行失败最终失败running蓝色或青色任务正在执行queued灰色任务已进入队列等待执行up_for_retry橙色或黄色任务失败但会重试up_for_reschedule浅蓝色任务被延迟调度传感器任务常用skipped浅灰色任务被跳过例如仅在特定条件下执行removed浅色运行后对应的 Task 定义被删除了scheduled深灰色任务已生成等待调度器分配upstream_failed暗红色上游任务失败导致当前任务未执行none白色无状态还没有进入相关执行状态这些颜色并不是单纯为了好看。在实际运维中打开 Grid 视图如果看到某一天整列出现大范围红色基本可以直接断定当天的数据链路有问题如果出现多个橙色说明存在失败重试需要检查上游临时抖动是否频繁。颜色让状态监控变成了一件非常直观的事情不需要逐个查询日志。4.3 重试、超时与失败处理Airflow 对任务失败的处理是可配置的。在default_args中常见的参数包括default_args { retries: 3, # 失败后最多重试 3 次 retry_delay: timedelta(minutes5), # 每次重试间隔 5 分钟 execution_timeout: timedelta(hours2), # 单次任务最长执行时间 }retries控制失败重试次数。retry_delay控制重试间隔。execution_timeout防止任务卡死超时后会强制标记为失败。重试机制虽然能缓解瞬时抖动但也会掩盖代码本身的缺陷。生产环境中建议区分“临时失败”和“确定性失败”对于网络超时、资源竞争等问题可以重试对于 SQL 语法错误、代码 Bug、参数错误等问题应该立即失败并告警而不是反复重试。5. 生产环境关键配置5.1 执行器选型Airflow 的 Executor 直接决定了任务在哪里执行。不同场景的选型建议如下执行器适用场景说明SequentialExecutor本地演示、单元测试串行执行任务不能并行不用于生产LocalExecutor单机多任务并行任务在 Scheduler 宿主机上以多进程方式运行适合中小规模CeleryExecutor分布式调度、多机执行需要 Redis/RabbitMQ 作为消息队列Worker 可以横向扩缩容KubernetesExecutor云原生环境、动态资源任务以 Pod 方式运行资源隔离强但部署复杂度高大多数中小团队从 LocalExecutor 起步足够。当任务量增长、单台机器负载过高时再迁移到 CeleryExecutor 是比较稳妥的路径。修改执行器的方式是在airflow.cfg中调整executor LocalExecutor修改配置后需要重启 Scheduler 进程才会生效。5.2 并发与资源控制Airflow 的并发参数较多常用的是下面几个# 每个 DAG 同时最多运行的 DAG Run 实例数 max_active_runs_per_dag 4 # Scheduler 并行处理的 DAG 数量 max_active_dag_runs_per_dag 4 # 全局任务实例并发上限 parallelism 16 # 单机最大可并行任务数 max_tasks_per_process 8参数设置过大会导致同一时间任务过多数据库和计算资源被冲垮设置过小则会出现任务长时间排队。建议在压测环境中观察任务实际耗时找到一个不过度浪费资源的平衡点。5.3 日志与监控Airflow 默认把日志写到本地文件系统。生产环境建议把日志集中存储例如使用 S3 或阿里云 OSS便于多节点统一查看和归档。在airflow.cfg中日志相关的配置包括# 基础日志目录 base_log_folder /home/your_user/airflow/logs # 是否将日志上传到远程存储 remote_logging False生产环境的监控告警推荐结合 Prometheus GrafanaAirflow 提供了 metrics 接口可以采集任务运行数量、调度延迟、任务耗时等指标。Web UI 本身适合查看单次任务详情不适合做长期趋势监控。5.4 元数据库选型Airflow 默认使用 SQLite但这只适合功能验证。生产环境推荐使用 PostgreSQL原因在于其并发性能更好、锁机制更成熟更适合调度器频繁写入任务状态。要切换元数据库先创建 PostgreSQL 数据库和用户然后修改airflow.cfg中的数据库连接串sql_alchemy_conn postgresqlpsycopg2://airflow:your_passwordlocalhost:5432/airflow修改完后重新执行airflow db upgrade来初始化表结构。生产环境一定要定期备份元数据库Airflow 的状态数据都存储在数据库中一旦丢失所有历史运行记录和 DAG 状态都会消失。6. 常见问题与排查思路6.1 问题速查表问题现象常见原因解决思路DAG 不出现在 Web UIDAG 文件路径不对或存在 Python 语法错误检查dags_folder指向执行python -m py_compile验证语法任务一直处于 queued 状态并发参数设置过小或 Worker 资源不足调大parallelism观察 Worker 日志Scheduler 启动失败元数据库连接异常检查sql_alchemy_conn确认数据库可连通任务运行后没有日志日志目录无写权限修改base_log_folder或提升目录权限时区显示不正确未配置default_timezone修改airflow.cfg中的default_timezone Asia/Shanghai修改 DAG 文件后 UI 不更新文件扫描有缓存等待dag_dir_list_interval到期或手动重启 Scheduler6.2 任务失败但原因不明确这是最常遇到的情况。建议按以下顺序排查打开 Web UI 对应的 DAG Run点击失败的任务节点。切换到 Log 页面查看任务输出的最后 50 行错误信息。区分是业务代码异常还是运维资源异常。检查是否有execution_timeout导致的超时退出。查看 Scheduler 日志确认任务是被正常提交还是被节点抢占。很多新手只看 DAG 变红就认为是代码问题实际上有一半以上的失败原因是资源不足、镜像拉取失败、数据库连接被拒绝等环境问题。6.3 大量任务同时触发导致数据库锁当一个 DAG 配置了catchupTrue且长时间没有运行重新部署后可能一下子生成大量 DAG Run导致元数据库写入压力激增甚至出现锁等待。解决思路# 通过命令行标记指定 DAG 的历史运行实例为失败 airflow dags reserialize更推荐的控制方式是把catchupFalse设为默认补数据时单独触发手动运行而不是批量自动补历史。7. 最佳实践与工程建议7.1 DAG 设计规范DAG 文件只负责流程编排不包含复杂业务逻辑。业务代码抽成独立的模块或包通过op_kwargs传入。任务单元尽量做到“一次执行、结果幂等”。同一个任务跑两遍效果和跑一遍相同。避免设计过于庞大的 DAG一个 DAG 控制 20 到 30 个任务以内比较合理。任务过多时可以拆分 DAG 或使用任务组TaskGroup。DAG 名称要能表达业务含义例如ods_user_log_sync、dwd_order_etl避免出现dag1、test2这种无意义命名。为 DAG 设置合理的execution_timeout防止异常任务拖住 Worker 资源。7.2 配置与安全管理数据库密码、API Token 等敏感信息不要直接写在 DAG 文件里。可以使用 Airflow 的 Variable 或 Connections并配合密钥管理服务。控制 Web UI 的用户权限避免所有用户都拥有 Admin 权限。生产环境建议按团队拆分 Role例如业务开发只拥有编辑 DAG 和查看状态的权限不赋予移除 DAG 的权利。定期更新 Airflow 版本。Apache Airflow 社区迭代很快旧版本可能存在安全漏洞和依赖兼容问题。升级前先在测试环境验证 DAG 兼容性。7.3 稳定性与容量规划给 Scheduler 的宿主机器预留充足的磁盘空间日志文件会持续增长建议配置日志轮转。监控 Scheduler 进程。Scheduler 如果挂掉定时任务不会按时触发而且 Web UI 上不一定有明显表现。数据库连接池参数需要根据任务量调整避免大量任务并发时数据库连接数不够。使用 CeleryExecutor 时Worker 节点需要保持与 Scheduler 相同的 DAG 代码版本否则会出现“UI 里能看见任务但 Worker 找不到对应代码”的问题。常见做法是使用镜像或制品库统一构建 DAG 包部署时所有节点同步更新。7.4 数据补偿与手动触发即使调度系统非常稳定上游数据源也可能出现问题。生产环境必须准备手动补偿数据的方案。在 Airflow 中可以通过 Web UI 的 “Trigger DAG w/ config” 功能手动触发一次 DAG Run同时传入 JSON 参数。例如传入需要重跑的日期范围由 DAG 内部的任务解析参数并处理指定数据。from airflow.models.param import Param with DAG( dag_iddata_reconcile_dag, start_datedatetime(2024, 1, 1), schedule_intervalNone, params{ business_date: Param(2024-01-01, typestring, descriptionBusiness date to process), }, ) as dag: # 任务内部通过 context 获取 params passschedule_intervalNone表示该 DAG 不会被自动触发只允许手动运行。这类 DAG 特别适合做数据补录和故障恢复。7.5 代码仓库与 CI/CDDAG 本质是代码必须纳入 Git 版本管理。推荐目录结构如下airflow-project/ ├── dags/ │ ├── ods/ │ │ └── sync_user_log.py │ └── dwd/ │ └── order_etl.py ├── plugins/ ├── requirements.txt └── README.md在 CI 阶段至少执行以下检查Python 语法检查。DAG 是否能被import可以使用python -c from airflow.models import DAG; from dags.xxx import dag简单验证。是否有未加execution_timeout的任务。是否包含明文密钥。8. 总结与学习路线Airflow 这个项目表面上是一个调度工具实际上一旦深入会涉及 Python 编程、数据库设计、消息队列、分布式任务执行、监控告警等多个领域。本文从安装部署起步介绍了 DAG 编写、调度逻辑、状态颜色含义和生产环境配置核心是为了让读者建立一条完整的技能链路从理解调度原理到能写出规范的 DAG再到能上线并维护一套可用的调度集群。如果你刚接触 Airflow下一步可以继续尝试使用 CeleryExecutor 搭建多节点分布式调度环境。用 Sensor 实现文件到达等待、数据分区就绪检测等功能。接入告警回调在任务失败时发送企业微信、钉钉或飞书通知。将 DAG 打包成镜像通过 KubernetesExecutor 运行。实际项目中最值得警惕的不是 DAG 写不出来而是任务运行后的“隐性失败”——任务退出码为 0但数据质量已经出现问题。因此在依赖 Airflow 的调度体系里一定要把数据校验任务也编排进 DAG并在关键节点配置告警。调度平台只是保证了“任务按时跑”真正保证数据可信的仍然是业务逻辑本身。希望这篇文章能帮你少走一些弯路也欢迎在实际部署中根据自己的环境灵活调整配置。
返回列表