ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Apache Airflow实战:从DAG定义到生产级任务编排

Apache Airflow实战:从DAG定义到生产级任务编排 如果你在工作中和“每天定时跑任务”这件事打过交道多半会产生一种复杂感受简单场景用 Cron 能应付稍微绕一点的需求靠脚本硬写也能扛可一旦任务数量上了两位数彼此还有先后依赖和数据传递失败要重试、漏跑了要补数据、半夜报警还要能查到当时发生了什么你就会认真考虑上一个正经的工作流调度系统。Airflow 就是我在处理这类问题里用得最顺手、踩坑也最多的一套工具。它本质是一套开源工作流平台由 Airflow 提出后捐给了 Apache 基金会核心是用代码定义“有向无环图DAG”把任务拆成节点、把依赖画成边交给调度器自动运行。它能做的远不止定时触发可以让几十个任务按顺序、按分支、按条件串起来单点失败自动重试历史区间一键回填每个执行实例的状态、日志、时长都留痕。适合谁用凡是需要跑数据管道、离线计算、报表生成、定时巡检、算法训练任务的人尤其是数据工程师和平台工程师基本都会在日常里和它打交道。刚开始接触 Airflow 的人往往先被一堆术语和概念搞晕DAG、Operator、Sensor、XCom、Executor……名字一大堆上手门槛看起来不低。但如果你只盯住一条主线——“它到底是怎么把一个任务跑起来的”很快就能理清脉络。下面我就按这个思路把 Airflow 的应用场景、核心机制和实操经验一次讲透。1. 先把 Airflow 放对位置它到底解决的是什么问题1.1 调度任务难的不是“定时”而是“编排”很多人一开始会问我有 Cron每分钟跑一次脚本不也能定时执行为什么还要上 Airflow答案是当你只需要“每隔多久跑一次这个命令”Cron 完全够用而且够轻。可真实业务里的任务是网状的不是线性的。举个最常见的场景每日报表要早上 8 点前产出。链路大概是凌晨 1 点抽数2 点清洗3 点完成汇总表最后 5 点渲染成报表并发送邮件。一旦抽数延迟后面的清洗和汇总都要顺延抽数失败之后的步骤做了也是白做。用 Cron 可以写五个任务但任务间的关系只能靠“时间猜”清洗必须比抽数晚但又不能太晚中间时间留多了浪费留少了容易踩空。更麻烦的是重跑某天源系统出了问题数据凌晨 3 点才补全后面四个任务得一个个手动触发顺序还不能乱。Airflow 把这种依赖关系变成了一张图。每个任务是一个节点谁先谁后、谁失败会影响谁画在 DAG 里一清二楚。调度器只在“上游成功”的前提下触发下游任务彻底告别时间错位带来的各种不确定性。漏跑了、失败了回填操作直接对整个历史区间重算不用再手动一个脚本一个脚本点。这才是它真正的价值所在不是在“定时”上做文章而是把“编排”这件事结构化、自动化、可视化。1.2 和 Cron、Jenkins、Temporal 放一起Airflow 的边界在哪选型的时候还得知道 Airflow 的边界。和这几类常见工具比一比会清晰很多工具擅长领域局限Cron单机定时任务简单直接无依赖管理无失败重试无可视化JenkinsCI/CD代码构建、发布流水线偏向部署场景任务编排能力弱进程模型较重Airflow数据管道、批量任务编排分支、重试、回放、日志完善不适合微秒级低延迟调度不适合流式计算Temporal / Prefect / Dagster长时运行、复杂状态流转新一代编排理念生态相对 Airflow 新既有组件和插件没这么全Airflow 最成熟的地方在于两点。一是生态Python 写的扩展插件、官方和第三方的 Operator 非常多数据库、云存储、消息队列、大数据组件基本都有现成封装。二是可观测性Web 界面能直接看到某个任务跑了多久、日志在哪、重试了几次、失败原因是什么排查问题比对着 Cron 日志手工翻效率高一个量级。反过来说Airflow 也确实不是一个适合所有场景的工具。它偏批处理调度频率太高会非常消耗资源通常分钟级以上才有意义它也不负责实时流计算你要做的是 Kafka 到 Flink 那一路不需要也没必要把 Airflow 塞进去。选不选 Airflow核心判断就是你的任务之间是否有明确的依赖关系以及你是否需要集中式的历史追溯和回放能力。如果答案都是“是”它就是对的工具。2. 核心概念拆解把 DAG 讲透你就掌握了 Airflow 的一半2.1 DAG 为什么是“有向无环图”Airflow 里的核心抽象就叫 DAG全称 Directed Acyclic Graph中文是“有向无环图”。拆开看三个词就能理解它的含义“有向”意味着每条边都有方向表示依赖方向“无环”意味着不允许闭环任务不会自己依赖自己形成死循环“图”意味着你的所有任务不是一长串命令而是一个结构化的网络。我在给团队成员讲的时候常用一个类比把 DAG 想象成一张施工流程图浇筑地基、砌墙、铺电路、装修每一步都有前置条件。你不能刚打完地基就刷漆更不能让装修环节倒过来决定是否打地基。每条边在施工图里就是一条“必须等上游完成”的规则。Airflow 做的就是把这张图变成可执行的东西图上每个节点对应一个实际任务调度器按照拓扑顺序一步步推进。在代码里DAG 的写法非常直观。一个最小示例长这样from airflow import DAG from airflow.operators.bash import BashOperator from datetime import datetime, timedelta default_args { owner: data_team, retries: 2, retry_delay: timedelta(minutes5), } with DAG( dag_idmy_first_dag, default_argsdefault_args, start_datedatetime(2024, 1, 1), schedule_interval0 2 * * *, catchupFalse, ) as dag: task_a BashOperator( task_idprint_start, bash_commandecho task start, ) task_b BashOperator( task_idprint_end, bash_commandecho task end, ) task_a task_b最后一行task_a task_b就是画边。和是 Airflow 提供的依赖运算符右侧的task_b依赖左侧的task_a。一个或多个任务之间可以这样连续连接比如task_a [task_b, task_c] task_d表达的是 task_b 和 task_c 都完成后 task_d 才开始。这种用代码“画图”的方式让整个工作流的逻辑能够被 Git 版本管理住谁改了什么依赖通过 diff 一眼看得出。2.2 Operator、Task 与 TaskInstance三层概念别混淆新手最容易绕晕的就是 Operator、Task、TaskInstance 这三个词。简单说它们的关系像是“类、实例、某次运行”的三层映射。Operator是执行单元的类型。比如执行 Shell 命令用 BashOperator执行 Python 函数用 PythonOperator查询 SQL 用 SQLExecuteQueryOperator等待某个文件出现用 FileSensor。它定义的是“这类任务怎么执行”。Task是 DAG 中的一个节点。你实例化了一个 Operator并给它起了 task_id它就成了一个 Task。Task 描述的是“在这个工作流里这一步在什么时候、依赖谁、失败怎么处理”。TaskInstance是 Task 在某次 DAG 运行中的具体实例。同一个 Task 今天运行和昨天运行各自拥有独立的 TaskInstance有独立的状态、开始时间、结束时间和日志。状态流转是排查问题的关键。TaskInstance 最常见的状态包括running、success、failed、skipped、upstream_failed、retry。其中upstream_failed是个很容易引起误会的状态它不等于本任务失败而是上游某个任务失败了自己压根没被允许执行。看到这个状态首先应该去翻上游的日志而不是盯在当前任务的日志里找原因。在实际操作中我最常提醒团队的一点是Operator 的重试参数只对任务本身的执行过程生效不会帮你自动解决上游数据缺失之类的问题。如果你在等待上游数据文件生成不要简单地在任务里加重试而是应该用 Sensor 去显式等待条件满足或者把重试时间拉长到超过数据延迟窗口。2.3 Scheduler 和 ExecutorAirflow 的双引擎Airflow 能跑起来靠两个核心进程协作。一个是Scheduler负责扫描 DAG 目录、解析文件、判断哪些任务该被触发、把任务放进队列另一个是Executor负责真正执行任务。两者各干各的职责分离是 Airflow 架构一个很重要的设计。你可以想象成 Scheduler 是大脑它思考“现在该干什么”而 Executor 是双手负责“把这些事干完”。两种角色解耦之后DAG 的调度频率、任务数量、并发量、执行机器都可以分别扩展。一开始跑在单机上两者可以共存一个进程任务量上来了可以把 Executor 拆出去变成多机 Celery 集群而 Scheduler 仍只需要狭义的调度判断。这种设计带来的最大好处是调度逻辑和任务执行隔离DAG 里有任何耗时任务都不会卡住调度器本身。换句话说你不用担心某个跑了 8 个小时的任务把后续任务的调度判断拖死因为判断和执行是两套体系。任务执行完状态回写元数据库Scheduler 再根据最新状态决定是否触发下游。2.4 时间、回溯与 catchup三个绕不开的调度细节Airflow 里围绕“时间”有不少容易搞错的概念。第一点是时区Airflow 默认使用 UTCstart_date和schedule_interval都基于 UTC 解释。国内团队如果直接写0 8 * * *看到任务每天早上 8 点UTC跑本地时间其实是下午 4 点。要么在配置里把default_timezone改成Asia/Shanghai要么在start_date里显式带时区否则时间一错所有调度的“肉眼预期”都会乱。第二点是 catchup 回填。start_date是一个历史时间点如果实际部署晚于这个时间并且 catchup 为 TrueAirflow 会把从 start_date 到当前时间之间每一个计划周期都补跑一遍。第一次部署 DAG 时这个行为很容易把人吓一跳因为一上线突然看到几十几百个 TaskInstance 同时处于排队状态。做法是默认把 catchup 设成 False只有明确需要跑历史数据时才手动触发回填。回填本身是个极好用的能力比如业务说“从 2024 年 3 月 1 日到今天报表全部重新计算”你只需要把日期范围填进回填参数Airflow 会按周期生成所有运行实例不用自己写循环去调接口。第三点是 depends_on_past 和 wait_for_downstream 这类约束。它们让你的任务在“上一次运行”结果未达到预期时不启动下一次。比如日报任务如果昨天的运行最终是失败状态你可能不希望今天的版本继续跑因为数据基础不完整。这两类参数在数据一致性要求比较高的管道里非常有用但也要小心一旦开启调度器在判断依赖时会更保守某个历史任务卡住后续所有周期都会跟着停住排障时需要把范围放大到整个历史链路上。3. 真实项目实操从零搭一个可以跑的 Airflow 管道3.1 场景设计一个标准的数仓 ETL 任务长什么样理论讲完我们来落地一个真实场景。假设你在给一家电商公司做销售数据分析每天需要把订单数据从业务库同步到数仓完成清洗和维度表关联再产出按商品、按区域聚合的销售报告最后发送给运营团队。整体任务链可以拆成四步从订单库抽取前一天的增量订单数据落地为本地文件或对象存储对原始数据做清洗去重、空值处理、类型转换将清洗后的数据写入数仓主表并关联商品维度和区域维度按汇总逻辑计算各指标生成报表文件并触发邮件通知。在这个链路里有个很典型的依赖细节第 2 步必须在第 1 步完全成功后才执行第 3 步要等第 2 步结果落地才有意义第 4 步依赖第 3 步。四条任务前一步失败后面不应该继续执行。这就是用 DAG 建模的基本盘。3.2 环境准备与基础配置实际部署 Airflow 并不复杂常见方式是 Docker Compose 或用官方 Helm Chart 部署到 Kubernetes。我先说本地快速验证的路径这适合刚开始接触的人。最简单的 Demo 环境是跑官方提供的 Docker Compose 文件。里面会启动 Scheduler、Webserver、Worker 和 PostgreSQL元数据库这几个服务。下载docker-compose.yaml后先建好日志和 DAG 挂载目录再把目录挂载给容器写 DAG 就不需要每次重新构建镜像编辑完刷新页面就能调度。一个最小命令是mkdir -p ./dags ./logs ./plugins echo -e AIRFLOW_UID$(id -u) .env docker compose up -d进入 Web 界面后在 Admin - Connections 里配置外部数据源连接例如配置订单库的 PostgreSQL 连接。在 Variables 里保存公共参数比如目标表名、业务日期偏移量、收件人邮箱。连接和变量都支持 Web 修改生产环境里也应尽量做到“程序代码不变、可调参数进配置”避免每次变更都要改 DAG 文件重新部署。3.3 手写第一个 DAG任务依赖与参数传递有了环境我们来写一个完整可运行的 DAG。这里我用 Airflow 2.x 推荐的 TaskFlow API它在任务间传参时比传统 XCom 方式简洁得多from airflow.decorators import dag, task from airflow.utils.dates import days_ago from datetime import datetime import json default_args { owner: data_team, retries: 2, retry_delay: 5, } dag( dag_idsales_report_daily, default_argsdefault_args, start_datedays_ago(1), schedule_interval0 3 * * *, catchupFalse, ) def sales_report_daily(): task def extract_orders(): # 模拟从订单库抽取增量数据 data {date: 2024-06-01, orders: [123, 456, 789]} return data task def clean_orders(payload: dict): # 清洗逻辑去重、过滤空值 orders [x for x in payload[orders] if x 0] return {date: payload[date], orders: orders} task def load_to_warehouse(payload: dict): # 写入数仓这里用打印代替真实写入 print(fLoading {len(payload[orders])} orders into warehouse) return len(payload[orders]) task def generate_report(count: int): # 汇总并生成报告 print(fReport count: {count}) extract_orders() clean_orders() load_to_warehouse() generate_report() dag sales_report_daily()TaskFlow API 比传统写法更贴近普通 Python 开发体验上游函数返回的值会自动作为下游函数的参数传入。你不需要手动使用xcom_push和xcom_pullAirflow 在背后完成 XCom 的传递。我在实际项目里更偏爱这种风格原因很简单可读性好业务逻辑抽出来是纯 Python 函数单元测试可以直接调用不必模拟 Airflow 环境。不过要注意一个“坑”TaskFlow 函数返回值如果体积很大不适合直接作为 XCom 传递。XCom 默认会存到元数据库里字段有大小限制大对象会撑爆数据库也会拖慢调度器。合理的做法是任务间传结果路径、任务 ID、或已写入表的主键范围真正的数据载体落到文件系统、对象存储或消息队列里。这个 DAG 上线后每天凌晨 3 点UTC会自动运行。你可以在 Web 界面看到一条从 extract 到 generate 的点状流点开每个任务能看日志、执行时间和重试历史。这就是最直观的“编排可视化”。3.4 上生产之前的部署形态选择本地 Demo 能用不代表直接照搬生产。生产环境里Airflow 部署形态要根据你的团队运维能力来选。如果你的团队熟练 Kubernetes用官方 Helm Chart 部署是最省心的路线。Scheduler、Worker、Webserver 各自独立 Deployment天然支持弹性伸缩Worker 可以根据任务量调整副本数。日志可以接 S3 或云存储避免节点磁盘膨胀。缺点是权限和网络策略要设计清楚元数据库要选托管实例稳定性会更高。如果团队没有 K8s但想保留分布式执行能力传统的 Celery 模式更合适。Scheduler Flower 多 Worker 部署在多台机器上扩展性也不错。Celery 模式需要维护队列组件故障排查会比 K8s 模式多一环但胜在架构直观、资料多。任务量并不大时单机 LocalExecutor 也足够应付几百个任务重点是别把 Scheduler 和执行任务塞在同一个容器里还要同时跑 Web尽量用 systemd 或容器编排分开进程保证一个进程挂掉不影响其他部分。部署形态本身没有绝对答案我的判断标准是核心业务数据管道优先选择托管 Kubernetes Worker稳定性最重要实验性任务和开发环境优先选择轻量部署成本和工作量都要最低。4. 我在生产环境里踩过的坑和爬出来的经验4.1 常见故障排查速查表这里把我在生产环境里遇到过的典型问题整理出来每一条都对应过真实的线上事故。表格形式更适合速查现象可能原因排查步骤DAG 一直不调度日志看不到任务执行Scheduler 未运行或 DAG 解析失败看 Scheduler 日志、看 DAG 列表是否有 parse error任务一直处于 queued 状态并发配置不足或 Worker 数量不够检查配置的 parallelism、dag_concurrency 和 Worker 数量任务成功但数据没写对代码里的时间参数用的是执行时间而非数据时间检查logical_date原 execution_date与data_interval_start的区别上游成功但下游报 upstream_failed上游实际产生了失败后又被重试成 func成功下游在失败窗口被推入失败态打开下游实例详情看“Upstream”页面里哪些任务状态不对调度延迟很大DAG 解析时间越来越长DAG 文件太多或太重import 了重量级第三方库拆分 DAG 目录、做 lazy import、压缩文件数量日志显示任务被 killed内存超过 Worker 限制或容器请求超限排查单任务内存占用调整 Worker 资源配置凌晨任务堆叠严重互相抢资源多个 DAG 的定时周期集中pool 规划不合理使用 Pool 隔离不同类型任务错峰排布这些坑里我最想强调的是时间参数。Airflow 2.x 里引入了data_interval_start和data_interval_end来更精确地表达某个运行实例对应的数据处理区间。以前很多人直接用execution_date甚至在 2.x 里发现该字段被废弃。实务上遇到每天一跑的日报任务运行时间是凌晨 4 点你真正关心的数据日期其实是“昨天”而不是“今天凌晨 4 点”。这一对对不上很容易导致报表少一天或者多一天尤其在跨年、跨月时更明显。建议统一封装一个from airflow.models import Variable from airflow.utils.context import Context def get_data_date(context: Context) - str: return context[data_interval_start].strftime(%Y-%m-%d)所有需要用到业务日期的函数都从 context 里取不要自己从系统时间推算。别看只是一个小封装在几十个任务里保持统一的语义能避免很多写数据时才发现日期错位的尴尬。4.2 并发配置、资源控制和稳定性优化Airflow 并发相关参数多很容易配乱。我按“从粗到细”的角度帮你理一遍parallelism整个 Airflow 集群最多同时运行的任务数是全局最大并发上限。dag_concurrency单个 DAG 内允许同时运行的任务实例数。注意这里说的是“同一个 DAG 里并发的实例数”不是“同时只能有一个运行实例”。max_active_runs_per_dag同一个 DAG 可以同时存在多少个运行实例。日报任务只允许一个还在跑防止上一个没结束下一个就叠上来这个参数要设置成 1 或较低值。worker_concurrency仅 Celery单个 Worker 进程里最多能同时执行多少个任务。Pool自定义资源池。比如把“写数据库的任务”单独放到一个 Pool限制为 2避免大量写库并发挤爆数据库连接。实际排查“任务排了一大堆但是不动”的情况优先检查顺序应该是全局parallelism→ 单 DAG 的max_active_runs_per_dag→ Pool 的容量 → Worker 数。按这个顺序走基本 10 分钟内能定位到瓶颈。除了并发稳定性还有两个让我印象深刻的点。第一是依赖库的版本隔离。多个 DAG 可能用到不同版本的包混在一个 Python 环境里总会出现“改了一个库另一个任务集体失败”。现在推荐的做法是使用 K8s 部署让不同 DAG 使用不同 Docker 镜像执行或者给同一批依赖相似的 DAG 建独立环境不要把所有东西堆在一个进程里。另一个是日志留存策略。生产环境任务日志增长非常快不设置轮转和归档几天就能占满磁盘。建议把日志目录切到对象存储并配置生命周期归档策略。关于任务幂等性我想单独强调几句。Airflow 的回填、重试机制让你很容易重复执行同一个数据周期。如果任务不是幂等的重复执行会写脏数据或者重复扣数据量。一个标准的做法是在写目标表时先按业务主键做删后插在任务开头加去重逻辑或至少依赖数据库的 unique 约束保证同一周期数据不会执行两次。没有幂等设计的数据管道用 Airflow 越用力反而越容易弄出脏数据。4.3 监控、告警与团队协作Airflow Web 界面自带的树形视图能看清运行状态但生产环境不能只依赖人肉盯页面。我通常至少补三套监控调度延迟监控记录“最后一个成功运行的 DAG 的logical_date是不是太久远”如果超过预设阈值说明链路卡住了。任务失败率监控实时统计最近一小时任务失败数量超过阈值后自动报警。注意排除掉还在重试中的实例否则重试策略会把你的告警轰炸到麻木。Scheduler 心跳监控Scheduler 进程是 Airflow 的中枢如果挂掉DAG 会全部休眠。用进程检查和元数据库里的最后一次心跳时间做兜底告警非常重要。告警通道方面我见过最实用的做法是 Slack 或企微机器人按 DAG 维度设置不同通知分组。可以把on_failure_callback统一封装成一个函数失败时自动带上 DAG 名、任务名、错误日志路径和下一次重试时间def alert_on_failure(context): dag_name context[dag].dag_id task_name context[task].task_id log_url context[task_instance].log_url notify(fDAG {dag_name} task {task_name} failed. Log: {log_url})封装好之后新增 DAG 时只要在 default_args 里带上这个回调函数所有任务就自动接入告警体系不用每个任务单独写冗余代码。团队协作方面我想说一个比较少被提及但很重要的习惯DAG 文件也要做 Code Review。很多团队把 DAG 当成“配置”而不是“代码”改起来很随意没人 review结果一个依赖关系的调整直接带崩一串任务。Airflow 的 DAG 本质是 Python 程序变更必须走 Git 流程至少做一次 review。DAG 文件的目录结构也应该分模块建议按业务线拆目录dags/ ecommerce/ order_report.py inventory_report.py finance/ settlement_daily.py common/ alert_helpers.py db_operators.py统一维护公共 Operator 和 Helper业务 DAG 只写业务流程不重复实现通用能力。这样团队新成员接手时看 DAG 就能快速理解业务链路不用钻进公共代码里猜前因后果。最后再分享一个我自己的实战习惯不要迷信并发调得越满越好。Airflow 这类调度系统通常“跑的顺”比“跑得满”重要得多。资源预留 20% 给重试、临时任务和调度本身的开销线上会更稳。对我来说衡量一个 Airflow 平台做得好不好不是看它能同时跑多少个任务而是看它是否做到了“该跑的任务不迟到、失败的任务能定位、历史的任务可回放”。把这三点守住这套调度系统就能扛得住业务很长时间的增长。
返回列表