
Airflow数据工作流编排实战从第一个DAG到生产环境一篇讲透【免费下载链接】airflow-guidesGuides and docs to help you get up and running with Apache Airflow.项目地址: https://gitcode.com/gh_mirrors/ai/airflow-guides每天早上 6 点上游数据表还没就绪你的处理脚本却因为调度时间到了而失败重试了三次中间某个任务挂了下游十几个任务全卡住而你还在手动盯着命令行刷日志。如果你也写过这种定时脚本crontab式的管道那么 Apache Airflow 就是为解决这类问题而生的数据工作流编排工具用 Python 代码把任务、依赖和调度写清楚它负责按计划执行、失败重试、并发控制和状态追踪。这篇文章按跑起来 → 写得对 → 排得动 → 撑得住的路线带你从零走完这套编排能力。五分钟搭出能跑的最小DAGAirflow 里一切工作流都叫 DAG有向无环图本质是一个 Python 文件。最小可用版本长这样from airflow import DAG from airflow.operators.empty import EmptyOperator from datetime import datetime, timedelta with DAG( dag_idmy_first_pipeline, start_datedatetime(2026, 8, 1), schedule0 2 * * *, catchupFalse, ) as dag: extract EmptyOperator(task_idextract) transform EmptyOperator(task_idtransform) load EmptyOperator(task_idload) extract transform load就是依赖箭头一行代码表达先提取、再转换、最后加载。把文件放进 Airflow 的 dags 目录默认每 5 分钟扫描一次打开 Web UI 就能看到它了。新手阶段只要记住 DAG 文件里的四个概念其他都是进阶内容详见 Airflow 核心组件 和 DAG 入门Operator一个任务节点比如PythonOperator跑一段函数EmptyOperator占个位HookOperator 连接外部系统数据库、API的底层胶水Sensor特殊的 Operator轮询等待某个条件成立文件出现了、分区写完了Connection / Variable集中存放账号密码和运行时参数别把密钥写进 DAG 文件⚠️ 坑点提醒start_date是必填项但它决定的是数据区间的起点不是现在跑。填得比当前时间远Airflow 会默认回补历史数据一次给你排出一百次运行。开发调试时记得配catchupFalse。DAG怎么写才对老手都在意的三件事DAG 是 100% 的代码写得好不好直接决定你凌晨三点要不要爬起来补数。三个最值得优先落实的原则完整版参考 DAG 最佳实践任务原子化。一个任务只做一件事。抽取转换入库塞进一个 task失败一次就得全部重跑拆成三个 task只需要重跑出问题的那一段。任务要幂等。同一个 DAG Run 重复执行结果应该和跑一次一样。用{{ ds_nodash }}这类模板变量确定处理日期而不是datetime.today()——后者在重跑历史日期时会算错是新手最常见的隐性 bug。增量处理。小时级管道每次只处理本小时的记录靠源系统的updated_at字段或自增 ID 过滤。单个区间失败不会污染全局重跑成本也最低。任务命名同样值得讲究fetch_orders_hourly比task_a在排查问题时友好得多。UI 的 Graph View 展示的是你写的依赖结构Grid View 展示的是每次运行的格子状态两个视图配合着看DAG 设计是否清晰一眼就能验证。任务没跑起来按这个顺序排查排障是新人上手 Airflow 最挫败的环节。记住这张排查路径表覆盖 90% 的任务不动了场景现象先查什么DAG 压根不在 UI 里文件是否在 dags 目录、是否过了 5 分钟扫描间隔、UI 是否有 Import Error多半是包没装或语法错误任务一直停在 scheduled对应 Pool 的槽位是不是满了、Executor 是否健康任务反复失败看 retries 配置是否在帮倒忙日志里第一个报错是什么上游失败下游全卡住检查失败策略考虑给失败任务配on_failure_callback而不是静默等待完整思路见 Debugging DAGs。两个高频细节值得单独说重试参数要会配。retriesretry_delay是default_args里最值得统一设置的字段网络抖动类错误多数靠它自愈。失败要能喊人。把on_failure_callback接到邮件或 Slack比事后翻日志高效得多参考 错误通知指南 和 日志体系。调度时间怎么设别只用 croncron 表达式schedule0 2 * * *能解决固定时点的问题但上游数据写完了再跑这种需求时间调度天生做不到。Airflow 给了三层递进的方案完整教程在 调度详解cron / timedelta固定周期最常用简单可靠。Timetable2.2用 Python 自定义任意调度逻辑cron 表达不了的时间规则如每月最后工作日都可以实现。Dataset2.4上游任务标记产出数据后自动触发下游 DAG实现基于数据可用性的编排不再猜数据几点就绪。另外DAG 之间还有一层显式依赖用 ExternalTaskSensor 或数据集把跨管道的先后顺序表达出来比把两条链路的 cron 时间错开半小时更稳。相关写法见 跨 DAG 依赖 和 数据集机制。并发与资源控制Pool 配置一分钟上手Airflow 默认会把所有任务丢进一个 128 槽的default_pool任务能并发就并发。大多数时候这是好事但当你有 50 个任务同时打同一个第三方 API 时下游接口先扛不住了。这时候就需要 Pool——给一批任务限流的令牌桶详见 Pools 指南call_api PythonOperator( task_idcall_api, python_callablefetch_orders, poolapi_pool, # 该池只开 3 个槽 priority_weight3, # 同池内排队时优先 )建池子在 UI 的Admin → Pools里加一条记录即可填名字、槽位数、描述。priority_weight决定同池任务谁先跑pool_slots能进一步控制单个任务占用的槽数。⚠️ 坑点提醒把任务指向一个不存在的 Pool任务会永远停在队列里UI 不报错也不提示。改 Pool 名称前先确认它真的存在这是血泪教训。生产环境怎么选 Executor任务调度只是半程谁来执行决定了系统能撑多大。Executor 就是干这个的四种常见选择对照如下原理拆解见 Executor 详解Executor适用阶段特点Sequential本地调试单进程串行零配置Local单机开发/测试单机多进程并行简单Celery生产主力分布式 Worker 消息队列弹性扩缩Kubernetes重资源/异构任务每个任务一个 Pod按需起停选型思路很简单开发机用 Local 就够任务量大、需要水平扩展上 Celery任务资源差异大有的要 GPU、有的跑几小时用 Kubernetes 更省。切换方式就是改airflow.cfg里的 Executor 配置UI 的Admin → Configurations可以随时核对当前生效值。配套的两个旋钮别忽略parallelism控制整个部署同时跑多少任务dag_concurrency限制单个 DAG 的并发两者配合 Pool 才能把资源用明白。动态任务任务数量不确定怎么办真实业务里任务数量经常是跑起来才知道的——比如要处理的 20 个城市列表、一批动态生成的表。Airflow 2.3 的动态任务映射expand/partial就是为此设计process_city process.partial(dagdag) # 固定参数 process_city.expand(citycities) # 每个城市生成一个并行任务上游任务返回列表下游按列表展开成一排并行任务跑完再汇给一个 reduce 任务就是完整的 MapReduce 模式。比在 DAG 解析期用 for 循环硬写 N 个任务优雅得多也省掉了代码改一处、发布一次的折腾深入用法看 动态任务指南 和 动态生成 DAG。上线前的三道质量关卡生产管道不能把跑完当成功还得回答数据对不对。三道关卡建议都配上管道内校验在 load 之后挂一个质量检查任务行数对不上、主键有重复就直接失败阻断下游。常用方案是 Great Expectations见 数据质量或 Soda见 Soda 指南纯 SQL 场景可参考 SQL 数据质量教程。重试与回调retries扛住偶发抖动on_failure_callback扛住必须有人知道的底线。可观测性UI 里盯 Graph/Grid View 之外的 Next Run 和日志配合 重跑机制 处理补数是日常运维的基本盘。下一步建议到这里从写第一个 DAG 到生产调度的链路已经完整了。建议的动手顺序先按第一节的模板本地跑通一个三任务 DAG → 用 Pool 和优先级给其中两个任务限流 → 把调度从 cron 换成 Dataset 触发。仓库里还有按主题整理的完整指南包括 Airflow 组件、UI 使用、连接管理、测试 DAG 等 50 多篇文档git clone https://gitcode.com/gh_mirrors/ai/airflow-guides把guides/目录当索引遇到具体问题时直接翻对应主题比从头背概念高效得多。 【免费下载链接】airflow-guidesGuides and docs to help you get up and running with Apache Airflow.项目地址: https://gitcode.com/gh_mirrors/ai/airflow-guides创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考