引言“如果你的代码需要按顺序执行、需要定时触发、需要处理失败重试你就需要一个调度器。”这是「每日一个开源项目」系列的第 177 篇。今天的项目是Apache Airflow—— 用 Python 编写和调度工作流的平台数据工程领域最广泛使用的开源工具之一。在谈技术之前先回答一个问题Airflow 到底能做什么简单说任何需要按顺序执行多个步骤步骤之间有依赖关系需要定时或按条件触发的工作都可以用 Airflow 管理。最典型的是数据管道——每天凌晨从数据库拉数据、清洗、写入数据仓库、生成报表。但它不止于此训练机器学习模型、发送邮件报告、自动化基础设施操作都有真实的生产用例。46,400 颗 StarApache 2.0版本 3.3.0全球数千家公司在生产环境使用。你会学到什么Airflow 能解决哪四类核心问题重点DAG 是什么为什么用图来描述工作流Operator 和 Provider600 内置集成覆盖哪些系统Airflow 3.0 的关键新特性事件驱动调度和 Asset 感知快速上手用 10 行 Python 写一个真实 DAG什么时候该用 Airflow什么时候不该用前提知识会写基础 Python了解定时任务的概念cron 等对数据处理流程有基本认知不需要是专业数据工程师Airflow 能做什么核心问题这是本文最重要的部分。场景一ETL/ELT 数据管道最主要的用途。数据工程师每天面对的典型问题每天凌晨 2 点 1. 从 MySQL 业务数据库拉取昨天的订单数据 2. 清洗和转换去重、格式标准化、补充维度数据 3. 写入 BigQuery / Snowflake 数据仓库 4. 更新 BI 报表 5. 如果第 3 步失败发 Slack 告警并重试 这 5 步有严格的执行顺序步骤 4 必须在步骤 3 成功后才能运行。Airflow 把这个流程写成一个 DAG配置好依赖关系每天自动触发每步失败都有记录重试逻辑可配置所有历史执行都在 Web UI 里可查。真实规模AirbnbAirflow 的创始公司最初用它管理每天数百个 ETL 任务。现在有公司在生产环境跑几千个 DAG每天执行数十万个任务。场景二机器学习训练流水线MLOps 场景越来越普遍每周一 1. 从数据仓库拉取最新训练数据 2. 特征工程归一化、编码、拆分训练/测试集 3. 训练模型可以是 Python 脚本、也可以提交到 Spark 集群 4. 评估模型指标精确率、召回率、AUC 5. 如果指标超过阈值自动部署到生产环境 6. 如果低于阈值通知数据科学团队步骤 5 和 6 是条件分支——Airflow 的 BranchPythonOperator 支持基于上一步结果做不同处理。场景三定时报表和数据同步不需要是复杂的大数据场景每天早上 9 点把昨天的销售数据发给管理层邮件 Excel 附件每小时把 CRM 里的新客户数据同步到营销平台每周五把各部门的 KPI 汇总后写入 Google Sheets每月 1 号生成上月财务报表并上传 S3这些工作以前靠 cron shell 脚本处理Airflow 提供了可观察性每次执行成不成功、哪步慢了、失败后怎么处理一目了然。场景四基础设施自动化Airflow 不限于数据场景每天检测 S3 里超过 30 天的文件并归档压缩生产环境数据库的定时备份和备份验证自动化的云资源扩缩容高峰期扩、低谷期缩CI/CD 管道中的集成测试编排核心概念DAG为什么用图来描述工作流DAG Directed Acyclic Graph有向无环图。每个节点是一个任务Task节点之间的边表示必须先完成 A才能开始 B的依赖关系。无环确保不会出现死锁A 等 BB 等 A。一个典型的 DAG数据管道 extract_data ──→ transform_data ──→ load_to_warehouse ──→ send_report ↓ validate_schema ──→ quarantine_bad_dataAirflow 用 Python 定义这个图fromairflowimportDAGfromairflow.operators.pythonimportPythonOperatorfromairflow.operators.bashimportBashOperatorfromdatetimeimportdatetime,timedelta# DAG 定义withDAG(dag_iddaily_sales_pipeline,schedule0 2 * * *,# 每天凌晨 2 点start_datedatetime(2026,1,1),catchupFalse,default_args{retries:2,retry_delay:timedelta(minutes5),},)asdag:# 任务 1从数据库提取数据extractPythonOperator(task_idextract_data,python_callableextract_from_mysql,)# 任务 2数据转换和清洗transformPythonOperator(task_idtransform_data,python_callableclean_and_transform,)# 任务 3写入数据仓库loadPythonOperator(task_idload_to_warehouse,python_callableload_to_bigquery,)# 任务 4发送报告reportBashOperator(task_idsend_report,bash_commandpython send_email.py --date {{ ds }},)# 定义依赖关系执行顺序extracttransformloadreport这 40 行 Python 就是一个完整的生产数据管道。Airflow 负责在凌晨 2 点触发它按顺序执行每个任务失败时重试在 Web UI 里展示每次运行的状态。调度方式Airflow 支持多种触发方式# 定时调度标准 cron 表达式schedule0 9 * * 1-5# 工作日早上 9 点# Airflow 内置快捷方式scheduledaily# 每天schedulehourly# 每小时scheduleweekly# 每周# 3.0 新增Asset 事件触发数据驱动调度fromairflow.sdkimportAsset scheduleAsset(s3://my-bucket/raw-data/)# 当这个数据集更新时触发参数化和动态 DAG因为 DAG 是 Python 代码可以用 Python 的全部能力动态生成# 用循环为 10 个地区各生成一个相同结构的任务forregionin[us-east,eu-west,ap-south,...]:PythonOperator(task_idfprocess_{region},python_callableprocess_region,op_kwargs{region:region},)Operator 和 Provider600 内置集成Operator 是 Airflow 的任务模板——每种 Operator 封装了与特定系统交互的逻辑。内置 Operator无需额外安装Operator用途PythonOperator执行任意 Python 函数BashOperator执行 Shell 命令BranchPythonOperator根据条件选择执行分支EmailOperator发送邮件HttpOperator调用 HTTP APITriggerDagRunOperator触发另一个 DAGShortCircuitOperator条件不满足时跳过后续任务Provider 包按需安装Provider 是针对特定平台的 Operator 集合pip install apache-airflow-providers-XXX安装云服务aws— S3、Redshift、EMR、Lambda、Glue、SageMakergoogle— BigQuery、GCS、Dataflow、Vertex AI、Pub/Subazure— Blob Storage、Data Lake、Synapse、Azure ML数据库postgres、mysql、snowflake、databricks、spark消息队列apache-kafka、rabbitmq、redis其他工具slack、github、http、ssh、docker、kubernetes实际写一个从 S3 读数据写入 BigQuery 的任务fromairflow.providers.amazon.aws.operators.s3importS3FileTransformOperatorfromairflow.providers.google.cloud.transfers.s3_to_gcsimportS3ToGCSOperatorfromairflow.providers.google.cloud.operators.bigqueryimportBigQueryInsertJobOperator# S3 数据 → Google Cloud Storage → BigQuerys3_to_gcsS3ToGCSOperator(task_ids3_to_gcs,bucketmy-s3-bucket,prefixdata/2026-08-03/,dest_gcsgs://my-gcs-bucket/,)bq_loadBigQueryInsertJobOperator(task_idload_to_bq,configuration{load:{sourceUris:[gs://my-gcs-bucket/data/*],destinationTable:{projectId:my-project,datasetId:sales,tableId:daily_orders,},}},)s3_to_gcsbq_loadAirflow 3.0 的关键变化Airflow 3.0 于 2025 年发布两个最重要的新特性事件驱动调度Asset Watchers3.0 之前Airflow 主要靠 cron 定时触发——每天固定时间跑不管数据是否已经准备好。3.0 引入了Asset数据资产概念DAG 可以订阅一个数据资产当这个资产更新时自动触发而不是等固定时间。fromairflow.sdkimportAsset,DAG# 定义一个数据资产raw_ordersAsset(s3://data-lake/raw/orders/)# 这个 DAG 在 raw_orders 数据集更新时触发withDAG(dag_idprocess_orders,scheduleraw_orders,# 数据驱动不是时间驱动):...还支持Asset Watcher持续监听消息队列Kafka、SQS 等实现近实时的事件驱动fromairflow.providers.standard.asset.watchersimportKafkaAssetWatcher my_assetAsset(orders-stream,watchers[KafkaAssetWatcher(topicnew-orders,...)],)DAG 版本控制3.0 给每个 DAG 增加了版本号。修改 DAG 代码后旧的历史运行记录保留原版本快照新的运行使用新版本。这解决了一个长期痛点以前修改 DAG 会导致历史记录对不上。Web UI可视化监控Airflow 的 Web UI 是它的核心卖点之一。启动后在浏览器里能看到Grid 视图每个 DAG 每次运行的所有任务状态按时间从左到右排列绿色成功红色失败黄色运行中。一眼看出哪个时间段出了问题。Graph 视图DAG 的有向图节点颜色反映当前运行状态点击节点可以看日志、重新运行单个任务、查看执行时间。Assets 视图3.0 新增数据资产的依赖关系图展示哪些 DAG 生产数据、哪些 DAG 消费数据。快速开始安装最简单方式# 创建虚拟环境python-mvenv airflow-envsourceairflow-env/bin/activate# 安装 Airflow约束文件确保依赖版本兼容AIRFLOW_VERSION3.3.0PYTHON_VERSION3.12CONSTRAINT_URLhttps://raw.githubusercontent.com/apache/airflow/constraints-${AIRFLOW_VERSION}/constraints-${PYTHON_VERSION}.txtpipinstallapache-airflow${AIRFLOW_VERSION}--constraint${CONSTRAINT_URL}# 初始化数据库并启动开发用单机模式airflow standalone浏览器打开 http://localhost:8080用admin/admin登录。生产部署生产环境推荐用官方 Helm chart 部署到 Kuberneteshelm repoaddapache-airflow https://airflow.apache.org helminstallairflow apache-airflow/airflow\--namespaceairflow\--create-namespace或使用托管服务Astronomer商业托管、Amazon MWAA、Google Cloud Composer。什么时候该用什么时候不该用适合 Airflow 的场景批处理工作流步骤之间有明确依赖需要定时或按事件触发需要可观察性哪次失败了、哪步慢了、要能回溯历史多系统集成数据在 MySQL → Spark → S3 → Snowflake 之间流动团队协作多个数据工程师共同维护大量管道不适合 Airflow 的场景流处理需要毫秒级延迟的实时流处理用 Kafka Streams 或 Flink简单 cron 任务只是定时执行一个脚本crontab 就够了不需要引入 Airflow纯 API 服务Airflow 是调度器不是 Web 框架极短间隔触发每秒触发的任务不适合 Airflow它的调度粒度是分钟级项目地址与资源GitHub: apache/airflow官方文档: airflow.apache.org/docs官网: airflow.apache.orgPyPI: pypi.org/project/apache-airflowSlack 社区: s.apache.org/airflow-slack总结Airflow 解决的问题用一句话概括把一串有依赖关系的任务从手工维护的脚本变成可以定时触发、失败自动重试、执行历史可查、团队共同维护的工程化工作流。它的核心价值不是帮你执行任务而是管理任务之间的关系——谁先谁后、谁依赖谁、失败了怎么办、成功后通知谁、数据准备好了自动触发。Airflow 3.0 把触发方式从定时扩展到数据就绪让整个管道从时间驱动变成事件驱动。这对数据仓库场景的意义很大不再需要在凌晨 2 点固定等待上游数据而是上游数据一到下游任务立刻开始。46,000 颗 Star17 年的积累生产环境跑了十年的验证——如果你的工作涉及数据管道、定时任务或多步骤自动化Airflow 是目前社区最成熟的选择。探索 PrimeSkills —— 精选 AI Agent 与技能的市场每一个都经过真实企业工作流验证去掉浮夸留下真正有用的。欢迎访问我的个人主页发现更多有价值的见解和有趣的产品。