ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

用Rust构建轻量任务流引擎:DAG调度与并发控制实战

用Rust构建轻量任务流引擎:DAG调度与并发控制实战 我们先聊个具体的场景你手里有几十个数据任务要按顺序跑有的能并行有的必须等前置完成中间还有失败重试、超时熔断、并发控制这些破事。你当然可以用现成的调度平台但很多场景下它们太重光是部署集群、配权限、写死配置文件就够你喝一壶的。我自己最后选择用Rust写了一个轻量的任务流引擎名字叫ruflo专门解决“我不想引入一套重型调度系统但又需要把任务编排、依赖、并发、重试都管起来”的痛点。这篇博文不打算只贴一堆代码我会把ruflo从需求拆解、核心设计、代码实现到性能调优和踩坑过程完整还原一遍。无论你是想抄一个类似的工具还是纯粹好奇一个任务流引擎内部是怎么转的这篇文章都有你能直接拿去用的部分。1. ruflo是什么这次为什么要自己拼一个任务流引擎1.1 背景我遇到的实际问题先说清楚我当时面对的事。团队里的数据管道是几个Python脚本用shell一个个串起来的靠cron定时触发脚本之间用“上一步成功写出文件再跑下一步”这种土办法来保证顺序。一开始数据量小没事后来任务多了问题就冒出来了任务A、B、C明明可以并行但串行脚本只能一个一个跑整体耗时翻了几倍。某一步失败以后后续步骤全部白跑而且没有自动重试半夜告警响了也没人处理。想要加一个新任务就得改shell脚本改完还可能影响旧流程特别容易漏。我当时也评估过Airflow、DolphinScheduler这类成熟框架但对我们这个体量来说太重了。部署要资源、学习要成本、配置模板一堆而我们其实只需要把“任务依赖”和“并行执行”这两件事做好。还有一个私心我一直想找个真实项目练手Rust与其用别人封装好的东西不如自己写一个顺手的小引擎于是ruflo就出来了。1.2 ruflo的设计目标与技术选型ruflo的目标我从一开始就定得很克制不追求做成一款通用分布式调度平台只做一个进程内可嵌入的异步任务流执行引擎。也就是说你的服务或脚本启动之后把任务节点和它们的依赖关系注册进来ruflo负责按照DAG有向无环图的拓扑关系去调度保证该等的一定等能并行的尽量并行。选型上我用了Rust加tokio运行时。为什么是Rust不是Go或者Python一个是性能Rust的异步任务在内存占用和调度开销上确实有优势另一个是工程体验所有权和类型系统能拦住一大批并发问题。编译期就把数据竞争、空指针这类坑消灭掉一大半写调度逻辑的时候心里踏实很多。tokio是Rust生态里最成熟的异步运行时任务调度、定时器、信号量这些基础设施都现成我只需要把流程编排的骨架搭起来就行。2. 核心设计思路从DAG到背压每个决策背后的为什么2.1 为什么用DAG做任务编排模型任务流本质上就是一个有向图节点是任务边是依赖关系。之所以限定为有向无环图是为了避免循环依赖导致永远无法终止。你不可能让任务A等B、B等C、C又等A这在逻辑上就是死锁。ruflo在注册节点的时候就会做环检测发现成环直接拒绝启动从源头上掐死这类问题。DAG的执行顺序靠拓扑排序来保证。拓扑排序的意思很简单每次挑出“所有前置任务都已完成”的节点来执行执行完就把它的后续节点的依赖计数减一减到零就可以进入就绪队列。这个过程反复进行直到所有节点完成。ruflo里我维护了一个依赖计数器数组每个节点记录还有几个前置没跑完前置跑完一个就减一减到零就通知调度器“我现在可以跑了”。这个方案简单高效复杂度是O(VE)V是节点数E是依赖边数。2.2 调度引擎怎么设计就绪队列加并发水位线核心调度循环我做得比较直接启动时先扫一遍所有节点把没有前置依赖的节点全部丢进就绪队列。就绪队列使用tokio的mpsc channel来实现容量可以配置。调度器从channel里拉出节点再根据当前并发水位来决定立即执行还是暂时挂着。并发水位线是整个引擎最关键的参数。它控制同一时刻最多允许多少个任务并行跑。我见过不少任务流工具并发控制做得很粗要么全局互斥退化成串行要么完全放开导致下游被压垮。ruflo用的是信号量方式一个tokio::sync::Semaphore实例初始化时设置最大并发数。每个任务要执行之前必须acquire一个许可执行完成再release。这样即使你有100个任务同时就绪也会被限制在预设的并发水位以内。参数选择思路是这样的多少并发合适取决于下游系统的承受能力。如果你要写数据库并发太高会把连接池打满如果只是CPU计算并发太高会导致上下文切换浪费。我一般建议先从2开始压测后逐步往上加直到延迟和错误率出现明显拐点那就是你的黄金并发数。这个值不需要很精准只要不把下游打爆就行。2.3 背压处理的细节有界通道和信号量怎么配合背压这个词听着绕实际就是“生产速度超过消费速度时怎么办”。ruflo里有两层生产消费关系一层是就绪队列的写入方前置任务完成回调和读取方调度器另一层是调度器和真正的任务执行。就绪队列我用了有界channel容量默认1024。有界意味着写入方在队列满的时候会等待而不是无限堆积。这样当任务产生的速度远超调度器消费速度时链条上游会自动慢下来不会内存爆炸。这个设计和消息队列里的消费者优先、队列满阻塞生产者是同一个道理。信号量相当于第二层保险它限制的是任务真正执行时的并发数。为什么两层都要因为有界channel管的是“等待调度的任务数量”信号量管的是“正在执行的任务数量”。如果没有信号量即使队列不爆也可能同时有几百个任务执行把CPU和下游都压垮。这两个配合起来一个管入口流量一个管并行执行量任务流才不会失控。3. 实操过程从零接入一个可用任务流3.1 安装与项目初始化ruflo以库的方式提供通过cargo引入就行。我发布在crates.io上版本号跟随语义化规范目前的稳定版本是0.4.x。在你的Cargo.toml里加上[dependencies] ruflo 0.4 tokio { version 1.0, features [full] }然后写一个最简单的入口use ruflo::{Flow, Node}; #[tokio::main] async fn main() - anyhow::Result() { let flow Flow::new(first_demo) .node(Node::new(step1, async || { println!(step1 running); Ok(()) })) .node(Node::new(step2, async || { println!(step2 running); Ok(()) })) .dependency(step1, step2)?; flow.run().await?; Ok(()) }Node::new的第一个参数是节点名字用来标识唯一性第二个参数是异步闭包真正的任务逻辑写在那里。dependency(step1, step2)的意思是step2依赖step1step1必须先跑完。整体跑起来以后控制台会先打印step1再打印step2顺序是确定的。3.2 定义你的第一个复杂Flow实际项目里任务不会只有两三个我举个更贴近真实情况的例子做一次数据同步需要先拉取源数据fetch然后清洗clean清洗之后两条分支并行一条做统计aggregate一条做备份backup最后等两个都完成再发送通知notify。用ruflo来定义这个流程就非常直观use ruflo::{Flow, Node}; async fn fetch_data() - anyhow::Result() { Ok(()) } async fn clean_data() - anyhow::Result() { Ok(()) } async fn aggregate() - anyhow::Result() { Ok(()) } async fn backup() - anyhow::Result() { Ok(()) } async fn send_notify() - anyhow::Result() { Ok(()) } let flow Flow::new(sync_flow) .node(Node::new(fetch, fetch_data)) .node(Node::new(clean, clean_data)) .node(Node::new(aggregate, aggregate)) .node(Node::new(backup, backup)) .node(Node::new(notify, send_notify)) .dependency(fetch, clean)? .dependency(clean, aggregate)? .dependency(clean, backup)? .dependency(aggregate, notify)? .dependency(backup, notify)? .concurrency(3) .build()?; flow.run().await?;执行的时候fetch先跑然后clean到aggregate和backup这里会并行启动都完成以后再触发notify。concurrency(3)表示最多同时跑3个任务这里两个分支并行完全够用。如果你想观察执行细节可以打开日志开启后每个节点开始、结束、耗时、失败重试都会有记录。3.3 超时、重试与并发参数的计算任务流引擎如果只负责顺序调度那和脚本串行没区别。真正提升可用性的是超时和重试机制。ruflo里每个节点都可以单独配置超时时间和重试策略Node::new(fetch, fetch_data) .timeout(std::time::Duration::from_secs(30)) .retry(3) .retry_backoff(std::time::Duration::from_secs(2))这个配置的意思是fetch任务最多跑30秒超时就直接判失败失败后自动重试最多重试3次每次重试之前固定等待2秒作为退避间隔。如果3次都失败整个流程会记录失败状态并且默认不再执行依赖它的下游节点。超时参数怎么定我一般会先测出任务正常情况下的P99耗时然后乘1.5到2作为超时值。如果正常需要10秒设15到20秒比较合理。设太短会误杀慢任务设太长又起不到保护作用。重试次数也不是越多越好重试3次已经是比较保守的上限如果3次都失败说明问题大概率不是偶发继续重试只会加重下游压力。还有一个容易忽略的参数是整体超时。如果一个流程的总执行时间有硬性要求比如必须在5分钟内完成可以在Flow上设置全局超时flow.overall_timeout(Duration::from_secs(300))这样即使某些节点重试还没结束整体超时一到ruflo会取消尚未完成的任务把流程标记为失败。这个机制在实时性要求高的场景下特别有用避免任务流卡在某个环节拖垮整个链路。4. 性能实测与调优记录4.1 基准压测效果理论讲再多不如直接看数据。我在一台4核8G的Linux服务器上做了压测机器配置很普通目标场景是模拟1000个任务节点依赖关系随机生成确保是一张合法的DAG。对比了串行执行和ruflo并行执行两种情况。串行执行1000个任务每个任务内部sleep 10毫秒总耗时大约是10秒多一点。换成ruflo并发数设成8同样1000个任务总耗时就掉到了2秒左右。这里几乎所有收益都来自并行化理论上限是1000乘以10毫秒除以8个并发约1.25秒实际2秒是因为调度本身、唤醒开销和信号量竞争还有一部分损耗。我又加大规模跑了一万节点并发保持8单节点耗时仍是10毫秒总耗时大约15秒。相比串行需要100秒收益还是很明显的。而且过程中内存占用稳定在150MB以内没有出现内存泄漏或无限增长的情况。这个内存表现主要得益于前面说的有界channel和信号量任务执行完立刻释放资源不会堆积。4.2 内存与调优参数压测过程中我也试过把并发数调得很大比如100情况就不太一样了。1000个任务、每个任务内10毫秒sleep的情况下并发调到100总耗时反而没有比并发8快多少因为任务太轻、CPU调度开销占比变大。但内存峰值却从60MB涨到了180MB。这个现象说明一个道理并发数不是越大越好要匹配任务的实际负载。如果你用ruflo跑的任务偏向IO密集型比如HTTP请求、数据库读写可以把并发值设高一些因为等待IO时CPU是空闲的。如果是CPU密集型任务并发数最好等于CPU核心数最多再留一两个给调度器自己用。我的经验公式是这样IO密集型并发数可以设为核心数的2到4倍。CPU密集型并发数设为核心数或核心数加1。混合负载从核心数开始压测后逐步往上加找到拐点。还有一个调优细节是channel容量。默认1024对绝大多数场景都够用除非你的DAG特别深、单层就绪任务特别多可以让容量跟着最大宽度走。最宽的那层有多少个节点容量就设多少避免调度器因为channel满而阻塞拖慢整条链路。5. 踩坑实录我用ruflo实际开发中遇到的典型案例5.1 问题一任务集体“卡死”最后发现是我把依赖配反了第一次用ruflo跑一个稍微复杂点的流程时我发现所有任务都卡住不执行控制台没有任何报错程序像是在等什么永远等不到的东西。排查了半天最后发现是我把依赖方向搞反了。我本意是A依赖BB先跑完才能跑A结果写成了dependency(B, A)。这样B就一直等A但A根本不在就绪状态形成了事实上的互相等待。这类问题用DAG环检测其实查不出来因为A依赖B、B依赖A确实是环但如果只有一条边配反了图可能是合法有向图只是拓扑逻辑反了。后来我在ruflo里加了一个提示机制如果启动后一段时间没有任何节点被调度就把所有节点的依赖关系和状态打印出来方便定位是不是“死等”。如果你自己调试类似引擎这个思路可以直接抄调度器空闲超时后输出诊断信息比白屏卡死好排查一百倍。5.2 问题二重试风暴把下游数据库打挂了有段时间生产环境任务成功率波动很大排查后发现是重试机制太激进。当时我对每个写数据库的节点设置了重试5次、退避时间0.5秒。结果某次数据库慢查询第一批任务失败后立刻重试0.5秒相当于没退避紧接着第二轮又把数据库打得更慢接着触发更多节点超时失败形成雪崩。后来我把写类任务的重试策略改成了指数退避也就是每次等待时间都翻倍第一次失败等1秒第二次等2秒第三次等4秒最多5次。同时把重试上限从5次降到了3次。这样即使下游出问题上游也不会无限施压。指数退避是分布式系统里的经典策略用在这里的原则是一样的给下游留出恢复时间而不是火上浇油。5.3 问题三内存飙升问题出在前置依赖回调上另一个印象深刻的问题是内存无限上涨。当时我在节点完成回调里写了这样一段逻辑节点执行完成后获得一个共享的广播通知然后把结果缓存到一个全局HashMap里。看起来没什么问题但忘记做清理节点越来越多结果缓存越来越大最后内存吃满。这个问题的根源是我把一个任务流引擎用成了“事件总线”。ruflo本身不会保存节点执行结果如果你想在节点之间传数据应该显式设计数据流要么通过数据库、消息队列要么在Flow内部维护一个受控的结果集并且用完及时清理。不要图方便搞一个全局缓存内存失控只是时间问题。后来我在ruflo里加了广播机制让节点可以通过topic订阅其他节点的完成事件这样数据传递变得更可控也避免了全局HashMap的陷阱。5.4 问题速查表现象可能原因排查与解决任务全部卡住不执行依赖关系配反或形成事实死等打开诊断日志查看各节点依赖状态检查dependency参数方向执行结果错乱共享了可变状态节点间数据串了确认没有共享可变全局变量节点间传数据用消息或结果集API重试导致下游被打爆重试次数过多或退避时间过短缩重试次数改指数退避给下游留恢复时间内存持续上涨结果缓存未清理或channel无人消费检查全局缓存使用后是否释放确认channel容量合理任务执行延迟高并发数过大导致CPU争抢调低并发水位按CPU密集/IO密集调整并发参数6. 个人经验总结这个引擎后续还能怎么玩写ruflo这件事最大的收获不是我“发明”了什么新算法而是把任务流引擎这块本来模糊的地带彻底盘清楚了。DAG拓扑排序、信号量限流、有界队列、超时重试每个概念单独说都不难难的是把它们拼成一个整体还能保持稳定。如果你也准备自己写一个类似的工具我建议先想清楚边界哪些功能必须有哪些可以不要。ruflo现在的定位就是“进程内可嵌入的轻量异步任务流引擎”不为分布式场景负责不为持久化负责这让代码量能控制在可理解的范围内出了故障也能快速定位。后续我打算在几个方向扩展ruflo。一个是把任务执行历史持久化到SQLite这样流程跑完以后还能复盘每个节点的耗时和状态对排查线上问题很有帮助。另一个是加一个简单的HTTP管理接口可以在不重启服务的情况下查看当前流程运行状态甚至手动触发某个节点重新执行。还有一个想法是做分布式协调不过那个水太深可能会直接依赖etcd的选主能力而不是自己从零实现。如果你只是需要一个趁手的任务编排工具其实不一定非要用我的库。你可以把这里面的设计思路抄走用你熟悉的语言写一个简化版。重要的是理解那几件事依赖怎么表示、并发怎么控制、失败怎么处理。这三件事想清楚你自己的任务流引擎就已经成功了大半。
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进