轻量级Python工作流引擎ruflo:JSON定义、状态机与审批流实践

发布时间:2026/9/9 12:09:15
轻量级Python工作流引擎ruflo:JSON定义、状态机与审批流实践 1. 轻量工作流引擎的设计初衷为什么我写了 ruflo我是在处理一个内部系统的审批流时第一次动了写一个工作流引擎的念头。当时团队的做法是每个流程用状态字段加一堆 if else 去硬写业务方提一个“加一层审批”的需求开发得改三四张表、五六个接口。改完还要担心历史数据的状态不兼容谁都怕动那一段代码。ruflo 这个名字拆开看就是 “rule” 和 “flow” 的组合。我的目标是做一个非常轻量的规则驱动型流程引擎它将业务流程定义为一份结构化的 JSON 文件运行时读取这份定义按节点推进状态把“人写的判断逻辑”变成“可配置的流程描述”。它不是我拍脑袋造的轮子而是我在对比了 Airflow、Temporal 这些重量级方案之后确认市场上缺一个“不需要部署集群、不需要写一堆 worker、能嵌入现有 Flask 项目”的这么一个小东西。1.1 核心需求解析ruflo 要解决的核心需求归纳下来有三个。第一个是流程可视化配置。产品经理和业务运营不应该每一次小改动都提工单如果流程定义可以用一份结构清晰的文件描述修改一个节点的指向只需改一处配置那迭代效率就上来了。第二个是运行状态可观测。流程跑到哪个节点了、当前在等谁审批、失败了是重试还是要人工介入这些信息必须能通过接口快速查询到不能靠翻数据库日志猜。第三个是无侵入集成。它不能约束主项目的技术栈最好是一个 Python 库pip install 就能用DB 表结构也尽量简单能跟现有数据库共存甚至支持 SQLite 起步。1.2 为什么不用现成的开源引擎评估过几个主流方案。Airflow 适合定时批量任务的调度但它要单独的 web 服务、独立的元数据库调度模型是 DAG对于“按用户操作触发的一次性审批流”来说太重了。Temporal 是微服务编排神器可靠性极高但你得额外部署集群再引入一套 SDK 的异步编程模型小团队真的负担不起。还有 Activiti 这条 Java 路线的框架BPMN 2.0 规范太庞杂配置文件要写一堆 XML和 Python 生态也不搭。ruflo 走的是一条极简路径。我只保留 Workflow、Node、Transition、WorkflowInstance 这四个核心模型状态流转矩阵固定在一个有限状态机里面节点执行逻辑用注册回调的方式暴露给开发者的业务代码。这样它既没有分布式事务也没有消息中间件依赖但你手里 90% 的会签、审批、条件分支、子流程类的需求它都能处理。2. 流程定义与节点模型一份 JSON 怎么跑起来ruflo 的流程定义文件是一个 JSON 数组最外层描述流程元信息内部是节点对象列表。每个节点有 id、type、next 三个基础字段复杂节点再附加 config。运行时引擎只认这个标准结构不关心业务字段是什么因此业务接入非常清爽。2.1 节点类型解析目前内置的节点类型有 start、task、condition、fork、join、approve、timer、end。start每个流程必须有且仅有一个起始节点作为 WorkflowInstance 创建的入口。task执行一个具体的业务动作例如调用外部 API、写数据库、发消息。它的执行逻辑要提前注册注册形式是 handler 函数。condition判断条件必须配置一个 condition 表达式或者一个返回布尔值的函数引用然后根据结果决定走哪个分支。fork把一个流程分裂成多个并行分支常用于会签、并行子任务。join聚合多个并行分支支持 wait_all、wait_any 两种聚合策略。approve这是最常用的审批节点它不自动执行任何业务逻辑而是挂起等待外部接口调用 approve 提交结果。timer定时等待节点可设置延时秒数或 cron 表达式常用于超时自动通过或超时提醒。end流程终止节点每个流程可以配置多个 end但只有第一个被触达的有效。2.2 节点连接关系定义节点之间的连接不是用通用 next 字段硬编码的我把转移关系独立成 transitions 数组。每条 transition 有 from、to、condition 三个属性condition 为空表示无条件转移有值则是一段 Python 表达式表达式运行时会自动注入全局上下文 context。例如条件节点这样配置{ id: check_amount, type: condition, config: { expression: amount 10000 } }然后 transitions 里这样写[ {from: check_amount, to: leader_approve, condition: amount 10000}, {from: check_amount, to: cto_approve, condition: amount 10000} ]这样设计的核心好处是流程流转逻辑和节点执行逻辑彻底分离。以后想改审批阈值只改 JSON 和触发接口的入参一行业务代码都不用动。我实测下来业务方对这种配置的接受度很高他们改 JSON 甚至比改代码还熟练。2.3 节点执行器的注册机制光有 JSON 定义引擎是不会干活的你得告诉它 task 类型的节点具体做什么。ruflo 提供 register_executor 接口用装饰器方式绑定。from ruflo import WorkflowEngine engine WorkflowEngine() engine.register_task(send_notify) def send_notify(context): # 这里写发邮件的逻辑 user_email context[requester_email] send_mail(user_email, 您有一条待办审批) return {status: sent}返回的字典会自动合并回 context后续节点都能读取。注册机制让业务方可以自由控制幂等性和重试逻辑这是 ruflo 比纯代码流程方案更近一步的地方。3. 安装与快速初始化从零跑通第一个审批流ruflo 目前发布在内部 PyPI 源安装很简单pip install ruflo装完之后先初始化一个工作目录作为流程定义和运行数据的存放位置ruflo init --home /data/ruflo执行完之后会在 /data/ruflo 下生成三个子目录definitions、logs、state。definitions 放流程定义 JSONlogs 放运行日志state 放当前运行中实例的状态快照。默认使用 SQLite 作为状态库连接串写在 ruflo.yaml 里。3.1 数据库表结构与状态机设计ruflo 的表很少核心就是 workflow_definition 和 workflow_instance 两张表。前者存流程定义的元信息和版本号后者存每个实例的当前节点、状态、上下文数据、创建时间和更新时间。状态机的状态流转是这样设计的当前状态可流转状态触发动作READYRUNNING节点开始执行RUNNINGSUCCESS / FAILED / WAITING节点执行完成 / 异常 / 进入审批WAITINGRUNNING / TERMINATED审批接口提交结果 / 超时或取消SUCCESSTERMINATED后续主动终止FAILEDRUNNING手动重试我会额外记录 transition_log 表记录每一次跳转的 from_node、to_node、触发时间、操作人作为审计链路。小系统也需要审计能力很多纠纷排查时这一张表能救你。3.2 定义一个请假审批流程跑一个最简单的流程场景是员工请假三天以内主管审批即可超过三天要总监加签。先写请假流程定义 leave_flow.json{ flow_id: leave_approval, name: 请假审批流程, version: 1, nodes: [ {id: start, type: start, name: 开始}, {id: fill, type: task, name: 填写申请, config: {handler: fill_leave_form}}, {id: check_days, type: condition, name: 判断天数}, {id: leader, type: approve, name: 主管审批, config: {assignee_expression: requester.leader}}, {id: director, type: approve, name: 总监审批}, {id: done, type: end, name: 结束} ], transitions: [ {from: start, to: fill}, {from: fill, to: check_days}, {from: check_days, to: leader, condition: leave_days 3}, {from: check_days, to: director, condition: leave_days 3}, {from: leader, to: done, condition: approve_result agree}, {from: leader, to: done, condition: approve_result reject}, {from: director, to: done} ] }注意 approve 类型的节点不会自己往下跳必须等外部调用接口提交审批结果。这一点和 task 节点完全不同我在设计时特意区分开避免业务逻辑把“发起请求”和“等待审批”混在一起。3.3 启动引擎与触发流程流程定义写好之后加载定义并创建引擎实例from ruflo import WorkflowEngine from ruflo.loader import load_definition engine WorkflowEngine() definition load_definition(definitions/leave_flow.json) engine.load_definition(definition)触发流程时传入初始上下文ctx { requester: {id: 1001, name: 张三, leader: 李经理}, leave_days: 2, reason: 家里有事 } instance engine.start_flow(leave_approval, ctx) print(instance.instance_id) # 返回唯一实例 ID引擎会从 start 节点开始自动流转到 fill 节点执行填单逻辑然后进入 check_days 判断最后走到 leader 审批节点挂起。审批接口的调用方式engine.approve(instance_idwf_xxxx, node_idleader, operator李经理, resultagree)如果上下文里缺了 operator 相关的字段引擎会在进入审批节点时报参数错误这点在对接外部系统时经常遇到后面我会单独展开。4. API 设计与外部系统集成把流程能力开放给业务方ruflo 的引擎本身是嵌入式的但如果你有多个服务需要共享流程能力我会推荐再加一层轻量 HTTP API。ruflo 自带了一个 FastAPI 的示例封装启动之后直接暴露几个核心接口。4.1 核心 REST 接口列表接口路径方法功能说明/api/flowsPOST创建流程实例/api/flows/{instance_id}GET查询实例状态/api/flows/{instance_id}/approvePOST审批提交/api/flows/{instance_id}/cancelPOST取消实例/api/flows/{instance_id}/retryPOST失败重试/api/definitionsGET列出已加载的流程定义/api/definitions/{flow_id}POST热加载新的流程版本接口的鉴权我建议用简单的 API Key 放到 header 里内部服务之间调用足够不要过度设计。要注意的是查询实例详情接口一定要返回当前节点名称、状态、上下文、流转日志这四个字段业务方前端画审批进度条全靠它们。4.2 回调任务与消息通知task 节点执行完业务逻辑以后如果希望异步通知其他系统可以在 handler 里把消息发到收发链路。我自己的做法是写一个 notify 装饰器在 handler 成功后自动打一条 webhookengine.register_task(fill_leave_form) notify(http://hr-service/api/leave/notify) def fill_leave_form(context): insert_leave_record(context) return {record_id: 233}这里有个小坑notify 的 webhook 调用一定不能阻塞流程主线程。ruflo 的设计是 handler 返回后立即提交状态变更webhook 进入待发送队列由后台 worker 消费。如果不这样做一次接口超时会卡死整个流程引擎。4.3 流程版本管理与热更新线上引擎最怕的是流程定义改坏了没法回滚。ruflo 的方案是版本号控制每次加载同 id 的新定义自动增加版本号运行中的实例继续执行旧版本。只有新触发的实例才用新版本。这个策略一开始就有意设计的否则审批一半的员工突然发现流程配置变了那直接乱套。热更新接口engine.load_definition(definition, active_versionFalse) engine.activate_version(leave_approval, version3)先加载但不激活验证没问题后手动激活或者做一个自动 diff检测节点 id 变更列表。我的习惯是每次上线前写个 pytest跑一遍全流程正向和反向用例确保新版本至少能走通正常路径。5. 状态机进阶会签、超时、驳回与动态节点基础流程跑通以后你会发现生产环境的需求像无底洞。第一个进阶需求大概率是多人会签第二个是超时处理。ruflo 的 fork/join 就是为会签设计的。5.1 并行分支与会签实现会签的意思是“多个人都要审批全部通过才算通过”。流程定义里用 fork 同时拉出多个 approve 节点再由 join 聚合。{ id: fork_start, type: fork, config: { branches: [approve_by_finance, approve_by_hr] } }, { id: join_all, type: join, config: {strategy: wait_all} }wait_all 策略下引擎会为每一条分支创建子实例父实例处于 WAITING 状态。每个审批节点被 approve 后子实例结束并回传结果。所有子实例 SUCCESS 以后父实例自动唤醒进入 join_all 的后续节点。wait_any 策略则恰好相反任意一个子实例成功后父实例立即唤醒其余分支会被标记为 SKIPPED。这里容易踩的坑是分支节点数超过二三十个比如公司全员投票场景SQLite 状态库写入会明显变慢。我的建议是单流程分支数控制在十个以内再大就拆子流程由父流程分阶段触发。5.2 审批超时自动提醒与自动通过审批节点挂起时timer 节点可以在旁边做配套。我的方案是审批节点配置 timeout_hours 字段引擎扫描到超时会触发一个 timeout_callback。比如主管审批超过三天没处理第四天自动催促第七天自动通过转交engine.on_timeout(leader, leave_approval) def handle_leader_timeout(instance_id): # 发提醒消息 send_reminder(instance_id) if aging_days(instance_id) 7: engine.approve(instance_id, node_idleader, operatorsystem, resultagree, autoTrue)自动通过是个危险操作一定要在审计日志里带上 operatorsystem 且 autoTrue 的字样。我之前遇到一个情况审批节点配了自动通过结果某个流程卡了两周没被发现系统自动批了一笔不该批的款项。后来加了两层保障第一自动通过只允许指定的节点类型和特定条件第二自动通过必须记录理由。5.3 驳回与重新提交驳回是审批流的家常便饭。我的实现是允许审批节点配置 reject_to 字段指定驳回后跳转到哪个节点。常见做法是驳回到发起人节点发起人修改后重新提交流程重新进入审批链。{ id: leader, type: approve, config: { assignee_expression: requester.leader, reject_to: fill } }重新提交的时候要注意 context 数据的处理。原上下文里的填单数据必须保留但审批意见要追加到记录里不能覆盖。ruflo 维护一个独立字段 approval_history 列表每次 approve 都追加一条。这样不但可以画完整的审批轨迹还能统计每位审批人的平均处理时长做绩效分析也方便。6. 实操中遇到的常见问题与排查经验ruflo 写出来之后团队内部用了快五个月GitHub issues 里收集了三十多个反馈。这里挑几个高频问题聊聊解法。6.1 节点执行成功但实例状态没有推进这是最典型的装配错误。task 执行器正常返回了但流程实例一直停在 RUNNING。排查时先看日志确认 handler 有没有抛出异常。如果 handler 正常结束那问题多半出在 transitions 配置上。下一步查 condition 表达式是否被正确求值常见的是浮点数比较没有做类型转换比如流程里 1 是字符串表达式里 1 是数字两个永远不相等。{condition: amount 10000}如果上下文里的 amount 是从数据库读出来的 Decimal 类型那比较没问题。但如果前端提交来的 amount 是字符串 8888这个条件永远 False。解决方式是统一在 handler 里做类型清洗或者 condition 表达式里显式转换float(amount) 10000。还有一个隐蔽的情况就是节点 id 和 transition 里 from/to 的 id 不匹配少写一个字母引擎静默跳过。我在新版本里加了启动时的节点引用完整性校验加载定义时如果发现挂空引用直接抛异常宁可启动失败也别跑去线上发现问题。6.2 审批节点一直 WAITING外部系统不回调审批节点的 WAITING 状态完全依赖外部接口回调。实际生产中调用方可能因为网络抖动、服务重启导致提交结果丢失。我的建议是设计一个补偿查询接口让调用方定时批量询问“有哪些实例在等我审批”。ruflo 提供了一个查询待办接口engine.get_pending_approvals(assignee李经理)返回当前所有 assignee 字段匹配、且状态为 WAITING 的实例列表。调用方只需要每个小时轮询一次这个接口把待办列表和本地数据库对比缺少记录就补发回调。这个方案简单粗暴但是非常有效我跑了快半年没有再因为网络丢包丢过审批记录。6.3 fork 分支过多导致的状态数据膨胀前面提过分支二十条以上会导致 SQLite 写入变慢我给的解决方向是拆子流程。ruflo 的 task handler 里可以再调用 start_flow 去创建子流程实例父流程在关键节点等着即可。比如一个采购会签流程预算超过五十万就不走普通分支了而是在 handler 里动态创建三个独立采购评审子实例子实例全部完成后父流程通过 join 节点聚合。这样既保证了扩展性也不会让单个流程定义膨胀到无法维护。6.4 流程定义改动了怎么保证存量数据安全这是非常关键的规则已经 RUNNING 的实例永远不要尝试迁移到新定义。每个 workflow_instance 表里都存 definition_version引擎运行时读取的是该实例当时的版本快照不是全局最新版本。如果你确实需要修改一个正在运行流程的未执行节点行为正确的做法是利用动态配置覆盖层的逻辑。在节点配置里加一个 condition 依赖某个开关变量的写法可以让判定结果实时变化。最典型的例子是审批金额阈值调整。定义里这样写{ type: condition, config: { expression: amount threshold(flow_idleave_approval, optionmax_days) } }threshold 函数从配置中心读取阈值这样只改配置中心的值不碰流程定义也能影响运行中实例的走向。把这类易变参数从流程定义中抽出来是我推荐的一个设计习惯。7. 生产部署与稳定性设计ruflo 的定位是嵌入式引擎你可以作为进程内库直接调用也可以独立启动 HTTP 服务作为流程中台。两种模式我都试过说出各自的场景适配。7.1 嵌入式模式适配中小项目如果你的项目是一个 Django/Flask 单体应用流程引擎跟主服务部署在同一个进程里。好处是调用无网络开销事务天然一致数据源共用一个库。坏处是流程执行逻辑不稳定的话会拖垮主服务。在这个模式下要把 ruflo 的 handler 函数统一包一层 try/except异常吞掉后标记 FAILED 并写日志绝不能往上抛。流程引擎问题不能影响业务主链路这条原则我写在团队开发规范里。7.2 独立服务模式适配多系统共享当你有多个系统都涉及审批、比如仓库管理系统和采购系统都要走同一个审批流那必须独立部署流程服务。数据库单独建库对外只暴露 REST API 和 Webhook 回调。同一个实例的创建、修改、审批全部走接口不要让人直接连数据库操作容易破坏状态一致性。部署时进程管理用 systemd 即可不需要 K8s。并发量级到不了那个程度一个 gunicorn 主进程加四个 worker 已经能扛住日均万级实例。SQLite 换到 PostgreSQL 连接串就能平滑迁移我建议线上直接用 PostgreSQL多个 worker 并发写 SQLite 容易产生锁等待。7.3 日志、监控与告警我做的日志分为三层。第一层是系统日志记录引擎本身的错误和警告。第二层是实例日志记录每个实例在每一步的状态变化、花费时长。第三层是审计日志记录所有审批操作人、动作、时间、结果这层数据保留至少一年应对内外审。告警方面写了两个指标第一个是 WAITING 超过 24 小时的实例数第二个是失败重试次数超过三次的实例数。只要这两个指标任一超过阈值企业微信机器人就会推送消息给值班负责人。实测下来超时告警能提前暴露没人处理的僵尸流程重试告警能及时发现外部服务不稳定。8. 扩展性设计如何接入你自己的业务动作ruflo 官方内置的 task handler 非常少因为每个业务的 handler 不一样。所以把这套东西接进业务系统的重点就是注册你自己的执行器和扩展节点类型。8.1 注册业务执行器前面演示过装饰器注册。这里要强调注册时机的问题执行器的注册一定要在引擎加载定义之前完成否则定义里配置的 handler 找不到对应实现引擎会在加载时报错。engine WorkflowEngine() engine.register_task(fill_leave_form, fill_leave_form) engine.register_task(send_notify, send_notify) definition load_definition(leave_flow.json) engine.load_definition(definition)如果你的执行器有依赖比如它要操作数据库、调用 Redis、访问外部 SDK建议用高阶函数的写法def build_fill_handler(db_session): def fill_leave_form(context): db_session.add(...) return {record_id: 233} return fill_leave_form然后在初始化引擎时把依赖注入进去这样测试时也能很方便地替换 mock 实现。8.2 自定义节点类型有些流程动作比如“从 Excel 导入数据后继续流程”“发送企业微信消息卡片”在任务里面做太笨重适合扩展成新的节点类型。ruflo 支持 register_node_type 接口自定义节点类型只要实现 execute 方法。engine.register_node_type(wecom_notify) class WecomNotifyNode: def __init__(self, config): self.webhook_url config.get(webhook_url) def execute(self, context): send_wecom_message(self.webhook_url, context) return {notify_status: ok}这样流程定义里就能直接用 wecom_notify 类型外面看起来跟内置的 task 节点无异。团队里的人后来自己还扩展了钉钉通知节点、短信验证码节点完全不需要改 rufo 核心代码。8.3 与消息队列rocket配合做异步处理当流程吞吐量变大以后task 节点的同步执行会成为瓶颈。我的建议是将重的任务节点改成“发出任务消息就返回后台 worker 消费完再回调引擎推进”的模式。相当于把 task 节点模拟成 approve 节点但发布的消息是给 worker 干活的。一个简化的写法是engine.register_task(async_job) def async_job(context): topic context[job_topic] message {instance_id: context[__instance_id], node_id: context[__node_id]} producer.send(topic, message) return {async: True, job_id: job_id}然后在 worker 消费完成后调用引擎接口把当前节点标记为执行完成引擎再继续往后跳。这套异步化改造做下来整个引擎的吞吐能力从每秒几十个实例提升到几百个完全够用到上万用户的规模。9. 我对 ruflo 后续迭代的三个设想第一增加可视化流程编辑器。现在 JSON 定义还是有一定门槛业务方直接改配置确实会用但不敢大调整。我想做一个拖拽式的 Web 画布节点连线、条件配置、审批人设置全在界面上完成最后导出标准 JSON。这个其实不难前端用现有开源组件改改就行。第二增加流程性能分析。基于 instance 表和 transition_log统计每个节点平均耗时、驳回率、超时率、审批人效率排名。这些数据对梳理组织权限、优化审批结构特别有价值。我见过一个客户他们发现某个节点驳回率高达 70%一查是审批人根本不清楚业务规则纯靠人工判断后来加了一个业务说明字段驳回率直接降了一半。第三支持多级子流程的动态编排。当前版本子流程创建以后父流程只能干等但我希望子流程结束以后能把结果回写父流程然后父流程根据结果决定下一步走哪个节点形成一个可嵌套、可动态扩展的流程树。这块实现上难度主要在于状态同步和并发控制但方向是对的。ruflo 这个项目的核心目标始终没有变就是让流程编排变得足够轻、足够透明、足够可控。轻是指架构上不引入无谓依赖透明是指每个实例每一步都有记录可查可控是指任何异常状况都有接口可以人工介入。它不追求大而全只希望你在需要流程能力的时候能花最少的时间把它揉进自己的系统里。我个人在实际使用中的体会是工作流引擎这类中间件稳定和可排查永远排在功能丰富前面。一个跑不崩、出了问题能快速定位的小引擎比一个功能全面但你根本不敢升级的大型框架有用得多。

关于本文作者

来自尧图内容编辑团队

尧图内容编辑团队 内容团队

尧图内容编辑团队

本文由尧图网络内容编辑团队执笔。团队由资深项目经理、前端工程师与设计师组成,所有内容均来自亲手交付的真实项目,先讲清问题、再给出可落地的解法。尧图深耕北京网站建设十年,服务过京华建材集团、智造科技等各行业客户,把一线经验沉淀为可复用的行业观察。

  • 十年建站经验,覆盖建材、制造、服务、文创等
  • 项目经理把关选题与事实准确性
  • 工程师与设计师联合撰写专业细节
  • 统一编辑规范,保证文风与排版一致
  • 每月复盘转化数据,迭代选题方向

延伸阅读

相关资讯与近期热门内容

深度阅读推荐

建站决策前值得细读的三篇

网站改版的5个关键决策
2024-08-12

网站改版的5个关键决策

什么时候该改版、改到什么程度、如何避免流量掉光,京华建材集团改版复盘给出答案。

获取专属建站方案

看完文章,把您的行业与预算告诉我们,免费获取一份量身定制的官网建设方案与报价。

立即免费咨询