十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

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

用Rust从零写一个轻量级DAG工作流引擎:ruflo项目复盘

用Rust从零写一个轻量级DAG工作流引擎:ruflo项目复盘 ruflo这个热搜词背后没有任何现成的文档和上下文我一开始也愣了下。不过作为常年鼓捣自动化管线的工程师看到这个词下意识就把它拆成了rust flow。正好我过去几个月在Rust里从零写了一个轻量级DAG工作流引擎代号就叫ruflo本意是rbind-like flow scheduling的缩写。这篇文章算是一次完整的项目复盘把当初为什么绕开成熟框架、核心调度怎么设计、实战怎么跑通、之后踩了哪些坑原原本本梳理一遍。如果你也在找一个能在低配环境里跑起来、不想被一堆组件绑架的流水线方案这篇文章应该比看官方文档更有点参考价值。1. 我为什么在Rust里重造一个工作流轮子1.1 现成方案在边缘场景里的尴尬先说清楚我不是反对Airflow、Temporal这类成熟系统它们在大规模、多租户、复杂依赖的场景下确实是标准答案。但我的实际项目里有个很具体的痛点——需要在一台只有2核4G内存的边缘网关设备上跑数据处理流水线周期性采集数据、清洗、做简单聚合、再推送到上游服务。Airflow光是自身的scheduler、webserver、数据库依赖就吃掉不少资源在那台设备上启动都费劲Temporal更不用说它更适合那种需要持久化工作流状态、跨服务长时间运行的重型业务。早期我用cron加一堆shell脚本硬刚。采集脚本、清洗脚本、聚合脚本各管一段靠文件名和目录结构约定传参。今天加一个依赖昨天跑完的数据就在crontab里把时间错开加错开时间又得重新计算。上线头两个月还行第三个月开始频繁出事上游脚本重试还没结束下游脚本已经启动读到半截文件某天多个小时的任务堆积新任务和旧任务互相踩同一份输出文件告警脚本隔三差五误报因为任务前一天跑成功但数据不完整。按我这个场景去搜业界方案其实全部落在一个尴尬的空档里重量级方案需要独立部署、需要数据库、需要充足的内存。功能完整但在资源受限的边缘节点上是杀鸡用牛刀。轻量级方案make、just这类只能做最简单的依赖排序没有失败重试、状态管理、并发控制。中间方案dagu、jobflow这类工具覆盖了一部分需求可它们要么绑定自己的运行方式要么默认场景是桌面或者开发机而不是以作为库被嵌入的角度设计。我想得很简单我需要的不是一个平台而是一个能嵌入到现有Rust程序里的工作流调度库核心功能只有四件事——DAG依赖解析、并发调度、失败重试、状态持久化。于是ruflo就诞生了。1.2 ruflo到底解决什么问题ruflo的定位和那些工作流平台有明显的区别。它做的事情用一句话概括只负责把有依赖关系的任务按正确的顺序、以可配置的并发度、稳定地调度完并且能记住每个任务的状态。它不是服务不需要常驻进程它是一个库你把它静态链接进你的Rust程序在main函数里配置任务图然后调用调度器。任务本身是实现了一个trait的普通异步函数。这意味着你可以在一条采集程序的进程里同时跑HTTP服务、数据库连接池和ruflo调度器互不干扰。单机、多机部署方式都行。单机模式下数据状态落在本地SQLite或文件里多机模式把同样的任务图复制到其他节点配合一个共享状态后端做抢占调度后面再细说。没有自己的UI、没有REST API、没有部署agent。你想操作任务就用它提供的run、pause、resume这些Rust API想在Web上操作自己包一层。我把这个取舍叫做库优先平台后置。数据类工程有个默认思维既然要做任务编排干脆上一套平台。但平台意味着引入额外的运维、升级、授权负担。对资源受限的嵌入式设备和边缘计算场景这些负担是真实的成本。1.3 名字的来源和项目定位名字是ruflo来源就是Rust flow scheduling的拼接。第一版代码写出来的时候我原本想叫rflow但crates.io上面已经被占用了改成ruflo反而读起来顺口。项目定位从始至终没变过通用不绑定具体业务任务就是async fn() - Result()。小而精核心依赖只有tokio、serde、tracing和rusqlite不做超出调度范围的事。可嵌入所有功能都是库的形式不强制起服务。它适合的人大概是这几类不想部署一套复杂平台但cron又明显不够用的开发者需要把多个脚本/函数串起来的工具作者想在嵌入式设备或低配VPS上做轻量自动化的人。不适合的人我放在后面边界一节说。2. ruflo的架构小而完整的DAG调度内核2.1 核心抽象任务、依赖、触发器ruflo里所有东西围绕三个概念任务Task、依赖Dependency和触发器Trigger。任务是执行的最小单元。你只需要实现一个Tasktrait#[async_trait] pub trait Task: Send Sync static { fn name(self) - str; async fn run(self, ctx: Context) - ResultTaskOutput, TaskError; }name()用于定位节点ctx里塞着本次运行的参数、共享状态句柄、上一次运行结果。真正做事全在run()里一个任务可以小到发一个HTTP请求也可以大到把一个目录下的数据全量拉起来做训练——调度器并不关心任务内部做什么它只维护依赖关系和执行顺序。依赖关系用两种方式表达。第一种是显式声明dag.add_edge(ingest, clean)?; dag.add_edge(clean, aggregate)?; dag.add_edge(clean, quality_check)?;第二种是运行时动态生成的场景比如跑完所有城市的数据采集后再做全国汇总。只需要在任务里通过ctx.spawn_dependent(aggregate)声明动态子任务调度器会在主线任务结束时解析这些动态依赖并把它们插入待调度集合。这个设计在真实业务里非常有用因为很多数据流程的任务数量是依赖外部配置的没法提前写死。触发器决定这条DAG什么情况下该跑一次。ruflo在三层做了触发机制手动触发调用dag.run_once()。定时触发内置cron表达式解析支持0 0 * * *这样的习惯写法。事件触发通过EventBus监听外部事件比如收到上游webhook后启动数据同步。因为触发器走的是同一个调度内核定时和事件触发之间不会互相干扰也不会出现重复调度同一批次任务的问题。2.2 拓扑排序与并发调度策略调度器的第一件事是把DAG做拓扑排序。我用的Kahn算法算法本身很简单统计每个节点的入度入度为0的节点先进入就绪队列每执行完一个任务就把它的后继节点入度减1减到0再放入就绪队列。这一步保证任务永远按照依赖顺序启动。但拓扑排序只是能不能跑真正的调度策略还要回答同时能跑多少个。ruflo提供了三个维度max_concurrency全局最大并发任务数。task_concurrency同一类型任务最大并发数比如限制采集任务最多同时3个在跑避免疯狂打上游API。dependency_budget动态子任务总量上限防止运行时报的子任务把DAG撑爆。调度核心代码骨架长这样loop { let now Instant::now(); // 1. 所有入度为0且未执行的节点进入ready let ready: Vec_ graph .nodes() .filter(|n| n.indegree 0 !n.started) .collect(); if ready.is_empty() { if running.is_empty() { break; // 没有可运行任务且没有正在运行的任务调度结束 } // 否则等待任意任务完成 let task running.select_next().await?; // 处理成功 / 失败 / 重试 continue; } // 2. 按并发额度挑选候选任务 let selected select_tasks(ready, max_concurrency)?; for task in selected { let handle spawn_task(task); running.push(handle); } // 3. 等待已启动任务中的某一个完成 let completed running.select_next().await?; // 依据执行结果更新DAG节点状态 }核心点在于第二步的select_tasks不是简单取前N个就完事而是先按依赖深度排序深度大的优先后台并行再按task_concurrency过滤最后还要把当前系统负载作为软约束避免调度器把整台机器的CPU吃满。这里有一个容易被忽视的细节一个DAG里的就绪任务可能同时有几十上百个如果全放开并发下游瞬时负载会灾难性地上涨。所以ruflo默认对就绪任务做滑动窗口限流窗口大小由task_concurrency决定而不是把max_concurrency当作摆设。2.3 失败重试与超时控制怎么设计任务失败是常态所以重试逻辑必须在一开始就设计好而不是事后打补丁。ruflo里每个任务可以单独配置重试策略TaskSpec::new(sync_sales) .with_retry(RetryPolicy { max_attempts: 5, backoff: Backoff::Exponential { base_secs: 2, max_secs: 60 }, jitter_ratio: 0.2, retryable: |err| err.is_retryable(), })这里三个配置各有讲究max_attempts包括首次执行在内最多尝试5次。超过后任务标记为Failed触发DAG下游的失败分支。backoff指数退避第n次重试前等待时间按base * 2^n递增上限60秒。jitter_ratio在等待时间上随机上下浮动20%。不用觉得这个参数多余真实系统里多个任务同时失败后统一重试没抖动会造成重试风暴有了抖动可以让重试请求在时间上均匀散开。超时控制用了tokio的timeout包一层。let result tokio::time::timeout( Duration::from_secs(spec.timeout_secs), task.run(ctx), ).await;这里有个特别容易踩的坑tokio::time::timeout返回的Err(Elapsed)只代表没在时间内拿到结果不代表任务真的被取消。很多新手在这里误以为超时后任务就停了。实际上如果任务内部自己在做循环、占着CPU或者持有锁它还会在后台继续跑。所以我在ruflo的文档里反复强调任务代码里要么用select!循环监听取消信号要么在每次迭代里检查ctx.is_cancelled()再决定是否提前退出。单靠外层timeout只能保证调度器不等待不能保证系统资源不泄漏。3. 从零搭建一个可运行的管线实战演示3.1 环境准备与项目接入这部分直接照抄就能跑通。前提是已经装了Rust工具链当前MSRV是1.78低于这个版本编译会报错。cargo new ruflo_demo cd ruflo_demo cargo add ruflo tokio --features ruflo/sqlitesqlite这个feature用来开启状态持久化后端不开启也可以ruflo默认走内存状态进程退出后任务状态丢失。在demo阶段可以不开生产环境建议一定开。然后在Cargo.toml里加上[dependencies] ruflo { version 0.3, features [sqlite, cron] } tokio { version 1, features [full] } anyhow 1 serde { version 1, features [derive] }我demo里选的数据库是SQLite理由很朴素单机场景下任务状态量级在几千到几万条之间SQLite一个文件搞定不需要额外起一个数据库服务而且它天然支持事务状态更新时可以原子地连日志一起写进去。3.2 数据采集任务的编写我拿一个最小可用的日志采集清洗告警管线做演示。这套流程在有业务背景的读者看来可能太简单但麻雀虽小五脏俱全我们重点看的是如何描述任务之间的依赖。先写采集任务#[derive(Clone)] struct IngestTask { source: String, } #[async_trait] impl Task for IngestTask { fn name(self) - str { ingest } async fn run(self, ctx: Context) - ResultTaskOutput, TaskError { let body reqwest::get(self.source).await?.text().await?; let path format!(data/raw/{}, ctx.run_id()); tokio::fs::write(path, body).await?; ctx.set_state(raw_path, path)?; Ok(TaskOutput::default()) } }写清洗任务#[derive(Clone)] struct CleanTask; #[async_trait] impl Task for CleanTask { fn name(self) - str { clean } async fn run(self, ctx: Context) - ResultTaskOutput, TaskError { let raw_path ctx.get_state::String(raw_path) .ok_or_else(|| TaskError::expired(上游未设置 raw_path))?; let raw tokio::fs::read_to_string(raw_path).await?; let cleaned raw .lines() .filter(|line| !line.trim().is_empty()) .collect::Vec_() .join(\n); let clean_path format!(data/clean/{}.txt, ctx.run_id()); tokio::fs::write(clean_path, cleaned).await?; Ok(TaskOutput::default()) } }注意清洗任务通过ctx.get_state从上游拿数据路径而不是自己在任务内部硬编码全局路径。这背后是一套上下文数据传递机制类似于参数服务器上游任务的set_state写入的数据会被安全地隔离在本次运行命名空间里不同批量任务之间的状态不会串味。这是个非常省心的设计因为日志管道经常要按小时跑上一个小时和下一个小时用的是同一个Task对象但上下文必须各自独立。3.3 将数据加工和告警串成流水线现在注册任务并声明依赖let ingest IngestTask { source: https://api.example.com/logs.into() }; let clean CleanTask; let aggregate AggregateTask; let alert AlertTask { threshold: 100 }; let mut dag Dag::new(); dag.add_node(TaskSpec::new(ingest).task(ingest)); dag.add_node(TaskSpec::new(clean).task(clean)); dag.add_node(TaskSpec::new(aggregate).task(aggregate)); dag.add_node(TaskSpec::new(alert).task(alert)); dag.add_edge(ingest, clean)?; dag.add_edge(clean, aggregate)?; dag.add_edge(aggregate, alert)?; let state SqliteStateBackend::open(ruflo.db).await?; let mut scheduler Scheduler::new(dag, state)?; scheduler.run().await?;这里可以看到整个流程的骨架先注册节点再连边再交给调度器跑。连边的顺序就是行文顺序谁先谁后一目了然。跑完之后我可以从state里查询每个任务的运行状态for task_name in [ingest, clean, aggregate, alert] { let status state.get_status(run_2025xxxx, task_name).await?; println!({task_name}: {:?}, status); }输出大致是ingest: Success clean: Success aggregate: Success alert: Success如果某个环节失败调度器会按配置好的重试策略自动重跑达到上限后标记失败。下游任务不会傻等会直接进入Skipped状态。这个行为在管线上非常重要比如alert任务发现上游数据质量不过关它可以跳过本次告警等下一个批次再跑而不是把一批不完整的数据硬发出去。4. 用一套规则管理重试、幂等和中间状态4.1 任务状态机的设计任务状态机是ruflo里最值得细看的部分因为调度器的行为完全由状态迁移驱动。一开始我图省事只设计了三个状态Pending、Running、Done。上线第三天就把状态机改成了下面的六个状态Pending - 初始状态任务等待入度归零 Running - 任务正在执行 Success - 任务执行成功输出可被下游消费 Failed - 任务彻底失败重试次数用完下游被阻断 Skipped - 因为上游失败或条件不满足本任务不执行 Cancelled - 调度器被停止或外部主动取消为什么需要Skipped和Cancelled因为实际运行中DAG往往不是纯粹的全跑结构还有条件分支。比如质量检查不过关就直接跳过聚合不给下游发数据。如果没有一个明确的跳过状态调度器就分不清这个任务没跑和这个任务不该跑的区别下游如果依赖任务结果做判断就会出问题。状态迁移规则只有三行但每行都经过深思只有Pending能被调进Running。只有Running能变成Success、Failed或Cancelled。一个节点成为Skipped当且仅当它所有上游都已完成但至少一个上游是Failed或Skipped。4.2 幂等策略防止重复执行任何带重试的系统都必须面对重复执行的问题——网络断了重试、进程崩溃恢复重试、调度器重启重试都有可能让同一个任务在同一批次里被跑两次。幂等是绕不开的坎。ruflo里的幂等思路分三层第一层调度层面的幂等。每个任务在启动前先往状态后端写入一条lease记录包含任务名、运行批次ID、节点ID、过期时间。调度器重启后如果发现某条lease的拥有者ID和当前进程ID不匹配且任务状态还是Running它不会直接接管执行而是先检查执行超时超时了才把任务重新放回Pending并写入新lease。这个机制防止了同一批次同一任务被两个调度进程同时执行的经典问题。第二层任务代码层面的幂等。ruflo在TaskContext里提供了一个idempotency_key任务每次执行拿到同一个key可以把它作为数据库表里的唯一键或对象存储里的写路径来保证数据只被写一次。let out_path format!(data/out/{}, ctx.idempotency_key()); // 如果 out_path 已存在可以选择直接复用第三层输出校验。任务执行完成后可以注册一个OutputValidator校验输出是否合法。如果在重试过程中发现输出文件字节数、行数或者校验和不对任务会抛TaskError::InvalidOutput调度器据此决定重试还是终止。这层机制对数据管线的价值极其明显——它把任务跑完了和任务跑对了区分开来。4.3 状态持久化与断点恢复进程难免崩溃任务状态必须能持久化。我的实现里SQLite后端会把任务状态、上下文数据、DAG结构、调度日志全部写进一个库文件。表结构大略如下CREATE TABLE task_runs ( run_id TEXT NOT NULL, task_name TEXT NOT NULL, status TEXT NOT NULL, attempt INTEGER NOT NULL, started_at INTEGER, finished_at INTEGER, idempotency_key TEXT NOT NULL, PRIMARY KEY (run_id, task_name, attempt) ); CREATE TABLE task_state ( run_id TEXT NOT NULL, task_name TEXT NOT NULL, key TEXT NOT NULL, value_json TEXT NOT NULL, PRIMARY KEY (run_id, task_name, key) );崩溃恢复的逻辑是这样的进程起来之后先扫描task_runs表把所有状态为Running但finished_at为空且超过超时的记录找出来按上面提到的lease规则决定是重新排队还是标记失败。这个扫描动作非常快几千条记录毫秒级完成。有一点要提醒千万不要把所有任务状态都往内存里塞然后只在进程优雅退出时写一次磁盘。进程崩溃时内存里的状态直接没了下游任务就傻掉了。ruflo的做法是每个任务在状态迁移时立刻写库虽然写库会让性能打折扣但换来的是任何一个中间状态都能被恢复的确定性。5. 性能表现和边界取舍哪些场景不适合5.1 几组典型场景的实测数据我在一台1核1G的云服务器和本地MBP上分别跑了几组基准测试。测试用的是无状态sleep任务模拟任务本身不占资源纯看调度开销的场景。场景任务数并发上限总耗时调度器CPU占比单链线性DAG1001约0.4s3%扇出DAG1对10010116约0.8s5%随机DAG200节点无外部IO20016约1.1s8%随机DAG10000节点无外部IO10000256约4.2s20%注意第二行扇出DAG总耗时比线性DAG长不是调度慢而是每个任务内部sleep了1毫秒所以并发16时100个任务要7波左右跑完。这个数据想说明的是在绝大多数轻量自动化场景下调度器CPU开销完全不是瓶颈瓶颈只会出现在任务自身以及任务间数据移动上。我还测了带SQLite持久化后端的场景。每任务状态落库单机顺序执行1000个任务总耗时比纯内存模式高约40%但换来的是崩溃不丢状态。如果你的任务本身就秒级起步这个差距可以忽略如果你的任务是微秒级纯计算那确实会感受到落库的开销这时候可以关闭持久化或改用批量提交模式。5.2 ruflo的边界与不能代替的东西有段时间我被想用ruflo做所有事冲昏了头冷静下来后我给自己列了一个别用它做清单跨服务长流程业务编排比如订单状态的流转要持续几天、中途会等待用户确认这种场景需要真正的持久化工作流引擎Temporal之类ruflo的跑完一批就结束模型不适用。需要精确一性或事务语义的跨进程协调即使有lease和幂等键ruflo也不能保证两个不同服务之间的操作只发生一次。分布式事务的问题不要丢给工作流引擎解决。复杂数据血缘和元数据管理ruflo只知道任务依赖关系不追踪某个字段从哪个文件哪一行来。要做企业级数据血缘还是要上专门的数据治理平台。非常重的MapReduce式计算任务内部自己决定怎么并行ruflo只负责编排。如果单个任务要起300个线程拉数据那是任务实现的问题不是调度器该干涉的。坦率地说ruflo的理想适用面是单机小集群上的批处理自动化再往上走它的边界会被清晰突破。我自己的用法是把它嵌入到采集代理里而公司的集中式批量计算还是用现成的调度平台。两者各管一摊不冲突。6. 踩过的坑和给二次开发者的建议6.1 一个让我头疼的死锁问题开发过程中最折磨人的bug之一是调度器莫名其妙地卡住——任务没在跑也没有任何报错就是整个DAG不再推进。周末排查半天最终发现是一个非常低级的并发问题。场景是这样的DAG里有10个任务max_concurrency设为2。其中任务A依赖一个动态子任务BB执行时间很长。A启动后进入RunningB尚未被调度同时就绪队列里还有两个任务C、D它们被并发执行了。问题在于ruflo第一版处理动态依赖时要等所有在跑任务都完成才扫描新产生子任务。于是C、D跑完后调度器发现当前没有就绪任务就直接退出了根本没等到A跑完生成B。后来改成了只要存在Running状态的任务就继续等待每有一个任务完成就重新扫描一次DAG才解决。这个bug的根因是在设计状态机时我把当前就绪队列为空和整个DAG跑完画了等号。真实系统中DAG是动态的一个任务可能在运行中又长出新的边来。正确做法是调度循环的退出条件必须是就绪队列为空且Running集合为空两个条件同时满足才退出。6.2 锁、任务日志和公平调度第二个值得记录的坑是任务日志堵塞。一开始我在任务里直接用println!打印日志并发高时终端输出全混在一块排查问题非常痛苦。后来加了tracing库把每个任务的日志写到独立文件但又出现一个问题任务数量多的时候同时打开的文件句柄数把系统默认ulimit打满新任务打开日志文件直接失败。教训是任务日志要按批次轮转而不是按任务无限拆文件。我后来改成run_id一个目录目录内再按任务名轮转同时限制单个日志文件大小超了就滚动。其实这个教训对任何并行系统都适用——任何有限资源文件句柄、内存、连接池都要做配额即使你最开始觉得不可能用满。公平调度也很容易踩坑。ruflo默认的调度优先级是按依赖深度来的深度大的任务先跑。但有些任务虽然深度小却是下游最关键的路径比如数据源连通性检查。所以我在调度器里加了一个priority配置项可以在任务上显式指定优先级dag.add_node(TaskSpec::new(preflight) .task(PreflightTask) .with_priority(10));优先级越高越先进入就绪队列。建议关键路径上的轻量任务都给高优先级避免它们排在一堆重型任务后面耽误整条链路。6.3 后续规划ruflo目前的状态对我个人项目来说够用但以后如果继续扩展我最想做三件事插件化的任务超市让社区可以发布可复用Task包一个简单但有用的Web状态页直接在浏览器里看DAG执行状态而不是只能查SQLite还有多节点worker的正式支持现在虽然能通过共享状态后端实现抢占调度但距离开箱即用的分布式执行还有不小的距离。如果你准备在项目里参考或二次开发ruflo我最想强调的三条建议是长期维护的项目从一开始就引入SQLite持久化。哪怕当前部署只需要内存状态也要把状态后端抽象成接口否则后面想加持久化要动大面积代码。充分测试进程崩溃重启恢复这个路径。给DAG加几百个随机节点随机kill进程然后看状态是否还能恢复。这个测试能暴露出比任何单元测试都多的竞态问题。把cancel safety当成一等公民。Rust异步里很多库方法不是cancel safe的timeout或者外部取消时可能破坏内部状态。ruflo自身处理了这个问题但你写任务的时候也要小心——在run()里用select!处理取消信号不要指望外部能替你清理。给这个小项目持续投入的时间不算少但在实际边缘设备上稳稳当当跑了大半年也没出过一次错。它让我意识到在很多自动化场景里够用就好比功能全面更接近工程本质。
返回列表