FEATURED · 精选文章

ruflo:轻量级本地数据管道编排工具,用YAML定义DAG实现可复现批处理

发布时间 / 2026/9/9 4:46:11
来源 / 创域科博编辑部
栏目 / 资讯中心
ruflo:轻量级本地数据管道编排工具,用YAML定义DAG实现可复现批处理 说实话最初冒出要做“ruflo”这个念头纯粹是被自己那堆乱七八糟的批处理脚本逼疯了。当时我正在同时维护三个数据清洗任务每个任务都由三到四个互不相关的 Python 脚本拼接而成中间靠手工改文件名传递中间结果。某次凌晨重新跑全量流程时因为少改了一个日期参数导致下游消费了三天前的旧数据等发现时已经晚了。那一刻我意识到问题不在我“不够小心”而在于整个流程缺少一种可以描述、可以执行、可以验证的确定性载体。于是 ruflo 诞生了。它本质上是一款极简的本地数据管道编排工具核心思路是用一个 YAML 文件定义整条数据处理链路把每个处理步骤封装成独立的节点节点与节点之间通过约定好的中间文件传递数据。它不依赖任何外部调度系统不搞容器化不引入消息队列甚至不需要常驻服务——一条命令跑完整个 DAG。正是因为足够简单它解决了我最痛的那个问题任何一条流水线无论隔多久回来看都能用同一个命令、同一份配置得到完全一致的结果。这里先说明一点ruflo 并没有重新发明什么轮子。市面上类似 Airflow、Dagster、Prefect 这类重型调度框架非常多但它们对个人项目和中小团队来说往往过于笨重——要装数据库、要起 Web 服务、要学一整套生态。ruflo 的定位完全不同它面向的是“本地批处理任务”这个被大厂框架忽视的角落你有几个脚本要串联你有定时生成的报表要统一管理你想让一次数据处理过程变得可复现、可追溯、可断点续跑。如果你的场景恰好是这个那这篇拆解应该对你有用。1. 项目背景从脚本堆到管道我到底在解决什么问题1.1 原始的“脚本接线”方案有多脆弱在没有 ruflo 之前我的典型工作流长这样采集脚本 output/raw_data.csv清洗脚本读取 raw_data.csv 再输出 cleaned_data.csv聚合脚本再读 cleaned_data.csv 最终生成 report.json。听起来没什么问题是吧问题出在“连接”上。第一个坑是隐式依赖。清洗脚本假定 raw_data.csv 一定存在且结构正确但这个假设从未被显式声明过。只要上游脚本某天改了输出列名或者因为网络问题产出了一个空文件下游脚本就会在某个莫名的报错中崩溃而排查这个报错往往需要打开好几个文件对比字段。第二个坑是重跑成本。当聚合逻辑更新后我必须手动确认是否需要重跑清洗步骤。如果能重跑我得先找到对应的输入文件再确认参数没变再小心翼翼地执行。这个流程就像手工操作一条生产线一旦中间任何一个环节被遗忘或错配结果就不可信了。第三个坑也是最隐蔽的参数漂移。同一个清洗脚本周一用--date2024-06-03跑过周三又因为临时修数据用--date2024-05-27跑过两条中间产物同时留在磁盘上文件名却是一样的。等真正要出报表的时候你根本不知道当时调用的到底是哪个参数组合。1.2 ruflo 想做成什么样ruflo 的目标很直白把数据处理的每一步变成一个有名字、有输入、有输出、有校验的处理单元并让这些单元之间的连接关系被显式地写下来。具体拆成三个能力显式声明整条流水线节点、依赖、参数、输入输出路径全部写在一个 YAML 文件里随代码仓库一起版本化管理。任何人在任何时间 checkout 到某个 commit看到的流水线定义都和当时执行的一模一样。确定性执行每次执行都基于配置文件中固定的参数和文件路径不允许隐式读取当前目录下的“碰巧存在”的文件。目标只有一个——同一份配置任何时候执行结果一致。断点恢复与结果校验如果第 3 步失败修复后可以只重跑第 3 步及其后续节点不必从头再来。而每个节点完成后会计算输出文件的 SHA256 校验和并存到执行记录里用于确认结果没有被意外改动。这样设计的原因非常朴素我不想再被“我上次怎么跑出这个结果的”这个问题折磨。只要配置在、输入在执行一次 ruflo 就能完整复现整个流程连中间产物的哈希值都能对上。2. 架构与数据模型4 个核心概念不多不少ruflo 的架构设计遵循了一个原则能用简单的数据结构讲清楚的事绝不引入复杂的服务。整条流水线在 ruflo 中被抽象成四个概念工程Project、节点Node、依赖Dependency产物Artifact。2.1 Project一切配置的根容器一个 Project 对应一个 YAML 文件通常放在项目的根目录下文件名是ruflo.yaml。它定义了项目名、默认工作目录、节点列表和全局环境变量。version: 1.0 project: demo-report workspace: ./work artifacts_dir: ./artifacts env: PYTHONUNBUFFERED: 1 nodes: - id: fetch ...这里的关键是workspace和artifacts_dir的分离。workspace是执行过程中的临时目录跑完可以清空artifacts_dir是最终产物目录会长期保留。这样分开的好处是重跑时不会把历史产物弄丢而临时文件又不会堆积。2.2 Node一个可执行的处理单元每个 Node 代表流水线中的一步必须声明三样东西怎么执行、输入是什么、期望输出什么。nodes: - id: clean_data type: command command: python scripts/clean.py --input {input.raw_data} --output {output.clean_data} inputs: raw_data: artifact: fetch.raw_data outputs: clean_data: path: cleaned/data.csv从这段配置可以看到几个关键设计command是实际执行的命令行里面用{input.xxx}和{output.xxx}做的变量插值执行前 ruflo 会自动替换成真实路径。inputs声明了本节点依赖的产物并且必须写成“哪个节点的哪个输出”这种形式。这就把隐式文件依赖变成了显式的节点间连接。outputs声明本节点会产生哪些文件路径是相对于项目根目录的。为什么不在命令行里直接写死路径因为路径一旦写死流水线就无法移动换个目录执行就挂了。通过变量插值ruflo 可以在运行时决定工作目录前缀让同一条流水线在不同机器、不同工作目录下都能跑。2.3 Dependency用 DAG 建立执行顺序节点之间的依赖关系完全由inputs声明推导。ruflo 每次执行前会构建一张有向无环图节点是顶点inputs中引用的其他节点产物是边。比如clean_data声明了fetch.raw_data作为输入那么就有了一条fetch→clean_data的边。图构建完成后ruflo 会做两件事拓扑排序确定节点的执行顺序。环检测如果配置里出现了循环依赖比如 A 依赖 BB 又依赖 A直接报错拒绝执行。依赖关系推导的一个额外好处是你永远不会因为忘记写“先跑 fetch 再跑 clean_data”而得到空输入。因为 clean_data 的输入已经绑定了 fetch 的产物fetch 没跑或者失败了clean_data 根本不会被调度。2.4 Artifact一切中间文件的第一等公民在 ruflo 里任何节点输出的文件都被称为 Artifact。它不像普通文件那样仅仅是磁盘上的字节而是带着元数据的对象——产出它的节点、产出时间、当时的参数、SHA256 校验和。这些元数据被记录在.ruflo/artifacts.json中是 ruflo 实现可复现性的关键。每次节点执行完成ruflo 会扫描outputs段中声明的所有文件计算哈希并记录。后续如果要校验结果只需要对比当前文件的哈希和记录中的哈希即可。还有一个细节Artifact 名称的生成规则是node_id.output_key比如fetch.raw_data。这种命名方式让你在配置里引用产物时一目了然——它是哪个节点产出的、在节点内部叫什么。这也是“可读性优先”的体现。3. 工作流编排的完整实现调度器、变量插值与哈希记录有了数据模型之后最核心的问题就是运行时如何把它们串联起来。这一节我会比较详细地拆解 ruflo 调度器的三个关键机制。3.1 调度器如何决定“先跑什么、后跑什么”调度器的输入是一张图输出是一个有序的执行队列。我用的是经典的 Kahn 算法做拓扑排序先把所有入度为 0 的节点入队然后逐个弹出、执行、更新相邻节点入度。def topo_sort(graph): in_degree {node: 0 for node in graph.nodes} adj {node: [] for node in graph.nodes} for node in graph.nodes: for dep in node.dependencies: adj[dep].append(node.id) in_degree[node.id] 1 queue deque([nid for nid, deg in in_degree.items() if deg 0]) order [] while queue: nid queue.popleft() order.append(nid) for neighbor in adj[nid]: in_degree[neighbor] - 1 if in_degree[neighbor] 0: queue.append(neighbor) if len(order) ! len(graph.nodes): raise RuntimeError(cycle detected in dependency graph) return order解释一下这段逻辑在 ruflo 里的实际意义graph.nodes是配置里所有节点的集合每个节点在加载时已经把自己的inputs引用解析成了“依赖的节点 ID 列表”。adj记录了从“被依赖节点”指向“依赖节点”的边。为什么要这个方向因为拓扑排序的本质是只有当一个节点的所有前驱都执行完它才能被调度。最后一步的环检测非常重要一个生产级的 DAG 调度器绝对不能容忍死循环。如果 ruflo 检查到环会直接终止执行并打印出环涉及的节点链条方便定位配置错误。注意当前版本的 ruflo 并不做并行调度所有节点都是串行执行的。这一点是刻意为之并行虽然能缩短总耗时但会给日志跟踪、中间产物管理和错误定位带来复杂度。对于本地批处理任务来说串行执行的十几分钟通常是可以接受的。3.2 变量插值配置不能出现“裸路径”在定义 Node 的command时我特意禁用了裸路径。比如python scripts/clean.py --input data.csv这种写法是不允许的因为data.csv的路径没有来源。正确的写法是command: python scripts/clean.py --input {input.raw_data} --output {output.clean_data}执行时ruflo 会构建一个变量表其中的键值来自两个维度input.key对应inputs中引用的上游产物路径。output.key对应outputs中声明的目标路径注意这里的路径是 ruflo 自动生成的你不需要手动创建父目录。这个设计带来几个直接好处路径统一无论你从哪个目录调用ruflo run所有节点拿到的输入输出路径都是绝对的不会因为“当前工作目录不同”而产生行为差异。参数可追溯命令中所有可变部分都来自配置文件执行时 ruflo 会把最终的命令完整记录到执行日志里。以后想看某一步到底怎么跑的查日志即可。父目录自动创建如果{output.clean_data}指向的子目录还不存在ruflo 会在执行前自动创建。这避免了很多脚本因为没有先mkdir -p而挂掉的窘境。3.3 执行记录与断点续跑每条流水线跑完后.ruflo/run_history/下会生成一个以时间戳命名的 JSON 文件记录每个节点的执行状态、退出码、开始结束时间、产物哈希。这个文件就是 ruflo 的“飞行记录仪”。断点续跑的机制依赖这套记录。当执行失败时ruflo 会把已成功节点的产物哈希记录保留下来。修复后如果执行ruflo run --resume调度器会做以下判断如果节点在上一次执行中成功且它的所有输入产物哈希与记录一致则跳过该节点直接将其产物标记为“已就绪”。如果节点未执行或执行失败则重新调度。这个判断逻辑非常实用。试想一下一个 8 节点的流水线跑到了第 7 步挂掉如果没有断点续跑重跑一次就要再等前 6 个节点跑完。有了断点续跑只需要几秒钟就能直接跳到失败节点继续。当然断点续跑有一个前提——输入不能变。所以我在--resume模式下额外做了一步校验对比上游产物当前哈希与上次记录的哈希不一致则放弃跳过强制重跑。宁可多花时间也不要用错误的中间产物生成错误的结果。4. 实战拆解从日志文件到周报 JSON 的完整管道理论讲再多不如直接看一个能跑的例子。下面我拿一个非常典型的需求来演示 ruflo 的完整配置和运行过程输入是一周的访问日志经过解析、清洗、聚合三步最终产出一个周报 JSON。4.1 第一步定义节点流水线一共四个节点fetch_logs从本地归档目录拷贝一周的日志文件到工作目录。parse_logs把 Nginx 日志解析成结构化 CSV。filter_abnormal清洗掉状态码异常或请求时长过长的记录。weekly_report按 URL 聚合访问量生成周报 JSON。对应的ruflo.yaml如下version: 1.0 project: access-log-weekly workspace: ./work artifacts_dir: ./artifacts nodes: - id: fetch_logs type: command command: bash scripts/fetch_logs.sh {output.log_dir} outputs: log_dir: path: logs/raw - id: parse_logs type: command command: python scripts/parse_logs.py --input {input.log_dir} --output {output.parsed_csv} inputs: log_dir: artifact: fetch_logs.log_dir outputs: parsed_csv: path: parsed/access_logs.csv - id: filter_abnormal type: command command: python scripts/filter_abnormal.py --input {input.parsed_csv} --output {output.filtered_csv} inputs: parsed_csv: artifact: parse_logs.parsed_csv outputs: filtered_csv: path: filtered/access_logs_clean.csv - id: weekly_report type: command command: python scripts/weekly_report.py --input {input.filtered_csv} --output {output.report_json} inputs: filtered_csv: artifact: filter_abnormal.filtered_csv outputs: report_json: path: reports/weekly_report.json从这份配置能清晰看到流水线的形状fetch_logs→parse_logs→filter_abnormal→weekly_report。每个节点都只关心自己的输入输出完全不需要了解其他节点的内部实现。4.2 第二步跑起来执行命令只有一条ruflo run执行过程中终端会打印每个节点的状态流转。我用一个简化的表格说明每个阶段看到的信息节点状态说明fetch_logsRUNNING → DONE日志文件被拷贝到 work/logs/rawparse_logsRUNNING → DONE生成了 8.3MB 的 access_logs.csvfilter_abnormalRUNNING → DONE过滤掉 4.2% 的记录weekly_reportRUNNING → DONE生成 weekly_report.json包含 137 个 URL 条目跑完后在artifacts/reports/下就能看到最终的 JSON 文件。同时.ruflo/artifacts.json会记录每个产物的绝对路径和 SHA256 值。4.3 第三步改造一个真实存在的“不可复现”场景这里我要重点说说参数问题。假设跑完第一次周报后产品经理说“上周数据里可能有爬虫流量能不能把 user-agent 包含 spider 的请求过滤掉”按照原来的脚本堆方式我会去改filter_abnormal.py里的过滤逻辑然后重新跑一遍全流程。可问题来了我到底要不要重新跑parse_logs如果 parse_logs 不会因为 filter 逻辑改变而变化理论上可以跳过。但我怎么确认在 ruflo 中这个问题被彻底化解了。我只需要改filter_abnormal节点的命令参数command: python scripts/filter_abnormal.py --input {input.parsed_csv} --output {output.filtered_csv} --drop-spider然后执行ruflo run --resume调度器会检查各节点输入的哈希发现parse_logs.parsed_csv的哈希和上次一样于是跳过fetch_logs和parse_logs直接从filter_abnormal开始重跑。如果我还想强制全量重跑只需要加--force参数。这就是 ruflo 给我的最大安全感改哪里就从哪里重跑没改的地方绝不会因为“碰巧重跑”而引入新的不确定性。5. 我踩过的坑与针对性设计ruflo 不是一次就写对的开发过程中踩了不少坑。这一节分享几个最有代表性的问题和最终的解决方案同时也是对前面设计逻辑的补充说明。5.1 坑一Shell 通配符和引号被错误展开最早的版本里我直接调用了subprocess.run(command, shellTrue)来执行命令。这带来一个问题如果配置里写了python scripts/parse.py --input logs/*.logshell 会自己展开通配符。看起来很方便但实际上很容易踩坑——当工作目录不对或者 logs 目录下文件数量为 0 时通配符会直接变成字面量脚本拿到一个不存在的路径报错信息还特别难懂。解决方案ruflo 默认使用参数列表模式shellFalse命令字符串会在内部按 shell 规则做安全的词法分割但不做通配符展开。如果你确实需要通配符需要显式在命令里包一层bash -c ...并自己负责语义。这个取舍牺牲了一点便利换来了行为的确定性。5.2 坑二输出文件未声明导致“跑完不知道产物在哪”早期版本的 output 声明是可选的节点可以只写 command不声明 outputs。结果就是脚本确实跑出了文件但 ruflo 不知道这个文件叫什么、在哪。后续节点想引用时又只能靠猜路径退化成了原来的隐式依赖。解决方案强制每个节点必须声明至少一个输出。虽然刚开始写配置时会觉得烦但习惯后你会发现每个节点的“契约”变得非常清晰输出就是这 1-2 个文件别的都不归 ruflo 管。这条规则让整条流水线的文件流变得完全透明。5.3 坑三校验和拖慢整体执行一开始我给所有产物都计算完整文件的 SHA256。大文件还好遇到几十 GB 的数据文件时光算哈希就要花几分钟严重拖慢了整条流水线。解决方案引入了一个开关artifacts.checksum: fast启用后只对文件的前 1MB 和末尾 1MB 做分段哈希。这样既能捕捉到绝大多数意外修改又不会因为全文件扫描而浪费时间。对于严格性要求更高的场景可以设置checksum: full恢复全文件校验。校验模式性能安全性适用场景fast高中日常跑批文件较大full低高归档、审计、复盘场景5.4 坑四失败重跑时的工作目录污染某节点失败后遗留了不完整的输出文件修复配置后再次运行脚本读到了这个半截文件导致再次失败。这类问题在数据处理场景特别常见。解决方案执行节点的命令前ruflo 会先清理该节点声明的输出文件如果有历史残留的话。这个设计可能很多人觉得多余——但它在实践中救了我很多次。每次执行都是从干净状态开始杜绝了“上次运行的残骸影响本次运行”的可能性。6. 运行性能与稳定性实测一个工具好不好用最终要看在真实负载下的表现。我给 ruflo 做了几组简单但有效的测试结论值得一说。6.1 一个中等规模批处理任务的实测数据测试环境是普通的 MacBook ProApple M1 Pro16GB 内存。流水线结构与第 4 节类似但数据量放大到约 15GB 的原始日志包含 1.2 亿行记录。整体表现如下阶段耗时说明配置解析 图构建 100msYAML 加载和依赖图构建开销可忽略fetch_logs3m 12s主要是磁盘 IO 拷贝parse_logs12m 48sCPU 密集型的逐行解析filter_abnormal4m 05s内存中过滤涉及大量字符串匹配weekly_report1m 58s分组聚合输出 2.3MB JSON调度器自身开销~80ms几乎可以忽略可以看到调度器本身的开销完全不是瓶颈。整体耗时基本被业务脚本吃掉这也是我选择“轻量调度外部脚本”路线的一个重要佐证对于数据处理管道性能的核心在于节点内的算法和 IO而非节点间的调度。6.2 失败恢复的实测我刻意在第 4 步写了一个会在某类数据上崩溃的脚本测试断点续跑首次执行到第 4 步失败耗时约 8 分钟。修复脚本后执行ruflo run --resume调度器在 500ms 内完成哈希校验跳过前 3 个节点直接从第 4 步开始。后续执行耗时 2 分钟总恢复时间仅为完整重跑的 1/4 左右。这个收益在节点越多、单节点耗时越长的时候越明显。如果你的管道有 20 个节点中途失败后用断点续跑可以节省 90% 以上的重跑时间。6.3 稳定性设计进程退出码与信号处理ruflo 对被调度节点的退出码做了严格处理退出码为 0节点成功记录哈希。退出码非 0节点失败立即终止流水线并输出该节点的标准错误日志位置。被信号杀死如 Ctrl-C标记节点为“被杀”保留已产生的输出文件但不会记录成功状态。这样设计后无论是正常失败还是手动中断整个流水线都处于一个可预测的状态。中断后跑--resume被中断的节点会重跑因为它的输出哈希没有被记录。7. 从 ruflo 中学到的那些“跟工具无关”的东西做了这个项目之后我对“工具”这件事本身有了很多新的认知。这些认知超出了代码范畴影响了我的工作习惯分享出来也许对你有参考价值。7.1 显式总是好过隐式即便写起来麻烦一点回顾最初踩坑的根源几乎全都指向“隐式假设”。脚本假设输入文件存在假设列名不变假设上次的运行结果无关紧要。ruflo 通过强制声明输入输出、强制校验哈希把这些假设一个个暴露到阳光下。这并不是 ruflo 的专利而是很多优秀工具的共同哲学。你在做任务编排、接口设计、甚至写函数参数时如果能做到显式大概率能避免未来 80% 的“莫名其妙出错”。7.2 轻量意味着更容易被坚持使用坦白说我也用过 Airflow。但在我个人项目的场景里它带来的维护成本远大于收益——数据库、调度器、权限、Docker 镜像……每一样都需要花时间打理。ruflo 只有一条命令、一个 YAML 文件没有服务端、没有数据库。这种极低的入场成本是它能被我在日常工作中持续使用的重要原因。有时我们把问题复杂化了。你想跑一个定时批处理真的需要一整套分布式调度平台吗大多数情况下不需要。你需要的只是“把几个脚本按顺序跑完并且出了问题容易排查”而已。7.3 结果可复现是一切数据工作的底气数据工作最怕的不是报错而是“无声的错误”——流程跑完了结果却是错的而你毫无察觉。ruflo 的哈希校验和严格的参数记录至少保证了“如果结果错了我能回到产生它的那个时刻看看当时到底发生了什么”。这种底气让跟业务方沟通时少了很多猜疑。当对方问“这个数字是怎么算出来的”时我可以直接翻出当时的配置、当时每个节点的执行命令、当时产物的哈希值而不是支支吾吾地说“好像是用那个脚本跑的”。7.4 接下来我准备怎么扩展 ruflo当前版本已经基本满足了我的日常批处理需求但后续还有几个想法在酝酿中支持节点级超时控制防止某个脚本卡死导致整条流水线挂起。加入一个简单的ruflo viz命令把依赖图渲染成 SVG方便给别人讲解流水线结构。增加一个“dry-run”预览模式显示本次将要执行/跳过的节点以及涉及的文件路径跑批前先确认一遍。如果你也经常被一堆脚本之间的隐形依赖搞得焦头烂额不妨试试类似的思路——先别急着上重型框架写一个 500 行的调度脚本可能都比在混乱中继续堆脚本强。ruflo 对我来说就是这个 500 行的调度脚本只不过它已经长得比我预期的大得多了。
RELATED — 相关阅读

相关资讯

LATEST — 最新资讯

最新发布

TODAY — 本日精选

新闻

WEEKLY — 本周精选

新闻

MONTHLY — 本月精选

新闻