如何快速上手 Apache Airflow 3:工作流编排、调度与监控指南
【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow
Apache Airflow 3 是 Apache 软件基金会旗下的工作流编排平台,让你用 Python 代码编写、定时调度和监控数据管道这类批处理工作流。如果你的团队有"先跑 A 再跑 B、失败了要重试、状态要一眼看清"这类需求,Airflow 是常见的选择之一。当前仓库中 3.x 系列的最新稳定补丁为 3.3.1(开发版 3.4.0),要求 Python 3.10–3.14。适合数据工程师、后端开发和运维人员阅读。本文先讲清它怎么工作,再给出跑通的最短路径、一个完整案例和上生产前的检查清单。
如何理解 Apache Airflow 3 的工作机制
先记住三个角色的分工:DAG 定义"做什么、什么顺序",调度器(Scheduler)负责"什么时候跑、按依赖顺序派发",任务实例(Task Instance)记录"每一次执行的真实状态"。
你的 Python 代码(DAG) │ 被扫描并解析 ▼ Scheduler ──派发──▶ Worker 执行任务 │ ▼ 元数据库(记录状态、日志、XCom) │ ▼ Web UI(http://localhost:8080)- 一个 DAG 就是一个 Python 文件,放在
AIRFLOW_HOME下的dags/目录里,Airflow 会周期性扫描并解析其中的任务与依赖。 - 每个任务每次运行都会生成一条任务实例记录,状态在"未开始 → 排队 → 运行中 → 成功/失败"之间流转,失败可配置自动重试。
- 任务之间用小量元数据传递用 XCom,官方建议任务保持幂等、不要把大数据集在任务间直接搬来搬去。
- 在本地开发时,
airflow standalone命令会把调度器、Web 服务、触发器等组件一并拉起,方便你理解各组件职责。
每个任务实例从创建到完成/失败的状态流转如下:
5 分钟跑通本地 Airflow 3 的最短路径
在 Python 3.10–3.14 环境下执行下面 5 条命令:
python -m venv airflow_env source airflow_env/bin/activate pip install apache-airflow==3.3.0 export AIRFLOW_HOME=~/airflow airflow standalone几点说明:
AIRFLOW_HOME是 Airflow 的主目录,配置、日志都落在里面,示例中设为~/airflow。airflow standalone首次启动时会在终端打印生成的用户名和密码,用它登录 Web UI。- 浏览器访问
http://localhost:8080即可看到 DAG 列表、触发运行、查看日志。 - 官方安装文档建议首次安装时配合 constraint 约束文件以保证依赖可复现,见 PyPI 安装说明 与仓库根目录的 INSTALLING.md。
启动成功后,在 DAG 列表页可以看到已解析的工作流并手动触发:
完整案例:用 TaskFlow 编写一个夜间日志清洗管道
下面是一个贴近实际的场景:每天凌晨 2 点,收集前一天的服务日志、清洗脱敏、汇总错误率。代码是仓库示例 DAG 同款的最小写法(airflow.sdk导入方式见 示例 DAG 目录):
from datetime import datetime from airflow.sdk import dag, task @dag(start_date=datetime(2024, 1, 1), schedule="0 2 * * *") def nightly_log_pipeline(): @task def collect_logs(): # 收集昨天的服务日志,返回文件清单(元数据,不是数据本身) return ["app-2024-05-01.log", "web-2024-05-01.log"] @task def clean_logs(files): # 按文件逐一清洗:去掉调试噪音、对敏感字段脱敏 return [f for f in files if f.startswith("app")] @task def report_error_rate(files): # 统计错误率并写入报表系统 print(f"待统计错误率的文件数:{len(files)}") files = collect_logs() clean = clean_logs(files) report_error_rate(clean) nightly_log_pipeline()关键约定:
schedule="0 2 * * *"是 cron 表达式,每天 02:00 触发一次;任务文件放到dags/目录即被自动发现。@task装饰的普通函数之间,用"函数调用"表达依赖:clean_logs(files)里传入的files会经 XCom 自动传给上游返回值。- 在 UI 里启用该 DAG 后,可以点 "Trigger Dag" 手动触发,也可以等调度时间自动运行;失败任务按配置自动重试。
仓库自带了逐步讲解的 TaskFlow 教程 和 完整教程索引,可以对照阅读。
Apache Airflow 3 生产环境部署检查清单
上生产前,逐项确认以下事项(依据仓库 README 与 Helm chart 文档):
- 元数据库:不要在生产使用 SQLite(官方明确不建议),改用 PostgreSQL(14–18)或 MySQL 8。
- 运行平台:只支持 POSIX 系统,参考镜像基于 Debian Bookworm;Windows 上开发请用 WSL2 或容器。
- 部署形态:中小规模可用 Docker 运行
apache/airflow官方镜像;Kubernetes 环境可用本仓库的 Helm chart(含 scheduler、workers、triggerer、api-server 等模板,见 chart 文档)。 - 组件拆分:生产环境通常把调度器、Web/API 服务、触发器等拆开独立部署,而非
standalone。 - 任务设计:保持幂等、大计算委托给 Spark 等外部引擎,Airflow 只做编排与状态记录;它不是流式系统,实时数据以拉取批次的方式处理。
- 可观测性:用 UI 的 Grid/Graph 视图跟踪运行,必要时接告警与日志收集,参考 管理文档目录。
下一步
- 读完 官方入门教程(fundamentals → taskflow → pipeline 三篇顺序即可上手)。
- 跑一遍 示例 DAG 目录 里的
example_simplest_dag.py,在 UI 里观察一次完整运行。 - 准备上生产时,先读 INSTALLING.md 与 Helm chart 生产指南。
【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考