ruflo:一个轻量级工作流引擎,让单机cron任务编排不再痛苦

ruflo:一个轻量级工作流引擎,让单机cron任务编排不再痛苦 深夜两点我盯着屏幕上的“今日无数据”四个字心里草泥马狂奔。原因是上周手动跑数的时候有个 Python 脚本抛了个异常我没有注意它后面依赖的三个报表任务全都静默跳过了。这种问题不是第一次了散落在 crontab 里的十几个任务之间没有任何依赖关系cron 只负责“到点启动”至于上一个任务有没有成功、数据有没有就绪它一概不管。也就是从那天起我认真找了一圈“能在单机上把任务编排起来”的工具最终写了个叫ruflo的轻量工作流引擎并且把它用进了我的每日数据处理里。这篇东西就当是项目记录加经验分享已经把坑都踩平了你可以直接照抄。1. 为什么我需要一个“中间地带”的调度工具1.1 散装脚本的痛点我忍了很久在引入 ruflo 之前我的日常是这样的crontab 里躺着十来条记录有些是 PHP 脚本有些是 Python 脚本还有一些是 shell 一行流。每天早上它们按照固定的时间点启动各自干各自的活。表面上看好像挺有秩序实际上全是手工作业A 脚本生成中间文件B 脚本在五分钟之后跑因为大家默认 A 五分钟内能跑完。某天服务器负载一高A 跑了十分钟B 读文件直接读到半个空文件整个链路的数据全错了。这只是其中一个问题。更麻烦的是失败处理。cron 的失败就是一条MAILTO邮件多数时候我根本不去看那台机器的 root 邮箱。任务的退出码、耗时、运行日志系统完全不给你归档想回查“上周三凌晨三点到底发生了什么”只有 journal 日志里一串没头没尾的记录。我需要的不是一个更复杂的调度平台而是任务之间能用第一性原生的方式画出依赖关系任务失败了能自动重试、能留痕我能随时打开它看到哪些任务成功、哪些失败、失败在哪一步。1.2 cron 为什么不够用不是说要废除 cron而是 cron 解决的是“周期性定时启动”它压根不是为工作流设计的。你要用它搭一个多任务依赖流程只能把顺序写进一个大 shell 脚本里./step1 ./step2 ./step3。这么干有两个问题任务之间完全没法并行。step1 和 step2 如果没有依赖关系理论上应该同时跑但在一个 shell 脚本里只能串行等待。任何一步失败了整条链停摆。后面任务不会跑但也没有任何机制告诉你“因为第一步失败所以后面的都取消了”更不会在你修好之后把后面补跑一遍。失败重试基本靠手。脚本挂了你去看日志改 bug然后自己重新执行整个管道。也许中间步骤是幂等的也许不是一切全凭记忆力。你当然可以在 shell 里加if、加循环、加 trap写得像一个迷你调度引擎但那就和拿着螺丝刀当凿子用一样能用但是别扭极了。重复代码多了之后维护成本直线上升。1.3 大厂的分布式调度小团队真的用不起我也认真评估过主流的开源工作流引擎。Airflow 功能强大有 Web UI、有调度器、有 executor、有丰富的 operator 生态。但代价是什么你要在服务器上维护一个 Python 环境要起 scheduler、webserver、worker 至少三个进程要建元数据库还要忍受 DAG 加载那套繁琐的 Python 类写法。Argo Workflows 更重一套 Kubernetes 集群跑在那里就为了每天跑三五个数据脚本实在说不过去。我估算过一次为了跑十几个脚本任务引入这套体系后的运维成本已经超过了任务本身的开发成本。这是典型的“为了补一个洞挖了一座山”。中间一定有某种产品形态单机可部署、配置简单、支持依赖编排和失败重试而这正是 ruflo 想填的位置。2. ruflo 到底是个什么东西核心概念与设计取舍2.1 名字由来ruflo 是个拼起来的名字取 Rust Flow 两个词的缩写。写这个项目的时候我明确了两条要求运行时必须是一个单一二进制文件最好用 Rust 实现性能不那么要紧但部署和分发必须极其简单工作流本身用 YAML 描述不要写代码因为很多日常任务脚本就是 shell 和 Python再让团队去学一套 SDK 真是不必要。ruflo 的发音也好听读起来像中文的“如流”我希望流程像流水一样顺畅不堵、不断地往前赶。2.2 三个核心抽象Workflow、Task、Triggerruflo 的核心只有三个概念没有更多Workflow工作流一个完整的流程定义包含所有任务、它们之间的依赖、触发规则定时、手动、事件、以及全局的重试和通知策略。对应一个 YAML 文件通常一个业务场景就是一个 workflow。Task任务工作流里的最小执行单元。绝大多数任务就是一个命令比如跑一段 Python 脚本、执行一条 SQL、调用一个内部接口。Task 之间通过depends_on声明依赖ruflo 根据这些声明构建一张有向无环图DAG执行时按拓扑排序推进。Trigger触发器决定工作流什么时候启动。当前支持schedulecron 表达式定时触发、manual命令行手动触发常用于补数据、webhookHTTP 请求触发。这三种基本覆盖了日常需求。为什么用depends_on显式声明而不是像 make 那样根据输入输出文件的 mtime 自动推断因为大部分任务操作的是一堆数据库表和接口文件依赖根本表达不了显式声明虽然笨一点但可读性强得多一眼能看到全貌。2.3 为什么能用 Rust 一个文件解决纯从技术角度讲这个项目用别的语言也能做。但 Rust 给我带来了三个非常实际的好处单静态二进制。没有 Python 解释器版本问题没有 Node 依赖目录没有 JVM 启动延迟。编译完 scp 到服务器上就能跑服务器上不需要装任何运行时。我在一台 CentOS 7 和一台 Ubuntu 22 上做了测试同一个二进制文件直接运行没有任何“缺 libc 版本”的问题。子进程管理更可控。每个 Task 最终都会 fork 一个独立子进程来执行命令。Rust 标准库对 child process 的封装虽然简单但配合 tokio 异步运行时可以精确控制超时、读取 stdout/stderr、处理退出信号。最重要的是任务如果在执行过程中变成了僵尸进程我能立刻检测到 timeout 并强制 kill这在 Python 里实现起来要多写不少代码。内存占用低、重启快。常驻内存 25MB 左右调度 1000 个任务时的启动开销控制在十几毫秒这对“一台小机器上跑日常调度”的场景已经绰绰有余。2.4 与同类工具的边界工具部署形态依赖编排失败重试状态可视化最适合的场景cron shell系统自带串行手写手动无三五个固定脚本无所谓顺序ruflo单文件DAGYAML 声明内置指数退避CLI/自绘 Web单机几十个任务需要依赖和重试AirflowPython 环境 调度器/Web/workerDAGPython 代码内置完整 Web UI有专职数据团队任务多到需要整体管理Temporal / Argo分布式集群代码 状态存储内置有平台微服务分布式编排远超日常脚本场景边界很清晰ruflo 不跟 Airflow/Temporal 抢大项目干的是“cron 不够用但引入全套框架又太沉”的中间活。如果你有超过一台机器需要统一调度或者你需要一个协同编辑的工作流平台别用 ruflo直接上正经平台。3. 半小时上手第一个 ruflo 工作流全记录3.1 安装从 GitHub Releases 下载对应平台的压缩包解压出ruflo可执行文件放到/usr/local/bin/下完事。wget https://github.com/yourname/ruflo/releases/download/v0.4.2/ruflo-x86_64-unknown-linux-gnu.tar.gz tar zxvf ruflo-x86_64-unknown-linux-gnu.tar.gz sudo mv ruflo /usr/local/bin/ ruflo version如果你坚持要从源码编译项目是标准 Cargo workspacegit clone https://github.com/yourname/ruflo.git cd ruflo cargo build --release ./target/release/ruflo version编译需要 Rust 1.70 以上版本编译时间三到五分钟优点是你可以改源码缺点是不太值得。实际上大多数用户直接用预编译版本就好这也是我发布 release 不鼓励自己编译的原因。3.2 写一个真正有用的工作流空谈概念没用直接上一个真实例子。假设我每天凌晨两点跑一个数据处理流程分三步拉取前一天的数据、清洗转换、写回数据库。id: daily_data_pipeline name: 每日数据处理管道 description: 拉取交易数据清洗后写入分析库 schedule: 0 2 * * * timezone: Asia/Shanghai max_retries: 2 retry_interval: 30 on_failure: notify: - type: webhook url: https://internal.alerts.example.com/ruflo-fail headers: Authorization: Bearer xxx env: PYTHONUNBUFFERED: 1 APP_ENV: production tasks: - id: fetch_data name: 从源库拉取数据 run: python3 /opt/scripts/fetch_trade_data.py --date {{date}} timeout: 600 max_retries: 1 - id: clean_data name: 清洗与格式转换 run: python3 /opt/scripts/clean_trade_data.py --input /data/raw/{{date}}.csv --output /data/clean/{{date}}.parquet timeout: 900 retry_interval: 60 depends_on: - fetch_data - id: load_to_warehouse name: 写入分析数仓 run: python3 /opt/scripts/load_parquet_to_dw.py --file /data/clean/{{date}}.parquet timeout: 1200 max_retries: 3 depends_on: - clean_data on_success: notify: - type: webhook url: https://internal.alerts.example.com/ruflo-success几个地方需要解释{{date}}是内置模板变量ruflo 在每次运行时根据date参数生成默认为前一天的日期也可以--date 2024-06-01手动指定。这个变量在补数场景里极其常见。每个 Task 可以独立设置timeout、max_retries、retry_interval覆盖全局默认值。比如load_to_warehouse最容易因为数据库锁而失败我给了它更多重试机会。env是全局环境变量Task 里的脚本都可以继承。实际运营中我们常用它统一注入PYTHONUNBUFFERED因为它能强制 Python 不缓冲 stdout让日志实时输出而不是攒到退出才打印。3.3 跑起来验证结果先把工作流文件存为/etc/ruflo/conf.d/daily_pipeline.yaml然后执行# 校验配置文件语法 ruflo validate /etc/ruflo/conf.d/daily_pipeline.yaml # 立即执行一次不受 schedule 约束 ruflo run --workflow daily_demo --date 2024-06-12validate子命令会检查 YAML 是否合法、depends_on是否引用了不存在的 task、整张图是否成环。成环是个非常容易犯的错误特别是后期改依赖时我见过有人把 A 依赖 B、B 依赖 C、C 再依赖 A 的环写出来过。ruflo 会在 validate 阶段就报错不会等运行到一半才卡死。执行过程中终端会实时刷出每个 Task 的状态变化pending、running、success、failed、skipped、canceled。跑完后会输出一个简单的总览Workflow: daily_data_pipeline (2024-06-12) Status: SUCCESS Total time: 42s Task Status Duration fetch_data success 12s clean_data success 25s load_to_warehouse success 5s这里注意一点failed和canceled语义不同。某个任务最终失败导致工作流中止时还没执行的后继任务标记为canceled表示“没跑不是你的错是前面的失败连累了你”。这个区分在排查问题时非常有用。3.4 调度与手动触发配置文件里有schedule字段的话运行ruflo daemon start会进入守护进程模式常驻后台。守护进程启动时会扫描配置文件目录把每个工作流的下一运行时间算好到点自动触发。manual触发无需杀守护进程直接在另一个终端执行ruflo run即可。这里有个工程设计细节每个 workflow 默认不允许同一个运行时刻出现多个实例重叠执行。如果上一个实例还在跑下一个定时点到了ruflo 会跳过本次触发在日志里记一条skipped due to already running。这条策略对数据任务非常关键否则两个实例同时写一行表轻则数据重复重则死锁。4. 生产环境绕不开的三件事重试、并发、状态4.1 失败重试与退避策略日常脚本任务失败的形态千奇百怪但绝大多数时候立刻重试有机会成功因为像数据库连接池打满、外部接口抖动、临时磁盘满这些问题过几十秒可能就恢复了。真正意义上的代码缺陷反而不需要重试重试也救不回来。ruflo 的max_retries和retry_interval支持一个简单的指数退避每次重试的等待时间为retry_interval * 2^n秒n 从 0 开始。比如配置retry_interval: 30第三次重试前的等待时间是 120 秒给下游留足了恢复空间。踩过的一个坑max_retries与retries都设置成很大的值而实际脚本里没有做幂等。比如任务里是“插入一行记录”重跑一次就会插两条。后面我把所有需要重试的任务全部要求做到幂等在数据库表里加唯一键或者利用INSERT ... ON DUPLICATE KEY UPDATE。这条经验我觉得比任何调度功能都重要因为没有幂等保护重试反而成了灾难。4.2 并发控制和任务级资源限制依赖关系决定任务何时启动但在同一时间可能有多个任务满足启动条件它们彼此不互相依赖。ruflo 默认会在一个工作流里并行执行这些任务同时提供一个全局并发上限max_parallelismmax_parallelism: 4这四个并发槽由整个工作流共享。如果fetch_data有 3 个子任务理论上可以同时跑但受限于 4 槽配置最多只能跑 4 个。这个配置的动机很简单小机器上 CPU 核心有限无脑并发会导致任务之间互相抢资源反而比串行更慢。我实测过一个场景10 个 Python 任务同时启动每个都要加载 500MB 内存 8G 的服务器直接 OOM。后来把max_parallelism调到 2问题消失总耗时反而降了因为没有了 swap 抖动。内存大户任务可以在 Task 内配置max_parallelism: 1强制独占。任务执行是在独立子进程里进行的。这一点很重要即使某个脚本malloc了 1GB 内存或者程序崩溃了影响被隔离在这个子进程内ruflo 主进程完全不受污染随时可以读状态、写日志、管理其他任务。4.3 状态持久化和重启恢复实现调度器一个容易忽略的大坑是“状态管理”如果任务跑到一半机器重启了怎么处理简单的方案是把状态保存在 SQLite 里ruflo 用的就是这个。每个 workflow 实例、每个 task 实例的执行状态都落盘到ruflo.db。机器重启后ruflo daemon start会扫描数据库中running状态的任务统一标记为interrupted不会自动重跑它们但会把这次中断完整记录在历史里。为什么不是自动重跑因为我发现自动重跑一个未知状态的任务风险太大也许它已经写完了结果、只是还没来得及把状态更新到数据库就断电了重跑会造成重复写入。所以我选择了“标记为 interrupted人工决定是否补跑”。可靠性保障手段还有一层Task 的run执行前ruflo 会先写一条“即将运行”的日志到存储里任务结束后再记“成功/失败”。这套预写日志write-ahead log的做法保证了任何瞬间掉电状态不会出现“其实跑完了但记录成 running”的矛盾。当然也不能说完全零风险操作系统层的崩溃和磁盘缓存这些不可控因素永远存在但日常排查已经足够用了。5. 八小时实战踩坑排查一个隐蔽的失败5.1 现象任务失败但没有日志有一次我收到告警某天的数据管道挂在了load_to_warehouse这一步。我把这次实例的状态和日志都翻出来发现详情页里一点有用的信息都没有退出码是 1stdout 为空stderr 为空。当时整个人是懵的因为脚本明明在命令行手动跑过毫无问题。5.2 排查链路我按下面这条线一步步查下去整个链路花了快八个小时虽然累但非常值。第一先确认真的是 ruflo 把命令执行错了还是命令本身退出了 1。我手动执行同样的命令行正常完成。所以问题从“脚本 bug”转向“运行环境差异”。第二看工作目录。手动执行时我在/opt/scripts/目录下脚本里用的是相对路径ruflo 创建子进程时默认继承守护进程的工作目录而守护进程的工作目录是/root/。脚本尝试读./config.ini结果找不到文件。但在第一眼看到的 stdout/stderr 里这个错误“消失”了。第三找到真正的原因。为什么 stdout/stderr 会为空因为脚本里有一段except Exception:直接裸吞了异常而且日志写到了一个相对路径的logs/error.log——它自己都因为工作目录不对而没被创建。循环放大到更上层看就是这个任务不仅代码写得糙连错误记录都依赖当前工作目录一环扣一环。最终解决方案在 Task 配置里显式声明cwd- id: load_to_warehouse run: python3 /opt/scripts/load_parquet_to_dw.py ... cwd: /opt/scripts这个字段的含义理解成“执行前先切到这个目录”就行。从那以后所有任务的命令都改成绝对路径工作目录要么不写继承要么写死。5.3 类似坑位清单还遇到过几个同类问题虽然不是机密但也值得拿出来当案例说时区不对服务器默认 UTCcron 表达式0 2 * * *按 UTC 触发实际比北京时间晚 8 小时。ruflo 里默认用系统时区要跑北京时间必须显式写timezone: Asia/Shanghai。我也因为这个在没配 timezone 的环境上半夜被叫起来过。YAML 数字解析陷阱在 shellrun里写--threads 3如果在某个值上不加引号YAML 解析器可能把“3”当作整数或字符串但传到 shell 命令里最终拼接结果没区分所以大部分情况没暴露。最坑的是2024-06-01这种字符串如果用了某些解析库会被当成日期对象除非显式加引号。稳妥的办法是统一用{{date}}模板变量不要手写日期。stderr 截断一开始 ruflo 对单个任务 stderr 只保留 4KB 末尾内容长日志问题根本看不到根源。后来我改成把完整 stderr 追加到日志文件同时在 UI 里做尾部截断显示。遇到类似问题的兄弟最好第一时间打开完整日志文件看别只看聚合面板。环境变量污染多个任务在同一台机器上跑如果脚本里 export 了全局变量会影响后续任务。ruflo 做子进程隔离后每个任务只继承全局env和任务自身env互相之间不会再污染。很多团队从 shell 脚本迁移过来时容易忽略这点。6. 性能、边界与后续玩法6.1 实测性能数据自己开发调度器我最担心的是空转时 CPU 占用和超大工作流图上的调度延迟。实测数据如下虚拟机 4C8G纯计算环境常驻内存约 25MB。调度 1000 个 task 的 DAG从触发到第一个 task 启动的耗时约 12ms。其中大部分时间花在读取 SQLite 状态和拓扑排序上这个量级对单机调度器来说毫无压力。单个空 shell 任务的进程冷启动开销约 2ms加上子进程 fork/exec 的整体损耗远小于任务本身的耗时基本可忽略。连续记录 10 万条执行历史后数据库文件约 80MB查询最近一周的执行记录仍然在十几毫秒内。按我的经验只要你的任务总数不超过 5000 条、单任务耗时不超过几小时、不需要跨机器调度ruflo 这种单机方案是完全够用的。6.2 什么场景别用 ruflo工具只有用在对的地方才有价值ruflo 也一样以下场景我觉得别硬套需要跨多台机器执行任务。比如你有 5 台服务器任务分散在不同机器上一台机器挂了还得能继续跑这就超出“单机调度”范围了应该上分布式任务队列或完整的 Airflow 方案。任务之间存在复杂的动态依赖。依赖关系在运行中才能确定比如“如果 A 任务输出了文件 X就跑到 B否则跑 C”这种动态分叉在 YAML 静态 DAG 模式下做不了必须用代码来实现。需要多人协作、有审批流的调度平台。ruflo 的设计哲学是“KISS”配置文件就是全部没有权限管理、没有版本回退、没有 UI 审批。真需要这些功能直接换平台不要自己二次开发。6.3 后面我想加的东西ruflo 目前的版本对我来说已经够用但我也记下了几个后续要做的方向第一个是远程执行器用 SSH 把任务下发到其他机器这能解决一部分跨机器需求但也把整体复杂度带高了优先级我排得不高。第二个是可视化面板用本地 WebSocket 推流做一个“只读”的 DAG 实时状态界面。不做编辑功能编辑还是靠 YAML因为线上环境改配置必须走 Git 审计可视化编辑反而不好管控。第三个是任务级别的资源监控收集子进程的 CPU、内存峰值这样以后排查 OOM 的时候不用再到系统层翻记录了。坦白说这几个功能我自己也知道不是所有人用的上但我做开源项目的原则一直是“先解决自己的问题再去想通用性”。如果你正卡在“cron 太裸、大框架太重”这个中间点上我建议你直接用类似思路搭一个小的依赖用 DAG、配置用 YAML、状态落 SQLite优先级就这么排至少比堆一堆 shell 脚本要强得多。ruflo 后续版本还会保持单文件分发和 YAML 优先的路线run 起来不折腾就是我对它最大的期许。