轻量级Python工作流引擎ruflo:JSON定义、状态机与审批流实践
2026/9/9 12:09:12 网站建设 项目流程

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 放流程定义 JSON,logs 放运行日志,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_id="wf_xxxx", node_id="leader", operator="李经理", result="agree")

如果上下文里缺了 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 成功后自动打一条 webhook:

@engine.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_version=False) engine.activate_version("leave_approval", version=3)

先加载但不激活,验证没问题后手动激活,或者做一个自动 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_id="leader", operator="system", result="agree", auto=True)

自动通过是个危险操作,一定要在审计日志里带上 operator=system 且 auto=True 的字样。我之前遇到一个情况,审批节点配了自动通过,结果某个流程卡了两周没被发现,系统自动批了一笔不该批的款项。后来加了两层保障:第一,自动通过只允许指定的节点类型和特定条件;第二,自动通过必须记录理由。

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_id='leave_approval', option='max_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 这个项目的核心目标始终没有变,就是让流程编排变得足够轻、足够透明、足够可控。轻是指架构上不引入无谓依赖,透明是指每个实例每一步都有记录可查,可控是指任何异常状况都有接口可以人工介入。它不追求大而全,只希望你在需要流程能力的时候,能花最少的时间把它揉进自己的系统里。

我个人在实际使用中的体会是,工作流引擎这类中间件,稳定和可排查永远排在功能丰富前面。一个跑不崩、出了问题能快速定位的小引擎,比一个功能全面但你根本不敢升级的大型框架有用得多。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询