做技术选型这事,最怕的不是项目复杂,而是方案多到不知道该从哪下手。这些年我在不同公司、不同团队里,把Airflow、Prefect、Dagster、Temporal这几个长任务编排工具都拉上生产跑过,每次换工具都是因为上一套方案在某个关键点上确实撑不住了。今天就把我实际选型、落地、排障的过程和心得整理出来,给正在纠结的朋友当个参考。
我的核心感受先说在前面:这四个工具压根不是一个物种。Airflow是批处理调度器,Prefect和Dagster是数据编排平台,Temporal则是通用的持久化工作流引擎。你如果非要用一个工具去覆盖所有场景,那大概率两边都将就,最后都不顺手。
1. 选型之前,先看清这四个工具到底在解决什么问题
很多技术讨论一上来就比功能列表、比star数量、比社区热度,我觉得这是本末倒置。工具只是需求的投影,你心里的需求到底是什么形态,决定了哪个工具能对得上号。
1.1 四种工具的本质定位:批调度、数据编排、工作流引擎
我们先从核心抽象说起。
Airflow的核心抽象是DAG(有向无环图)。你定义一组任务,以及它们之间的依赖关系,调度器按时间触发DAG运行,把每一个TaskInstance调度给执行器去跑。注意,Airflow本身不执行任务,它只负责"决定何时触发哪个任务",真正跑Python函数、跑Spark作业、跑SQL的是Executor背后的worker、pod或云资源。所以它的本质是一个"批处理作业的调度和依赖管理器",优化目标是把一堆已经确定的批任务按时、按序地跑完。
Prefect的核心抽象是Flow和Task。Flow是你用Python代码写出来的完整流程,Flow内部可以包含条件分支、重试、动态任务生成。Prefect 2.0之后,这个概念大幅简化,去掉了原来1.x里容易让人迷惑的状态机、映射规则,改成"代码即流程",再用Deployment把Flow发布出去,由Prefect的调度服务按计划触发。它更强调"开发者写起来爽",所以上手体验比Airflow顺滑很多,原生支持动态、按运行时结果变化的流程结构。
Dagster的核心抽象是Asset(数据资产),强调软件定义资产(Software-Defined Asset)。你不再先定义任务、再定义依赖,而是直接定义"数据从哪来、经过什么产出哪张表/哪个文件"。Dagster会自动从资产函数之间的输入输出参数推断出一条物化管线,天然带有数据血缘、可观测性、资产级回放等能力。它解决的不仅是"任务怎么调度",更是"这份数据的生命周期如何被管理和追踪"。
Temporal的抽象则完全不同。它没有DAG这个概念,它的核心是Workflow和Activity。你直接用代码写业务流程,Workflow就是一段用Python、Go、Java等语言编写的工作流逻辑代码,中间可以调用Activity去执行具体的外部操作(比如调API、发消息、做模型推理)。Temporal的杀手锏是持久化执行:整个Workflow的执行状态会以事件流的方式持久化保存,任何一个worker宕机、网络分区、进程被杀,恢复后Workflow都能从最后一次成功的事件点继续跑,而不是从头再来。
所以你看,Airflow、Prefect、Dagster解决的是"数据批任务怎么按时按序跑完",而Temporal解决的是"一段长时间运行的业务流程怎么保证可靠地执行到底"。前者偏数据工程,后者偏分布式系统。
1.2 一个很实在的判断框架:从需求反推工具
我之前选型时用过一套框架,就是先问自己五个问题,回答完了基本就知道该选谁。
第一,任务形态是固定批处理还是动态业务流?如果每天凌晨2点跑同样的ETL,跑完之后生成报表,那Airflow或Dagster都是好选择。如果任务链路是"用户下单后触发一系列事件,每一步都不知道下一步要干什么,要等外部回调或者人工审批",那Airflow就明显不合适了,Temporal才是正经答案。
第二,你关注的是"调度稳定性"还是"运行时可靠性"?这句话很关键。我见过不少团队抱怨Airflow跑任务不及时、任务失败后要手动干预,但如果你的核心诉求只是"按计划触发+失败重试+告警",Airflow完全够用。可如果你需要"一个任务中间要等30分钟外部系统响应,响应回来继续往下执行",Airflow没有耐心等,也不该让它等,Temporal专治这种问题。
第三,团队的技术背景和运维能力如何?Airflow和Dagster要自己运维,调度器、数据库、执行器组件数量不少。Prefect有托管的Prefect Cloud,本地开源的server模式已经简化了运维。Temporal更是有名地难运维——它本身就是一个分布式系统,光server端就有前端、历史、匹配三大服务,还要配数据库和Elasticsearch(可选)。如果团队没有专职的SRE或平台工程师,Temporal自建会变成一场灾难。
第四,需要静态图还是要运行时动态?Airflow对动态图支持很差,DAG结构必须在解析期确定,想在任务A跑完之后根据结果动态生成任务B、C、D,这事在Airflow里要么用复杂的动态Task Mapping,要么干脆做不到。Prefect和Dagster对动态的支持好一些,而Temporal写起来就是普通代码,if-else、for循环、递归都可以随便用,天然就是动态的。
第五,可观测性要求到什么级别?是要"知道任务成功还是失败",还是要"知道每一步耗时、数据血缘、每次运行产出的资产清单"?Airflow的UI能看日志和实例状态,够用但朴素。Dagster的UI会直接展示资产依赖图和每次物化的事件时间线,体验完全上了一个档次。
这五问过完,我通常就能划出一个非常清晰的边界:纯批数据调度选Airflow,注重数据治理和血缘选Dagster,想快速上手且体验现代化选Prefect,但凡有长时间运行、可靠执行、步骤编排需求,直接上Temporal。
2. Airflow:数据批处理的"行业标准",但别让它超载
Airflow在我职业生涯里用得最多,也最熟悉它的脾气。虽然现在新的项目我经常推荐别的工具,但不可否认,很多公司线上批处理调度仍然是Airflow的天下,生态成熟程度是其他几个工具暂时追不上的。
2.1 Airflow的核心模型与适用场景
先快速过一遍Airflow的工作原理。你用Python代码声明一个DAG,描述这个DAG里有哪些Task,Task之间怎么连边,什么时候由Scheduler触发一次DagRun。Scheduler是一个常驻进程,默认每隔5到15秒轮询一次,检查有没有到等待时间的DAG,然后把对应的TaskInstance放进消息队列或者直接提交给Executor。Executor有多种,常见的是LocalExecutor(单机多进程)、CeleryExecutor(分布式队列)、KubernetesExecutor(每个任务一个Pod),实际生产里后两者居多。
Airflow最擅长的场景有三个。一是定时批量ETL,比如每天凌晨同步各业务库数据到数仓,再启动dbt、Spark、Flink等作业做数据加工,最后生成报表数据。二是跨系统的数据同步编排,比如定时从S3拉取日志,加载到ODPS或ClickHouse,再调用模型训练脚本。三是作为数据平台的任务总控,对接其他调度系统,统一做权限、日志和告警收敛。我待过一家电商公司,以前几十条离线数仓管道全挂在Airflow上,高峰期每天有上万次TaskInstance在跑,稳定性其实相当不错。
但是我要提醒,Airflow适合的是"流程预先可知"的工作负载。DAG是一张静态图,调度器每次DagRun都跑同一张图,只是参数不同。这个特性是Ansi性能好的原因,也是你后期想让它"智能一点"时的最大障碍。
2.2 生产环境里的Airflow常见瓶颈
Airflow最常见的坑不是调度器本身,而是"你想让它做它不该做的事情"。
第一个是把Airflow当实时任务平台。有些同事会写一个DAG每隔1分钟跑一次,或者写一个Long-Running Task挂着做实时监听。Airflow压根不是干这个的,它每次DagRun都要留下一堆元数据和日志,短周期DAG会让Scheduler和MetadataDB压力骤增,task的执行延迟也会被放大。实际上我把调度频率压到1分钟,Airflow元数据库的连接数和查询延迟就开始异常,最后不得不拆出去给别的系统做。
第二个坑是让Airflow执行真正"跑很久"的任务。Airflow有个合理的执行时间预期,绝大多数任务应该在分钟级跑完,亚小时级也还凑合,但如果你有一个任务需要等待外部系统半小时以上的回调,或者要连续运行好几个小时,Airflow的task会占住worker资源,而且中间任何一次机器重启、worker崩溃,整个task直接失败,没有恢复机制。这种场景你应该写成一个轮询外部状态的循环(短任务),而不是真让它长跑。
第三个坑是DAG解析慢导致调度延迟。Scheduler会周期性重新解析所有Python文件来发现DAG结构,如果你在上面放了很多复杂的模型定义、写了很多动态生成Task的逻辑,或者在import区做了吃IO的操作,解析时间就会暴涨。我自己遇到过,一个Airflow库里200多个DAG,因为某位同事写了非常重的import和文件扫描逻辑,导致半分钟以上才能解析完一轮,调度自然就变得越来越不准点,最后只能定期手工清库、优化解析代码。
第四是XCom传递大数据。Airflow的XCom默认把数据写进元数据库,如果你用它传输DataFrame或者大JSON,数据库很快会被撑爆,任务之间的耦合度也直线上升。正确做法是不要用XCom传大对象,把中间结果写到对象存储或临时表,下游重新读。
2.3 实操建议:什么时候继续用Airflow
如果你已经有Airflow在线上跑,而且任务形态是规律的、定时的、靠重跑和告警能兜底的,那就没必要为了所谓"先进"去迁移。我建议在以下几种情况,继续用Airflow是完全正确的。
第一个情况是团队已经积累了大量Airflow DAG,换工具的迁移成本远高于工具本身带来的收益。第二个情况是你需要对接非常成熟的生态,比如Airflow的Operator几乎涵盖了主流大数据组件——Hive、Spark、BigQuery、Snowflake、Redshift、Kubernetes,而且很多云厂商托管服务都兼容Airflow的API和UI,换工具反而增加了集成成本。第三个情况是你就想用一个稳定的、社区活跃的、不会被几个核心维护者绑架的框架,Airflow是很稳妥的选择。
不过,如果是全新项目,我会认真考虑是不是还要选它。Airflow的上手体验确实比Prefect、Dagster差一些,写复杂分支和动态任务时让人憋屈,加上UI的老旧感,很多年轻工程师更愿意用新鲜工具。而且Airflow迁移到Kubernetes之后,任务级Pod带来了资源隔离,但调度器本身还是单点(虽然可以高可用),大量Pod启动慢、镜像拉取慢的问题依然存在。这些事你要心里有数。
3. Prefect与Dagster:数据编排的新一代选择
这两个工具放一起对比很自然,因为它们都长着"下一代数编排"的面孔,都强调Python开发者体验,都更重视可观测性。但它们的理念差异也挺大,选之前建议想清楚自己是"想把流程跑好"还是"想把资产管好"。
3.1 Prefect的现代数据编排思路
Prefect 2.x之后,整个设计思路回归到一个核心信条:你写的Flow就是一段简单的Python代码,调度、重试、日志、缓存这些能力用装饰器和配置去叠加。我试过一次之后,最大的感触是"这不就是正常的Python开发体验吗"。
一个典型的Prefect Flow长这样:
from prefect import flow, task @task(retries=3, retry_delay_seconds=60) def fetch_data_from_api(url: str) -> dict: # 你的业务逻辑 return {"data": "whatever"} @flow(log_prints=True) def etl_flow(api_url: str): data = fetch_data_from_api(api_url) print(f"Got: {data}") if __name__ == "__main__": etf_flow.serve(name="daily-etl", cron="0 2 * * *")注意几个点:fetch_data_from_api是一个Task,etl_flow是Flow,两者都是普通Python函数;retries参数让任务失败后自动重试;最后的serve方法把Flow注册成一个常驻运行的部署,可以通过cron表达式定时触发。整个流程里没有Airflow那种"你在文件里声明DAG,然后等调度器扫描"的割裂感,代码即定义、即执行。
Prefect的另一个优点是Deployment模型比较灵活。你可以把Flow打包后部署到远端运行,也可以在每个工作机上直接serve,还可以通过UI手动触发一次运行。它支持Storage Block来管理代码和数据存储,支持Work Pool来管理动态worker,说白了就是给流动环境里的流程提供了一层"发布平台"的抽象。我自己用过Prefect Cloud的免费额度,体验确实顺滑,UI比Airflow好看得多。
但它也有短板。Prefect server自托管版本和Cloud版本之间存在一些能力差距,比如自动化规则(Automation)某些高级功能只在Cloud提供,开会后你得掂量清楚到底用哪个版本。另外它的大数据生态集成没有Airflow那么全,虽然市面上通常有现成的集成库,可长尾需求还是要自己写代码对接。
3.2 Dagster的数据资产视角
Dagster和Prefect走了不同的哲学路线。Dagster认为,数据平台的核心不是"任务",而是"数据资产"本身。它的写法是我先把"有哪些表、有哪些指标、这些资产怎么产生"定义清楚,然后框架自动推导出物化管线和计划。
一个简单的Dagster资产定义:
from dagster import asset, Definitions, ScheduleDefinition @asset def upstream_table() -> str: # 产出上游表 return "select * from source_data" @asset def final_report(upstream_table: str) -> str: # 依赖 upstream_table 产出最终报表 return f"create table report as {upstream_table}" daily_schedule = ScheduleDefinition(job=Definitions(...).get_job_def("__ASSET_JOB"), cron_schedule="0 3 * * *")注意final_report函数参数upstream_table直接指向了上游资产函数名,Dagster会通过这个函数签名自动构建依赖图。你定义的每个函数就是一个资产,每次运行会记录AssetMaterialization事件,UI上能看到哪张表在什么时间点被哪个代码段产出,血缘关系一目了然。
这种抽象带来的好处是在做数据治理、表依赖分析、重跑影响范围排查时非常舒服。你不需要翻DAG源代码去猜"这张表到底被谁改了",直接在UI上输入资产名就能看到上下游。Dagster还内置了类型系统、资源系统和配置系统,可以复用连接器配置、切分环境,写起来比Airflow优雅很多。
不过要注意,Dagster的这套模型是有学习门槛的。它不像Prefect那样"写了一堆装饰器就能跑",你需要先理解Asset、Op、Job、Schedule、Sensor、Code Location这一整套概念。我第一次上手时,光搞明白"代码库(Repository)"和"部署位置(Code Location)"的区别就花了两天。而且Dagster更偏数据领域,你要是想编排的是微服务业务流,用它也很别扭——它没有Temporal那种持久重放机制,本质上仍是一个"调度+执行数据任务"的平台。
3.3 两者对比与选型建议
我直接给结论:如果你的团队主要写Python,希望保留普通开发的流畅体验,又不想被Airflow的DAG模型束缚,Prefect是最容易落地的。它的文档干净,示例多,排错直观,新人基本一天就能上手。
如果你的数据团队越来越大,表之间依赖混乱,经常需要回答"这张表哪来的、被谁改了、重跑会影响到谁"这类问题,Dagster的资产中心模型会帮你省掉很多沟通成本。它适合体量中等以上、有数据治理诉求的团队。
实时情况是,Prefect和Dagster都还在快速迭代,生产环境普及率不如Airflow高,这带来一个潜在风险:出了问题,可参考的社区案例和第三方插件比Airflow少。你在白嫖它们新特性带来的爽感的同时,也得准备好自己挖坑填坑。我之前在公司引入Dagster时,就遇到过一个老版本的调度问题,最后是去翻GitHub issue才找到规避方案,确实没有Airflow那种"和社区一起成长"的厚实感。
4. Temporal:被低估的长任务工作流引擎
很多做数据工程的人看到Temporal会有点懵:它既不是DAG调度器,也不是数据血缘平台,那它到底算什么?我一开始也懵,直到被一个在线业务团队的需求逼着去研究,才真正搞懂它真正的价值。
4.1 Temporal能做什么:从ETL到业务长事务
Temporal的定位是"持久化执行的工作流引擎",核心解决的是分布式应用中长流程的可靠性问题。它不关心你是不是在跑数据任务,它关心的是一段可能持续几小时、几天甚至跨年的业务流程,如何确保不会因为进程崩溃、机器宕机、网络抖动而中断。
举个很典型的业务例子:用户发起退款申请,流程要经历风控审核、优惠券回收、账务扣减、通知用户、日志归档等多个步骤,其中风控审核可能等待外部系统回调,账务扣减要调支付网关,通知要发消息队列。这种流程如果用传统代码硬写,每一步都要自己处理重试、状态保存、失败补偿,代码会越来越乱,而且一旦进程挂掉,整个流程的状态就丢了。Temporal把这些脏活全接管了。
它的工作方式通俗解释就是:你用普通代码编写Workflow函数,Temporal的client SDK负责把Workflow的每一次函数调用、每一个变量变化都以事件的形式写入Temporal Server的历史存储。执行过程中,这台worker挂了也没关系,新worker接手时会从历史事件中"重放"这段代码,把状态恢复到挂掉之前的点,然后接着往下走。
注意"重放"这个词,它是Temporal的核心机制,代价也很高,要求你的Workflow代码必须是确定性的,不能依赖本地时间、随机数、外部请求等不确定来源——与外部系统的交互一律封装在Activity里。这也是很多人刚上手时最不习惯的地方。
4.2 核心概念与可靠性原理
Temporal的关键概念有四个。
Namespace是租户隔离,你在里面定义工作流队列和配置,生产环境一般按团队或业务线单独建Namespace,便于做权限和配额管理。
Workflow是一个用编程语言编写的持久化可重放函数。Workflow只能做计算和决策,不能直接调外部API,也不能做IO。它的运行单元是"事件",每一步执行都会被记录。
Activity是实际干活的单元,它可以读写数据库、调API、发消息,可以失败和重试,Temporal会按你配置的retry policy自动重试Activity。Activity执行完把结果返回给Workflow。
Worker是一个常驻进程,它通过SDK轮询Temporal Server,领取Workflow任务和Activity任务来执行。你的业务代码跑在Worker里,通常一个服务启动时注册多个Workflow类型和Activity类型。
还有一个常被忽略的机制是Signal和Query。Signal允许外部系统向正在运行的Workflow发送消息,比如"用户点击了取消,请停止当前步骤";Query允许查询Workflow当前的状态,而不用修改流程逻辑。这两个能力让Temporal可以应对实时、交互式的长流程,这在Airflow那个模型里根本没法想象。
序列化和事件历史:Temporal Server会把Workflow执行产生的所有事件存到数据库里,一个Workflow最多能积累多少事件取决于你的配置和存储能力。短期无感,但长跑型Workflow以及重放频率很高的话,事件膨胀会变成性能瓶颈。我做过一个长时间运行、每步都发Signal的Workflow,跑到后面明显感觉Event History变长,Worker重放变慢。
4.3 适用边界与坑
Temporal绝不是一个拿来平替Airflow的工具。我用下来,它有自己明确的适用边界。
第一,如果你只是需要每天跑一组固定的SQL和Python脚本,Temporal不是最佳选择。它不是为批处理调优的,每次触发一个Workflow执行、记录事件、持久化状态,相比Airflow直接submit一个task,成本高不少,而且它的UI远远不是"调度监控"的形态,不适合当数仓任务总控面板。
第二,Temporal的运维成本被很多人低估了。自建的话,你要部署Temporal Server的前端、历史、匹配等组件,还要搭配数据库和可选的Elasticsearch,更别提高可用、扩容、监控告警。别指望装个docker-compose就能上生产。我认识一个朋友所在的公司,初始demo跑得很欢,后来要搞生产就卡在基础设施上,最后还是选了托管服务。所以如果你没有足够的分布式系统运维经验,建议优先考虑Temporal Cloud或按需自建。
第三,确定性约束是个隐形地雷。你写Workflow代码时必须遵守"代码即数据"的规则——不能直接用time.Now()、不能读环境变量、不能调随机函数这些不确定的函数,否则重放时会出现不一致,导致流程状态错乱。这些细节SDK的静态检查不一定能全部拦下来,真正能依赖的是你自己对Temporal运行原理的理解。好在官方示例和文档已经把常用模式都覆盖了,照着左边喝酒右边砸键盘的新手期不会太长。
第四,Temporal非常擅长"业务编排",这一点在当前微服务和跨团队协作场景下非常值钱。我之前负责的一个系统里,订单生命周期、风控流程、资源审批这些长事务,全用Temporal梳理成一个个Workflow,每个团队用自己熟悉的语言写Activity,由Temporal统一编排,代码简洁度和线上稳定性都明显提升。要是你还在一遍遍手写状态表、循环重试、分布式锁,我真建议你看看Temporal。
5. 生产环境选型决策表与混合架构实践
前面讲了各自的优劣势,但真正到选型落地的时候,一张能直接对照的速查表往往比长篇大论更有用。下面这张表就是我自己做技术方案时的常规参考,供你拿过去按需调整。
5.1 一张表看清四个工具的差异
| 维度 | Airflow | Prefect | Dagster | Temporal |
|---|---|---|---|---|
| 核心抽象 | DAG(静态图) | Flow/Task(动态Python) | Asset(数据资产) | Workflow/Activity |
| 触发方式 | 时间Cron/外部触发 | 时间/事件/手动 | 时间/传感器/事件 | 代码调用/信号/定时 |
| 运行模型 | 调度器+执行器 | 调度服务+Worker | Daemon+Run Worker | 独立Worker重放执行 |
| 动态流程 | 弱,DAG解析期确定 | 较强,运行时动态 | 较灵活,但以静态资产为主 | 极强,普通代码风格 |
| 数据血缘 | 弱,需插件 | 一般,缺乏深度 | 强,资产级血缘 | 弱,不关注数据资产 |
| 长时运行 | 不支持 | 有限支持 | 有限支持 | 原生支持 |
| 失败恢复 | 任务级重试 | 任务级重试 | 任务级重试 | 工作流级持久恢复 |
| 运维复杂度 | 中高 | 中(Cloud版低) | 中高 | 高 |
| 上手曲线 | 中 | 低 | 较高 | 高(需理解确定性重放) |
| 适合场景 | 定时批ETL、数仓调度 | 现代化Python数据编排 | 数据治理、血缘、数据平台 | 微服务编排、长事务、业务流 |
这张表的核心信息就一句话:Airflow是"调度"的王者,Prefect/Dagster是"数据编排"的进化版,Temporal则是"流程可靠性"的终极答案。四个工具之间不是迭代关系,而是各自占据不同赛道。
5.2 混合使用:何时让多个工具共存
我见过不少团队在"只有一个编排工具"的思维里打转,其实现实中,成熟平台往往是多个工具各管一段,形成混合架构。
一种常见组合是"Airflow + Temporal"。Airflow继续负责数据仓库的离线批量ETL,每天凌晨跑一堆规规矩矩的批任务;Temporal负责业务侧的长流程,比如订单处理、审批流、跨系统调用。两者之间通过消息、API或数据库表做连接,互不干扰。这个组合的好处是,每个工具都只用在自己最擅长的场景,避免了"一个工具万能"的妥协。我现在的团队就是这么用的,Airflow只管数仓,Temporal管业务编排,职责边界写进研发规范,效果很好。
另一种组合是"Dagster + Temporal"。如果你比较新,不想碰Airflow的老生态,可以用Dagster做数据资产层的编排和血缘管理,同时沉淀数据平台自身的治理能力;业务侧的长事务交给Temporal,两边通过Dagster的Sensor或Temporal的Signal联动。这个组合适合数据产品和业务产品并行迭代的公司。
还有一种更轻的"Prefect单干"路线。如果你的业务复杂度还没到需要Temporal的程度,又不想被Airflow的DAG模型憋死,Prefect可以同时承担一部分轻量工作流编排。不过我不建议在Prefect里写太多需要持久恢复的长流程,它毕竟不是专为长时可靠执行设计的。
核心原则很简单:不要让一个工具背所有锅。选型时,先看主场景是哪一类,再决定主引擎是谁,其他场景用辅助工具或直接写服务代码,不强行收敛。
5.3 迁移实操经验
如果你决定从Airflow迁到新的编排平台,我劝你千万别做"一次性大爆炸迁移"。Airflow DAG是团队积累了几个月甚至几年的逻辑资产,直接全部重写,风险太大,业务等不起。
我的建议是三步走。
第一步,先做批次画像。把现有DAG按"凌晨批ETL""事件驱动任务""跨系统依赖任务"分类,统计每个DAG的平均运行时长、失败率、上下游依赖关系。画像做完,你才知道哪些任务该保留在Airflow,哪些该走新平台。
第二步,选一个业务影响小、失败容忍度高的试点任务,在新平台上用Pyhton重写一遍,跑一段时间灰度验证。重点看新平台的任务成功率和运行耗时,对比旧平台基线。我当时用一个小时级的报表任务做试点,在Dagster上跑了两周,跑完发现血缘视图真香,但调度准确率也和Airflow持平,这才放心扩大范围。
第三步,分批迁移+双跑回退。每批任务迁移后,保留旧Airflow里的DAG,但暂停自动触发,手动触发一次做对比。新平台跑出同样结果、下游消费方无感知后,再彻底关停旧的调度。回退接口要做好,一旦新平台出问题,能在几分钟内把调度切回Airflow。
每次迁移最大的风险不是工具本身,而是团队习惯和运维监控体系的切换。建议提前把新平台的告警、日志、权限接入公司现有体系,别让工程师在新平台上"裸奔"。
6. 常见问题与排查技巧实录
最后分享一些我在实际生产和迁移过程中踩过的比较有代表性的坑。这些内容不好从官方文档里直接找到,属于典型的需要"实操现场记录"的部分。
6.1 问题排查速查表
我把高频问题整理成一张方便比对的表,你在排查时可以直接对照。
| 现象 | 可能原因 | 排查方向 | 解决建议 |
|---|---|---|---|
| Airflow调度不准点,DagRun延迟 | DAG文件解析过慢、Scheduler过载 | 检查scheduler日志解析时长、元数据库连接 | 优化DAG import,拆分大DAG,调整scheduler参数 |
| Airflow任务卡在running状态不结束 | worker存活检测失效、心跳丢失 | 查看celery/k8s executor日志 | 升级容错策略,配置pod独立运行逻辑 |
| Prefect Flow运行失败但UI看不到详细日志 | 日志输出级别不够或没接入日志存储 | 检查日志配置和默认Log handler | 手动配置Log打印或接入统一日志平台 |
| Dagster资产长时间不物化 | Daemon崩溃、调度器未启动 | 查看daemon日志、心跳状态 | 重启daemon,配置高可用 |
| Temporal Workflow重放报错 | Workflow代码使用了不确定函数 | 检查代码里的time.Now、random等 | 改用Temporal提供的time和random封装 |
| Temporal Activity无限重试导致下游压力大 | retry policy配置过于激进 | 检查Activity的retry policy | 设置最大重试次数和指数退避 |
| 混合架构下任务互相阻塞 | 两套工具使用了同一个资源池/数据库 | 查询运行队列、锁情况 | 分配独立资源或拆分队列,加熔断 |
6.2 排查过程的真实案例
举一个我印象比较深的例子。有一段时间,Airflow里一个晚间的整表同步任务频繁延迟,调度器显示TaskInstance已经启动,但实际执行日志一直不更新,持续几小时不结束。一开始我以为是数据源网络慢,查了一圈发现不是,后来看worker机器的进程列表,发现这个任务把内存几乎吃满,导致worker进程假死,心跳发不出去。Airflow的task如果长时间没收信号,在调度器眼里它还是"running"状态,可实际上已经卡死了。
这个问题的深层原因,是任务代码在同步大表时使用了不合理的fetch逻辑,把全量数据先拉进内存再落库,而不是流式写入。优化代码改成游标分批拉取之后,问题就解决了。这件事给我的教训是:Airflow之外的执行环境,内存和网络边界如果不在设计时预判,调度的表象再健康也救不了任务本身。
另一个案例是我们在使用Temporal时,把一个外部API调用直接写进了Workflow代码里。在开发环境一切正常,直到一次worker容器被调度到新节点,重放Workflow时,外部API被重新调用了一次,结果那边生成了两条重复的订单记录——这个问题就是典型的确定性约束被破坏。后来我把所有外部交互都搬进Activity,再配合幂等键,把重复调用的问题从根上解决了。从那之后,我审计新Temporal代码时,第一件事就是在Code Review清单里加一条:Workflow里禁直接做IO,外部副作用一律走Activity,且Activity要设计成天然幂等。
还有一个我低估了很久、直到线上才发现的坑:Temporal的事件历史增长。当时我们有个常驻型Workflow,会监听大量外部信号并据此做步骤更新,跑了一周之后,这条Workflow每次重放的时间从毫秒级涨到了秒级。原因是每次Signal都会在事件历史里追加记录,历史越长,重放成本越高。排查后,我们对这类"无限期常驻+高频交互"的Workflow做了重构,拆成短生命周期的多个Workflow,用Signal和ChildWorkflow做联动,重放成本直接降了下来。
这些事如果只读文档,你可能永远不会踩到,但一旦踩到,又往往让人措手不及。所以我的建议是:上生产之前,一定要做充分的故障演练,模拟worker宕机、数据库抖动、网络分区、事件历史膨胀等场景。编排工具的价值在于可靠性,可靠性必须靠实验验证,不能靠信仰。
结尾一点真实体会
做了这么多年技术选型和平台建设,我逐渐认同一个观点:没有最好的工具,只有最合适的边界划分。Airflow、Prefect、Dagster、Temporal,每一个都是某些场景下的最优解,也都不是万能的。你真正要做的,是先想清楚自己手里的是什么问题,再去挑工具,而不是让工具反过来决定你的架构。
如果非要说一条最重要的建议,那就是"别让一个工具承载所有幻想"。Airflow解决不了长事务,Temporal也解决不了数据血缘,混着用不是耻辱,反而是工程成熟的标志。我自己踩过的坑,大多数都源自"试图用一个工具把世界的复杂性全都收纳进来",后来又都得靠拆分和边界划清来收拾。
最后送大家一个小习惯:每当你准备引入一个新的编排工具,先花半天时间把这个工具最核心的抽象模型前前后后想明白,再用一张小流程图画出自己的主数据流,如果这半天想明白了还觉得合适,那大概率就是真的合适。想不明白的时候,不要硬上。