ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

长任务编排选型指南:Airflow、Prefect、Dagster、Temporal生产实践

长任务编排选型指南:Airflow、Prefect、Dagster、Temporal生产实践 做数据平台或者后端基建的朋友这两年肯定被反复问过一个问题长任务编排到底用什么Airflow太老Prefect太新Dagster太超前Temporal又不像一个调度器……每次技术选型群聊最后都会变成一场各说各话的辩论。我这些年前后在几家公司把这四个工具都实际用过一轮从第一代用Airflow撑起数仓调度到后来在业务团队里拿Temporal重做了一套审批和支付异步链路中间还顺手用Prefect和Dagster做过实验性项目踩过的坑足够写一本小册子。这篇文章不打算做那种功能对比式的表格堆砌而是想从生产环境的角度把这四个工具的底层逻辑、适用边界、以及选型时真正要问自己的几个问题讲清楚。下面的内容基于我个人在真实业务里的实践不一定代表所有场景但至少能给你一个靠谱的思考框架。1. 先用一句话说清楚这四个工具分别解决什么问题很多人第一反应是把Airflow、Prefect、Dagster、Temporal放在同一张表里比功能但这四个东西其实不完全是一个物种。Airflow是最经典的DAG调度器它的核心贡献是把定时跑批这件事做成了行业标准。你用Python代码描述任务依赖调度器按时间触发执行器把任务发给worker元数据库记录每次运行的状态。这套模型是2005年左右Airbnb为了解决数仓ETL问题设计的十几年下来积累了海量用户和成熟生态。但它的代价是DAG在调度器启动时解析任务之间的数据传递靠XCom这个小管道动态性和实时性都很弱。如果你只做T1的离线数仓Airflow至今仍是稳妥选择。Prefect和Dagster都是冲着Airflow的痛点来的。Prefect的口号是Python原生工作流它允许你在运行时动态生成任务支持任意的Python控制流不需要像Airflow那样把整个图静态化。Dagster则换了一个角度它把数据资产作为一等公民强调你关心的不是某次任务跑没跑而是某张表、某个模型、某个指标是否可信可用。这两个工具都比Airflow现代但它们的受众其实有差异Prefect更像给开发者的便捷工具Dagster更像给数据团队的数据治理平台。Temporal又是完全不同的逻辑。它不关心定时调度也不关心数据血缘它关心的是任何一段业务逻辑能不能被可靠地持久化执行。你把一个工作流写成普通代码Temporal负责在代码执行到一半机器宕机时从某个事件点恢复重放负责在网络超时、下游返回异常时按你的策略重试负责把每个步骤的进度和结果完整记录成历史事件。它本质上是一个分布式执行引擎而Airflow、Prefect、Dagster本质上是调度器加任务执行管理器。这四个工具的底层差异决定了它们在不同生产环境里的表现完全不同。下面我从技术哲学、分水岭问题、典型场景和落地经验四个维度展开讲。2. Airflow、Prefect、Dagster的技术分叉一条线上长出的三个方向2.1 Airflow的霸权与债务先得承认Airflow能火这么多年不是没有原因的。它的DAG模型非常直观写一个Python文件定义任务和依赖提交给调度器就行。团队里哪怕是刚入门的数据工程师在看懂两个例子之后就能上手。生态方面更是没得说Hive、Spark、Snowflake、S3这些常用系统的Operator基本都有人维护遇到问题Google一下就有答案。生产环境里的BI报表、数仓建模、数据质量检查用Airflow跑得很顺。但Airflow的痛点在使用一段时间后会集中爆发。第一是调度器的性能瓶颈DAG太多、任务太碎的时候调度器会成为瓶颈你需要调max_active_dag_runs_per_dag、max_active_tasks_per_dag这些参数要么牺牲并行度要么无限加worker。第二是DAG的静态性Airflow在调度周期启动时会解析整个DAG文件如果你根据上游结果动态生成下游任务官方推荐的做法是Dynamic Task Mapping但这套机制在1.x时代很难用2.x时代有所改善但依然不自然。第三是XCom的局限任务之间传数据只能用XCom而XCom存数据库传大数据时会撑爆元数据库。实践中我们不得不在任务里把中间结果写到对象存储再用下游任务去读。这种问题不是Bug而是架构决定的。Airflow的设计前提是运行计划已知、任务边界清晰适合稳定、可预测的批处理流程。但这恰恰是长任务编排最不喜欢的特征——长任务往往意味着运行时间长、中间可能动态扩展子任务、失败后希望精确断点恢复而不是把整个DAG从头再跑一遍。2.2 Prefect如何重做Airflow没做好的事Prefect 2.x现在叫Prefect 3.x最核心的改动是把DAG从静态定义改成了运行时编排。你在Python函数上加上flow和task装饰器工作流就是普通的Python函数调用里面可以写for循环、if-else、也可以根据上游输出动态生成新任务。这个模型对开发者非常友好你不用再纠结DAG文件该写在哪个目录、变量怎么传。Prefect的另一大优势是缓存与重试机制设计得比较精细。它把每一个task run的状态、缓存key、重试策略都记录在后端数据库里同样的参数重复跑同一个task可以直接走缓存结果。这在数据开发和ML特征计算的场景下很实用。Prefect还自带了比较完整的UI可以看flow run日志、任务耗时甘特图、失败任务的重试记录对于中小团队来说省掉了搭一套监控面板的时间。不过Prefect并不是没有代价。它把很多逻辑放到了客户端侧你用装饰器装饰的函数会在工作进程里执行所以它对worker的稳定性、网络环境、Python环境一致性要求比较高。在实际落地中我发现Prefect适合跑那种开发频次高、任务数量中等、依赖关系经常变的业务。假如你是一个有5000个稳定定时任务的数仓平台Prefect并没有比Airflow高明太多反而因为社区生态不如Airflow丰富很多现成的connector要自己写。2.3 Dagster把资产放在了任务前面Dagster是我个人觉得理念上最反直觉、但一旦理解就回不去的工具。它不让你先定义任务再定义调度而是让你先定义资产asset比如一张表、一个训练好的模型、一份指标报表然后让框架反向推导出这个资产由哪些任务产出、依赖哪些上游资产。这个设计对数据团队有非常大的现实价值。第一它天然支持数据血缘的可视化和数据质量检查的挂载你可以给每个资产挂上freshness check或者自定义断言一旦资产新鲜度不达标上游所有依赖它的任务都会受牵连第二Dagster的launchpad和UI对非工程师群体也友好业务分析师能看到我今天看的这张报表依赖哪些底层表不用去翻调度代码。第三Dagster在分区partition和传感器sensor上做得比Airflow好它有一套明确的数据版本概念处理增量数据场景时你不再需要自己维护一堆分区变量。但Dagster的短板也很明显它对团队成员的要求比Airflow高。你不能再随心所欲地把一堆任务塞进一个全局DAG里而是要按资产来组织代码和配置触发方式也从定时变成了资产物化思维。对于没有专门数据工程团队、只能靠开发顺手维护的公司来说Dagster的抽象层显得厚重。3. Temporal是另一种生物它不是调度器而是分布式执行内核3.1 从画DAG到写代码Airflow、Prefect、Dagster再怎么现代化本质上都是定义一张有向无环图然后调度节点执行。Temporal直接把这张图扔掉了。你在Temporal里写工作流就是写一个普通的异步函数函数里await某个activity的执行结果然后继续下一步。这个工作流代码会被Temporal SDK编译成一个可重放的状态机由Temporal Server记录每个事件的进度。这是两个完全不同的心智模型。用Airflow你要思考的是我的任务A结束后任务B什么时候可以开始任务B失败了任务C还要不要跑用Temporal你思考的是我要给用户创建一个订单先扣库存再申请发票如果发票服务超时5秒后重试最多三次期间订单状态保持为处理中。后者其实就是写业务代码只不过这段代码跑在了一个不会因为机器宕机就丢失状态的平台上。3.2 持久化执行、确定性重放和历史事件Temporal能保证不会因为机器宕机丢失状态靠的是事件溯源Event Sourcing加确定性重放。工作流的每个步骤完成后SDK会把步骤产生的命令比如执行activity X参数Y发给ServerServer以事件形式追加到工作流的历史里。如果worker进程崩溃Temporal会从最近的事件点重新恢复工作流执行SDK重放历史事件直到恢复到崩溃前的状态再从那里继续。这个机制有一个硬性要求工作流代码必须是确定性的。你不能在里面用time.Now()、不能读随机数、不能直接访问外部API因为这些行为在重放时会产生不同的结果导致状态不一致。正确做法是所有不确定操作都封装在activity里由Temporal保证activity执行进度和结果也被记录下来。这个约束刚上手时会很别扭但习惯之后你会觉得这是解放——因为一旦满足了确定性重试、恢复、观测都变成了平台能力而不是你业务代码里的try-catch。3.3 和前三者最本质的区别有人说Temporal也可以做定时任务啊它有Schedule API这话对但低估了它。Temporal真正擅长的是那些生命周期长、状态复杂、需要跨服务协调的工作流比如订单支付流程先锁库存→扣款→通知发货→更新积分任意一步失败都要有补偿、机器学习训练任务启动容器→等待训练完成→上报指标→触发评估训练可能要跑几小时、音视频转码上传→转码→截图→审核→发布单条任务可达十几分钟甚至小时级。这些场景如果放在Airflow里你需要把每一步设计成独立任务再用传感器轮询外部状态稍微一复杂就变成一堆while True: sleep(30)的垃圾任务。放在Prefect里动态编排是解决了但Prefect的容错和故障恢复能力并没有专门为跨服务长事务设计worker的断线、任务的长时间运行会带回很多边界问题。放在Temporal里这就是它的主场activity心跳、超时、取消、信号signal、子工作流child workflow这些机制全部是为长事务而生的。4. 生产环境选型时要跨过的五个分水岭问题4.1 你的任务是计划性的还是事件性的这是我觉得第一道分水岭。所谓计划性任务就是每天凌晨2点跑前一天数据、每个整点同步一次订单表这类任务周期固定、依赖明确用Airflow或Prefect的schedule就能解决。所谓事件性任务是当某笔支付回调到达时开始执行一个多步骤的账户入账流水、当某个工单被创建时启动一个审批和通知流程。事件性任务往往没有固定时间表它是被外部消息触发的状态多、持续时间不定这时候更适合用Temporal。一个项目里两种任务都存在也很常见。你可以让Airflow负责定时跑批让Temporal负责事件驱动的业务编排两者之间用消息中间件对接。这比硬生生把业务事件都塞进Airflow要合理得多。4.2 任务图是静态的还是动态的Airflow的DAG是解析期确定的虽然支持Dynamic Task Mapping但它骨子里不是为运行时无限分支设计的。如果你预测性地说我要根据每个用户的不同状态生成不同后续步骤最好直接用Prefect的动态工作流或者Temporal的代码内分支。反过来如果你们的任务图几十年不变新增任务只是照着既有模板复制粘贴那Airflow的稳定性和生态优势就是加分项。实际经验是数据平台里大多数ETL任务是静态的业务中台里的长流程大概率是动态的。这个特征直接决定了你是选调度器还是选执行引擎。4.3 失败重试和状态恢复你希望谁负责这个问题问的是当任务失败后你希望重新执行整个流程还是从失败点精确恢复。Airflow默认是在task层级重试如果失败发生在task 3重试时只会重跑task 3但前提是task 3本身设计成幂等的你可以重放它的输入和参数。这个模型在数据同步类任务里基本够用。Temporal走的是事件重放activity精确结果复用的路线如果你的activity具备确定性同样的输入得到同样的输出失败后Temporal能保证只重新执行失败的那一个activity其余已完成步骤直接复用历史记录中的输出。这对长流程来说太关键了——一个跑了两小时的转码任务如果在最后一步网络抖动失败你是想整个流程重新来一遍还是只重试最后一步后者才是生产环境真正需要的。Prefect和Dagster在重试上做得比Airflow好一些但它们的恢复粒度仍然在task/op层级和Temporal的事件级恢复不是一回事。4.4 你拿什么做可观测性和审计可观测性在长任务编排里不是加分项是保命项。Airflow有标准的状态页面、日志聚合、任务持续时间图表但任务内部发生了什么它管不了。Prefect的UI在flow内部细节展示上更友好可以看到每个task run的参数、缓存命中情况、重试次数。Dagster在资产血缘、分区覆盖状态、数据新鲜度这些维度上最强适合对数据质量治理要求高的场景。Temporal则提供了一套非常强大的历史事件查看器Event History可以精确看到工作流每一个步骤在几毫秒内发生的事件、输入输出、超时原因。这种级别的审计能力在排查订单状态为什么卡住这类问题时无可替代。4.5 团队愿意为学习曲线买单多少这是最实际的问题。Airflow的Python DAG模型几乎人人能上手资料多、踩坑贴多招聘成本最低。Prefect和Dagster需要团队接受新的执行模型尤其是Dagster的资产优先思想很多老数据工程师转过去会骂娘。Temporal的学习曲线更陡不但要理解确定性和重放而且要懂一点分布式系统概念如namespace、task queue、worker identity前端业务团队学起来需要一到两周的集中培训。所以我的看法是选型不是选最好的工具而是选最匹配你们团队当前认知水平和业务形态的工具。5. 按团队画像对号入座我的推荐矩阵下面这个矩阵是我根据实际项目经验总结的不是官方文档结论仅供参考。团队/场景首选工具理由离线数仓、BI报表、定时ETL任务量大且稳定Airflow生态最成熟、稳定性有保障、招聘维护成本低数据平台ML特征任务动态生成需频繁改动工作流PrefectPython原生、动态任务、缓存和调度体验好数据治理驱动关注血缘、分区、数据质量Dagster资产模型能把数据质量和调度统一业务长流程如支付、审批、订单调度、IT自动化Temporal持久化执行、精确重试、历史审计能力最强再补充几个混合架构的实例都是我实际见过或做过的方案方案AAirflow TemporalAirflow专门负责所有定时任务Temporal专门负责业务事件流。两者通过Kafka或数据库表解耦。一个典型链路是Airflow每天凌晨跑完数据同步后生成一批待处理事件Temporal工作流消费这批事件执行一系列需要人工审批的后续操作。好处是两边都不越界坏处是两套系统需要两拨运维。方案BPrefect先落地后续按需引入Temporal如果团队之前没有用过Airflow反而更容易接受Prefect。等Prefect跑起来之后把其中那些跨系统调用、需要精确重试、状态复杂的任务逐步迁移到Temporal。我见过一个推荐系统团队就是这样做的Prefect负责每天的批量特征计算Temporal负责在线推理和反馈回路里的长流程协调。方案CDagster做数据平台Temporal做业务平台这是我目前比较推崇的数据中台范式。数据侧用Dagster管理资产和血缘业务侧用Temporal管理跨服务流程。两边用数据湖/仓库完成对接Dagster负责把数据变成可信资产Temporal负责把资产变成业务动作。6. 落地与迁移过程中的实操经验6.1 部署与运维的几个硬指标先说Airflow。生产环境一定要把scheduler单独部署不要和webserver挤在同一个进程里元数据库别用SQLite用PostgreSQL并做好定期备份。用CeleryExecutor时要注意broker里的任务积压监控用KubernetesExecutor时要考虑Pod启动延迟对任务调度时延的影响。我见过一个团队把所有DAG文件放在同一个仓库几百个任务全在一个DAG里调度器直接跑崩后来拆成按域隔离才稳定。Prefect部署比Airflow轻很多核心是prefect server加worker。要注意的是版本兼容性Prefect 2.x和3.x的API有不小差异升级前一定要读迁移文档。另一个坑是task的幂等性设计Prefect的缓存命中粒度很细参数、输入、缓存key稍微不一致就失效导致缓存形同虚设写work pool和deployment时要统一日志和缓存的key规范。Dagster在部署上特别依赖一个集中式的code location管理。它的代码不是一个文件而是一个repository需要暴露给Dagster daemon去加载。多个团队协作时建议用分支环境隔离不同project否则每个人新增一个asset都可能影响全局图。Dagster的UI已支持资产图搜索但资产太多上千个时加载会慢需要开启asset_group分组合并展示。Temporal部署最重的其实是Server端依赖Cassandra/PostgreSQL和Elasticsearch可选用于增强可见性。好在官方提供了temporal server开发模式本地能快速跑起来。生产环境建议至少3节点集群namespace按照团队或业务线隔离。worker侧要注意task queue的命名规范不同业务团队不要混用一个task queue否则消息会串。6.2 迁移路径设计从Airflow迁移到Prefect或Dagster我建议先新后旧双轨并行。不要做一锤子大迁移那个风险太高。先把新开发的流程在新工具里跑通再挑两张怕脏数据影响的表试着迁移跑稳定一个季度后再扩大范围。Airflow生态里很多Operator的工作流本身就是Python代码迁移到Prefect时大部分逻辑可以复制只需要改装饰器和trigger方式迁移到Dagster时你需要把原有的DAG思维转成资产依赖思维这一步要花比较多时间所以建议先拿一个业务域做试点。从Airflow迁移到Temporal的情况比较特殊因为两者解决的问题交集很小如果你拿Temporal去替代全部Airflow一定会很痛苦。反过来Temporal也不能完全替代Airflow的定时批量触发能力。迁移时最好按业务价值来排优先级哪些流程的失败恢复成本最高、哪些流程手动补偿最频繁、哪些流程你永远不知道卡在哪一步这些优先迁。6.3 常见坑和应对建议我在实战中遇到的坑有几个值得分享Airflow的时区与会话坑start_date和schedule_interval的配合逻辑很多人搞错导致任务提前或推迟执行。不要试图自己推断直接打开官方文档的DAG run页面观察里面对logical date和execution date的解释。Prefect的flow代码里用了阻塞调用Prefect有一个特性是每个task run会占用worker的一部分资源如果你的task里用了requests.get这种同步调用建议用task装饰器包一层线程池或者直接用anyio异步处理否则并发上来会白占资源。Dagster的多仓库环境混乱多个团队共用Dagster时建议每个repo的asset key命名加业务前缀否则在全局资产图里会出现很多重复的源表节点非常难排查。Temporal workfow里使用非确定性代码前面说了time.now()、随机数、直接改全局变量都是雷。我见过有同事在workflow里直接连数据库结果replay时数据库连接无法重放整个workflow报non-deterministic error。所有网络IO和可变逻辑必须塞进activity。把worker和业务服务混部Temporal的worker最好独立部署因为它有背压机制worker负载过高时会导致heartbeat超时你以为任务挂了其实是worker线程阻塞了。独立的worker pool 单独的监控告警是标配。生产环境的选型没有银弹判断的标准永远是这个系统在你的团队里能不能长期稳定运行、能不能让人愿意维护。Airflow、Prefect、Dagster、Temporal每一套都有人用得好也每一套都有人用得痛苦。我个人的体会是先想清楚你要解决的核心矛盾——是批量调度的稳定性还是复杂流程的可靠性是数据的血缘和治理还是业务代码的执行韧性——再选工具而不是反过来拿工具去套业务。希望这篇长文能帮你少走一些弯路。
返回列表