ARTICLE DETAIL

资讯详情

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

Airflow、Prefect、Dagster、Temporal选型实战:从批处理到长任务编排

Airflow、Prefect、Dagster、Temporal选型实战:从批处理到长任务编排 做技术选型这事最怕的不是项目复杂而是方案多到不知道该从哪下手。这些年我在不同公司、不同团队里把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里有哪些TaskTask之间怎么连边什么时候由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(retries3, retry_delay_seconds60) def fetch_data_from_api(url: str) - dict: # 你的业务逻辑 return {data: whatever} flow(log_printsTrue) def etl_flow(api_url: str): data fetch_data_from_api(api_url) print(fGot: {data}) if __name__ __main__: etf_flow.serve(namedaily-etl, cron0 2 * * *)注意几个点fetch_data_from_api是一个Tasketl_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 fcreate table report as {upstream_table} daily_schedule ScheduleDefinition(jobDefinitions(...).get_job_def(__ASSET_JOB), cron_schedule0 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 一张表看清四个工具的差异维度AirflowPrefectDagsterTemporal核心抽象DAG静态图Flow/Task动态PythonAsset数据资产Workflow/Activity触发方式时间Cron/外部触发时间/事件/手动时间/传感器/事件代码调用/信号/定时运行模型调度器执行器调度服务WorkerDaemonRun 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也解决不了数据血缘混着用不是耻辱反而是工程成熟的标志。我自己踩过的坑大多数都源自试图用一个工具把世界的复杂性全都收纳进来后来又都得靠拆分和边界划清来收拾。最后送大家一个小习惯每当你准备引入一个新的编排工具先花半天时间把这个工具最核心的抽象模型前前后后想明白再用一张小流程图画出自己的主数据流如果这半天想明白了还觉得合适那大概率就是真的合适。想不明白的时候不要硬上。
返回列表