Apache Airflow 名字里的 “Airflow” 是气流Logo 是几组彩色编织线Colors寓意明确大量任务可能彼此交叉、依赖、等待最终要像气流一样稳定有序地运转。如果你正在找一套能解决“多任务编排、定时调度、失败重试、状态统一可视”的生产级调度器Airflow 基本是绕不开的答案。这次我们直接看 Apache Airflow它是什么、怎么安装部署、怎么写第一个 DAG、怎么用 REST API 触发任务以及常见的坑和资源占用观察方法。文章最后会给出生产环境落地建议内容较长建议先收藏。1. 核心能力速览能力项说明项目类型开源工作流调度与编排平台开源来源Apache Software Foundation 顶级项目核心能力DAG 任务编排、定时调度、依赖管理、失败重试、告警、REST API、Web 可视化开发语言Python 为主前端为 React运行平台Linux 为主Windows 可以用于开发调试生产建议 Linux推荐部署方式Docker Compose 或独立 Python 环境 进程守护硬件要求学习环境建议 2 核 4G 起步生产按任务并发和调度频率评估调度方式Cron 表达式、固定间隔、事件触发传感器、手动触发、数据集触发是否支持 API支持Airflow 2.x 提供稳定 REST API 和 Python Client是否支持批量任务支持可一次回填多个日期、动态生成 DAG、并行执行任务适合场景ETL 数据管道、数据仓库同步、报表生成、算法模型训练调度、业务自动化从材料看Airflow 的关注点不在“单个任务执行多快”而在“很多任务之间如何可靠地编排”。它最适合的场景是任务之间有明确先后依赖又需要按时间周期反复运行。2. 适用场景与使用边界2.1 适合谁数据工程师负责数仓 ETL每天要从多个数据源抽取数据清洗后写入数仓。后端开发需要把多个服务接口串联成一条业务流水线比如定时拉取三方数据、调用内部接口、生成报表文件。算法工程师需要串联数据预处理、模型训练、模型评估、结果导出并且要记录每次运行的日志。运维/平台同学需要统一的作业调度入口而不是每个业务各自写一套 crontab。2.2 能解决什么问题任务依赖B 必须等 A 成功后执行C 需要 A、B 都成功再执行。失败重试某个接口临时超时可以自动重试 3 次。运行历史每次任务运行有日志、有状态、有耗时。补数据上周的 ETL 跑失败可以按日期回填。多人协作DAG 文件放到统一代码仓库后可以走代码评审、灰度发布。2.3 不适合什么场景对极致低延迟敏感的场景比如需要在毫秒级响应用户请求Airflow 更偏向离线调度不适合在线服务。超轻量单任务如果只有一个脚本每天跑一次用 crontab 成本更低Airflow 会显得重。复杂实时流计算实时流处理应该交给 Flink、Spark Streaming 等系统Airflow 一般负责周期调度这些任务而非做实时计算本身。2.4 合规边界使用 Airflow 编排任务时任务大概率会访问数据库、内网接口、对象存储等资源。需要注意对敏感数据做好脱敏日志避免输出明文密码、身份证号、手机号等。任务权限遵循最小化原则不要让调度账号拥有所有数据库写权限。需要调用第三方平台接口时确认平台授权范围。涉及人脸、声音、个人隐私等数据的处理任务必须确认数据来源合法和授权完整。3. 环境准备与前置条件在安装部署 Airflow 前先确认机器满足基础条件。3.1 操作系统与 PythonAirflow 核心是 Python 项目生产环境通常部署在 Linux 上。常见组合是Ubuntu 20.04 / 22.04 或 CentOS 7 / Rocky LinuxPython 3.8 到 3.12 之间具体版本需要参考对应 Airflow 版本的官方说明如果使用 Docker Compose 部署则本机不一定需要预装 Python较低版本的 Python 可能无法安装新版 Airflow建议先确认好版本兼容关系再执行安装命令。3.2 数据库选择Airflow 需要使用数据库保存 DAG 元数据、任务运行实例、日志索引等。用途推荐数据库快速学习/功能测试SQLite生产环境PostgreSQL 12 或 MySQL 8.x学习阶段使用 SQLite 最简单但 SQLite 不适合高并发调度生产环境必须切换到 PostgreSQL 或 MySQL。3.3 磁盘与内存Airflow 安装后Web Server、Scheduler、Worker 等组件都会占用内存。任务日志默认写在本地磁盘所以学习环境建议至少 2 核 CPU、4G 内存、20G 可用磁盘。生产环境按每日任务数量、并发上限、日志保留周期来评估通常建议 8G 内存起步。3.4 网络与端口Web 界面默认端口是 8080。如果机器上已有其他服务占用 8080需要修改默认配置或改用自定义端口。部署前先检查端口# Linux ss -lntp | grep 8080 # 如果端口被占用会看到对应进程信息4. 安装部署与启动方式Apache Airflow 提供多种启动方式下面按常见程度介绍三种pip 安装、Docker Compose、生产进程拆分。4.1 pip 安装方式这是最简单的本地起测方式。先创建虚拟环境避免和系统 Python 环境冲突。# 创建虚拟环境 python3 -m venv airflow_env # 激活虚拟环境 source airflow_env/bin/activate # 安装 Airflow版本号请按官方发布情况选择 pip install apache-airflow安装完成后需要初始化数据库# 初始化元数据库 airflow db init初始化完成后创建一个管理员账户airflow users create \ --username admin \ --firstname Admin \ --lastname User \ --role Admin \ --email adminexample.com \ --password admin接着启动 Web Server 和 Scheduler。# 启动 Web 服务默认端口 8080 airflow webserver --port 8080 # 新开一个终端启动调度器 airflow scheduler启动后浏览器访问http://127.0.0.1:8080登录账号就是刚才创建的 admin/admin。这样一步只能用于学习验证生产环境通常要拆成独立服务并用进程守护工具托管否则终端关闭服务就退出。4.2 Docker Compose 方式Docker Compose 是当前最常见也最省心的部署方式。Airflow 官方提供完整 Compose 配置包含 Postgres、Web Server、Scheduler、Worker、Triggerer 等组件。先提前说明Compose 配置内部需要初始化数据库、创建用户所以官网提供的完整 yaml 文件比这里展示的简化版更完整。更稳妥的做法是直接从官方文档获取标准配置然后修改少量参数。下面给一个可运行的简化模板适合理解组成结构version: 3.8 x-airflow-common: airflow-common image: apache/airflow:latest environment: AIRFLOW__CORE__EXECUTOR: LocalExecutor AIRFLOW__DATABASE__SQL_ALCHEMY_CONN: postgresqlpsycopg2://airflow:airflowpostgres:5432/airflow AIRFLOW__CORE__LOAD_EXAMPLES: False volumes: - ./dags:/opt/airflow/dags - ./logs:/opt/airflow/logs - ./plugins:/opt/airflow/plugins depends_on: - postgres services: postgres: image: postgres:13 environment: POSTGRES_USER: airflow POSTGRES_PASSWORD: airflow POSTGRES_DB: airflow restart: always airflow-webserver: : *airflow-common command: webserver ports: - 8080:8080 restart: always airflow-scheduler: : *airflow-common command: scheduler restart: always这里要特别说明使用apache/airflow:latest不便于版本锁定生产环境建议换成具体版本号例如apache/airflow:2.9.0风格。简化版没有包含数据库初始化、用户创建步骤首次启动需要进入容器执行airflow db init和airflow users create。如果觉得简化版启动麻烦可以直接去官方文档找完整 Compose 文件。完整启动流程参考如下# 创建目录结构 mkdir -p ./dags ./logs ./plugins ./config # 下载官方 Compose 文件按当前 Airflow 稳定版路径调整 curl -LfO https://airflow.apache.org/docs/apache-airflow/stable/docker-compose.yaml # 启动全部组件 docker compose up -d初始化数据库docker compose run airflow-worker airflow db init创建管理员用户docker compose run airflow-worker airflow users create \ --username admin \ --firstname Admin \ --lastname User \ --role Admin \ --email adminexample.com \ --password admin启动后依然通过http://127.0.0.1:8080访问。4.3 生产环境组件拆分生产环境一般把组件拆成独立进程进程职责webserver提供 Web UI 和 REST APIscheduler调度 DAG生成任务实例worker真正执行任务triggerer处理可延迟任务Deferrable OperatorsflowerCelery 模式下的任务监控面板常见执行器搭配LocalExecutor单机多进程并发适合中小规模。CeleryExecutor分布式执行适合任务量较大的生产集群。KubernetesExecutor每个任务动态拉起 Pod适合已经容器化的团队。4.4 启动后检查启动服务后通过下面几个方式验证基础环境# 检查 Web 端口 ss -lntp | grep 8080 # 检查进程 ps aux | grep -E airflow|gunicorn # 查看完整版本信息 airflow version如果能看到 Web 登录页说明部署已经成功。5. 功能测试与效果验证环境启动后最重要的是验证“DAG 能不能跑起来”。下面从零开始写一个最小 DAG 并触发。5.1 编写第一个 DAG在dags目录下新建文件csdn_demo_dag.pyfrom datetime import datetime, timedelta from airflow import DAG from airflow.operators.python import PythonOperator default_args { owner: csdn_writer, depends_on_past: False, start_date: datetime(2024, 1, 1), retries: 1, retry_delay: timedelta(minutes1), } def print_message(): print(Apache Airflow DAG run success) return done def print_done(): print(all tasks finished) return done with DAG( dag_idcsdn_demo_dag, default_argsdefault_args, schedule_intervaldaily, catchupFalse, tags[csdn], ) as dag: task_init PythonOperator( task_idprint_message, python_callableprint_message, ) task_done PythonOperator( task_idprint_done, python_callableprint_done, ) task_init task_done这段代码里dag_id是全站唯一标识。schedule_intervaldaily表示每天跑一次。catchupFalse表示不补跑历史日期。task_init task_done明确了执行顺序。5.2 让 Scheduler 识别 DAG将上面的 Python 文件放入 DAG 目录后Scheduler 默认每隔一段时间扫描一次。如果希望快速刷新可在 Web UI 里点击该 DAG 对应的刷新按钮或者直接等待自动扫描。在 Web UI 首页搜索csdn_demo_dag看到 DAG 名说明已被识别。5.3 手动触发 DAG点击 DAG 名称进入详情页右上角 “触发 DAG” 按钮即可手动触发。触发后重点观察观察项预期结果DAG 运行状态出现一次运行记录状态最终为成功Task 实例状态两个任务依次成功任务日志能看到 print 输出的文字耗时单条 Python 任务通常几秒内完成打开任务日志的路径DAG 详情页 - Graph View - 点击任务节点 - Log看到Apache Airflow DAG run success说明任务执行成功。5.4 测试计划时间对调度的影响Airflow 的调度逻辑中start_date、schedule_interval和logical_date是最容易绕晕的地方。简单理解Airflow 在一个调度周期结束后才会调度下一个周期。例如schedule_intervaldailystart_date为2024-01-01那么日期为2024-01-01的 DAG Run 会在2024-01-02 00:00之后被创建并调度。手动触发时可以在 Web UI 中指定运行日期。测试计划时间时推荐把start_date设置为过去时间catchupFalse观察首次调度是否正常。5.5 测试失败重试在任务函数中人为抛出一个异常观察 Airflow 是否按retries参数重试def print_error(): raise ValueError(simulate task failure)将任务替换成该函数后触发 DAG 运行预期效果任务先变为失败。到达重试间隔后自动重试。重试次数达到上限后任务状态为失败。通过这个测试能判断失败处理机制是否正常工作。6. 接口 API 与批量任务Airflow 2.x 的 REST API 是接入自动化平台的重要入口。6.1 API 基础信息REST API 默认和 Web Server 共用 8080 端口API 基础路径是/api/v1使用 API 前需要确认已创建用户并且该用户有对应权限。基础环境默认开启 Basic Auth。6.2 查询 DAG 列表curl -X GET http://127.0.0.1:8080/api/v1/dags \ -u admin:admin \ -H Content-Type: application/json返回 JSON 中包含 DAG 列表、是否暂停、文件路径等基本信息。6.3 暂停和恢复 DAGPATCH 接口可以修改 DAG 状态例如暂停某个 DAGcurl -X PATCH http://127.0.0.1:8080/api/v1/dags/csdn_demo_dag \ -u admin:admin \ -H Content-Type: application/json \ -d {is_paused: true}6.4 通过 API 触发 DAGcurl -X POST http://127.0.0.1:8080/api/v1/dags/csdn_demo_dag/dagRuns \ -u admin:admin \ -H Content-Type: application/json \ -d { dag_run_id: api_trigger_001, logical_date: 2024-01-01T00:00:00Z, conf: {} }触发成功后返回信息里包含dag_run_id后续可用它查询运行状态curl -X GET http://127.0.0.1:8080/api/v1/dags/csdn_demo_dag/dagRuns/api_trigger_001 \ -u admin:admin \ -H Content-Type: application/json结果中state字段会显示queued、running、success、failed等状态。6.5 Python Client 方式安装官方客户端库pip install apache-airflow-client然后通过 Python 调用import airflow_client.client from airflow_client.client.api import dag_run_api from airflow_client.client.model.dag_run import DAGRun configuration airflow_client.client.Configuration( hosthttp://127.0.0.1:8080/api/v1, usernameadmin, passwordadmin, ) api_client airflow_client.client.ApiClient(configuration) dag_run_client dag_run_api.DAGRunApi(api_client) dag_run DAGRun( dag_run_idpython_client_001, logical_date2024-01-01T00:00:00Z, ) api_response dag_run_client.post_dag_run( dag_idcsdn_demo_dag, dag_rundag_run, ) print(api_response)这里需要注意logical_date需要传入字符串具体字段格式以当前客户端版本为准。如果导入报错先检查包版本是否与服务端匹配。6.6 批量任务与回填Airflow 对“批量任务”的支持主要体现在几个方面第一回填。如果某个 DAG 已经配置了schedule_interval可以通过命令行一次生成多个日期的任务实例airflow dags backfill csdn_demo_dag \ --start-date 2024-01-01 \ --end-date 2024-01-07这条命令会生成 7 个日期的 DAG Run并且按日期顺序调度执行。第二动态 DAG。可以根据上游数据列表动态生成多个任务。例如遍历多个表名为每张表生成一个加载任务。TABLE_NAMES [table_a, table_b, table_c] with DAG( dag_iddynamic_etl_dag, ... ) as dag: for table in TABLE_NAMES: task PythonOperator( task_idfload_{table}, python_callableload_table, op_kwargs{table_name: table}, )第三分支任务。同一个 DAG 里可以并行执行多个互不依赖的任务只需要在依赖定义时不加连接即可。7. 资源占用与性能观察Airflow 的资源占用受执行器、任务量、日志保留策略影响很大不同环境差异明显。部署后可以从下面几个维度观察。7.1 进程内存# 查看 Airflow 相关进程的内存占用 ps aux --sort-%mem | grep -E airflow|gunicornWeb Server 基于 Gunicorn默认会启动多个 worker 进程。Scheduler 的内存占用和 DAG 文件数量、调度频率直接相关。如果 DAG 很多建议给 Scheduler 分配独立内存。7.2 数据库连接数Scheduler、Web Server、Worker 都会访问元数据库。使用 PostgreSQL 时可以查询数据库连接数select usename, count(*) from pg_stat_activity group by usename;如果连接数超过数据库上限需要调整 Airflow 配置中的连接池参数sql_alchemy_pool_enabledTrue sql_alchemy_pool_size5 sql_alchemy_max_overflow107.3 并发参数Airflow 有多个层级控制并发配置项作用parallelism全局最大并发任务数max_active_tasks_per_dag每个 DAG 同时活跃任务数max_active_runs_per_dag每个 DAG 同时运行的 DAG Run 数worker_concurrencyCelery Worker 并发任务数这些参数会影响任务调度速度和系统负载。参数值调得越高调度越激进对数据库、执行节点的压力也越大。7.4 如何降低资源占用将默认的load_examples设置为 False避免加载官方示例 DAG。设置日志清理机制定期清理旧日志文件。合理设置retries不要对所有任务设置过高重试次数。DAG 中避免大量高频率调度例如每分钟调度一次要非常谨慎评估。使用daily之类的大粒时间间隔除非业务真的需要分钟级调度。8. 常见问题与排查方法Airflow 部署和使用中经常遇到的问题这里整理成一张排查表。问题现象可能原因排查方式解决方案Web 页面打不开8080 端口被占用或 Web Server 未启动ss -lntp | grep 8080查看进程日志修改端口重启服务或换用 8081 端口DAG 在 UI 不显示DAG 文件不在dags_folder目录或文件语法错误查看 Scheduler 日志检查目录配置把文件放到正确目录修正语法后等待扫描DAG 显示但无法调度Scheduler 未运行ps aux | grep scheduler启动 Scheduler 进程任务一直处于 queued 状态并发上限已满或 Worker 未启动检查 parallelism、max_active_tasks 参数提升并发参数或扩展 Worker任务报错找不到模块Python 环境中缺少依赖包查看任务日志检查 worker 进程的 Python 环境在运行环境中安装对应依赖SQLite 数据库锁表SQLite 不支持高并发写入查看 Scheduler 日志中的 database locked切换到 PostgreSQL 或 MySQL时区不对系统时区与 Airflow 配置不一致airflow config get core default_timezone配置default_timezone并使用 Asia/Shanghai 等指定时区API 返回 401用户名密码错误或用户无权限检查认证方式测试登录重新创建用户或配置正确凭据API 返回 404DAG 不存在或 API 路径不对确认 DAG ID查看 API 路径使用GET /api/v1/dags确认路径修改 DAG 后没生效Scheduler 缓存或文件未被扫描等扫描周期查看任务日志在 UI 手动刷新 DAG或重启 Scheduler8.1 DAG 文件语法检查在部署 DAG 前先本地检查语法python -m py_compile dags/csdn_demo_dag.py如果语法有问题这里会直接报错。8.2 查看 Scheduler 日志Scheduler 日志是排查“为什么任务没跑”的关键# pip 方式安装时日志默认在 ~/airflow/logs/ tail -f ~/airflow/logs/scheduler/latest/scheduler.log # Docker Compose 方式 docker compose logs -f airflow-scheduler8.3 清理残留进程多次重启后可能残留旧进程导致端口冲突或任务重复执行pkill -f airflow scheduler pkill -f airflow webserver注意生产环境不要随意pkill建议用进程守护工具统一管理。9. 最佳实践与使用建议9.1 DAG 设计经验DAG 应该保持“短小清晰”。一个 DAG 只负责一条完整业务链路不要把几十个无关任务塞进同一个 DAG。任务层级不宜过深否则定位问题会很痛苦。变量分离不要将数据库连接串、接口地址、明文密钥直接写在 DAG 文件里。Airflow 提供 Connections 和 Variables 功能连接信息通过 UI 或环境变量维护export AIRFLOW_CONN_MY_POSTGRESpostgresql://user:passwordhost:5432/dbnameDAG 中通过Connection获取from airflow.hooks.base import BaseHook conn BaseHook.get_connection(my_postgres)9.2 幂等设计任务要支持重复执行而不产生脏数据。典型做法是写入前先清理目标分区。使用唯一主键重复写入覆盖。每次运行使用独立的运行 ID 或日期参数。Airflow 每次调度都会传入业务日期通过logical_date区分不同批次from airflow.decorators import task from airflow import DAG import pendulum with DAG( dag_ididempotent_dag, start_datependulum.datetime(2024, 1, 1, tzAsia/Shanghai), scheduledaily, catchupFalse, ) as dag: task def load_data(dsNone): # ds 是业务日期字符串如 2024-01-01 print(fprocessing data for {ds})上面的ds参数由 Airflow 自动注入代表 DAG 当前运行的计划日期。用这个日期来区分每天的数据批次能有效保证幂等。9.3 发布流程DAG 文件本质上就是代码建议走 CI/CD 流程分支开发合并前测试 DAG 语法。通过自动化任务扫描 DAG检查是否包含敏感信息。发布到目标环境时先暂停 DAG避免旧版本执行一半。发布后观察一段时间再恢复调度。9.4 告警配置任务失败后需要及时通知。Airflow 支持在 DAG 级别配置on_failure_callback常见做法是接入钉钉、邮件或消息队列。下面是一个最小失败回调示例def send_alert(context): dag_id context[dag].dag_id task_id context[task].task_id print(ftask failed, dag{dag_id}, task{task_id}) default_args { on_failure_callback: send_alert, }实际生产环境建议由告警平台统一接收不要每个 DAG 各自接一种通知方式。9.5 性能与成本控制设置合理的日志保留周期例如保留 30 天。对长时间运行的任务拆分成可重试的小步骤。使用传感器时要设置timeout避免无限制等待。如果同一周期内大量任务相互独立可以调整并行度提升吞吐。分布式场景下给不同任务配置不同队列避免所有任务挤在同一队列。10. 总结与下一步Airflow 最值得尝试的点不是“调度一个 Python 函数”而是“把复杂的任务依赖关系用 DAG 可视化地表达出来”。先验证的功能很简单写一个最小 DAG触发一次看任务日志再通过 REST API 触发一次确认接口可以通。最容易踩的坑集中在三个阶段安装阶段Python 版本和 Airflow 版本不匹配。启动阶段Scheduler 没启动任务一直不执行。使用阶段对start_date和logical_date的理解偏差导致调度时间不符合预期。从项目背景看Airflow 的 Logo 用彩色编织线表达任务交错依赖这个设计非常贴合它解决的问题。任务再多、依赖再复杂只要有清晰的 DAG、权重合理就能保持有序执行。后续可以继续扩展的方向包括接入 Celery 做分布式执行、使用 KubernetesExecutor 实现动态资源分配、利用数据集 Data-Aware Scheduling 做事件驱动调度、结合数据观测平台打通血缘和日志分析。如果你的团队已经有大量 cron 任务需要统一管理或者多条 ETL 流程经常互相等待Airflow 值得投入时间试跑。建议先在一台测试机上用 Docker Compose 部署跑通一个最小 DAG再逐步迁移真实任务。别急着把所有任务一次性搬过来先跑通一条链路再谈完整迁移。