资讯动态

自研轻量级规则流引擎ruflo:告别if-else,用图编排业务流程

发布时间:2026/9/9 11:00:04 来源:尧图企业网站定制
如果你也在维护一个写满了 if-else 的业务模块大概率能理解我下面要聊的这件事。去年我在订单系统里被一段八百行的判断逻辑折磨到不行改一个分支要顺着调用链翻半天还总担心改完会影响哪个下游。后来我干脆自己写了一个非常小的规则流编排引擎取名叫ruflo把流程和规则从业务代码里拆出来整个维护体感完全变了。ruflo 的定位不是 Drools 那样的大而全规则引擎也不是 Camunda 那种重量级流程平台它就解决一件事用图来表达业务流转用规则来决定走向。简单说流程里的每个步骤是一个节点节点之间靠规则连接规则返回布尔值引擎根据规则命中情况决定走哪条路。如果你也在维护一堆嵌套 if-else或者正打算自己封装一个轻量编排工具这篇应该能帮你省不少弯路。下面我从设计思路、核心实现、业务接入到踩坑经验完整拆一遍。1. 为什么我会自己写一个叫 ruflo 的规则流编排引擎1.1 规则引擎和流程引擎之间存在一个没人管的空档先说说我当时的处境。订单系统里有大量业务判断黑名单、会员等级、优惠券门槛、地区限制、库存状态、风控标记全都写在 service 层的一个大方法里。每来一个需求就往里面加一个 if每个 if 后面还挂着几个 return。代码能跑但没有一个人敢改。我评估过市面上现成的方案Drools这类规则引擎很强适合大规模规则推理但对一个普通业务项目来说学习成本不低规则语法也有自己的心智负担有点杀鸡用牛刀。Airflow / Temporal / Camunda这类流程引擎适合重流程编排、长任务调度但要么要部署平台要么要引一整套 SDK对一个单机应用来说太重了。市面上的规则库、决策表工具能解决单个条件命中的问题但表达不了先做 A再根据结果决定 B 还是 C最后汇合做 D这种带顺序和分叉的流程。我需要的其实很小用几行配置或代码描述一个流程让规则决定走哪个分支每一步执行都有日志出问题能定位到具体节点。这个需求正好卡在规则引擎和流程引擎之间的空档。既然没有合适的小工具那就自己写一个。1.2 ruflo 的设计原则流程靠图规则即数据ruflo 这个名字很简单ru 取 rule 的前两个字母flo 取 flow 的前三个字母合起来就是 rule flow规则流。名字定了之后我把这个工具的设计原则也顺带压缩成两句话流程靠图。所有业务步骤抽象成节点节点之间的连接关系构成有向图。有了图流程就变成可分析、可渲染、可拖拽的数据而不是硬编码的函数调用栈。规则即数据。一条规则就是一条记录名称、条件表达式、目标节点、优先级。规则不再散落在 if-else 里而是作为独立的实体存在。改规则不需要动业务代码只是改一条配置或者一行数据。这两句话听起来简单但落地的时候会牵扯出不少细节节点怎么定义、规则怎么挂在边上、上下文怎么传递、引擎怎么调度。下面逐一说清楚。2. ruflo 的五个核心设计决策2.1 节点是最小执行单元在 ruflo 里所有业务动作都是一个节点。判断用户是否合法、计算订单折扣、调风控接口、写审计日志全部可以拆成节点。节点的输入输出都走同一个上下文对象节点之间不直接互相调用。节点的抽象我用 Python dataclass 表示核心字段就四个from dataclasses import dataclass from typing import Optional, Callable, Any dataclass class Node: id: str # 节点唯一标识 name: str # 节点展示名日志里能看懂 handler: Optional[Callable[[Context], None]] None # 真正执行的逻辑 extra: dict None # 扩展字段比如超时、重试次数每个节点尽量只做一件小事。比如检查用户是否在黑名单是一个节点调用频次限制服务是一个节点计算最终价格是另一个节点。节点拆得足够细后面才能自由排列组合复用。如果节点粒度太粗把三个动作塞进一个节点里那不叫编排换汤不换药。2.2 规则不是 if-else而是一等公民规则在 ruflo 中是挂在两个节点之间的边上的。引擎执行完当前节点后会遍历这个节点的所有出边按优先级排序取第一个规则结果为真的边跳到目标节点。规则就是一个接收上下文、返回布尔值的函数。但和直接写在业务代码里的 if 不同规则会被注册到规则表里拥有自己的名字。这样设计带来的最大好处是流程路由和规则实现解耦了。from typing import Callable def is_vip_user(ctx) - bool: return ctx.get(user_level) 1 def is_new_user(ctx) - bool: return ctx.get(register_days, 0) 7 def is_amount_over_threshold(ctx) - bool: return ctx.get(order_amount, 0) 5000这几个函数本身什么都不做只负责回答当前上下文满不满足某个条件。真正决定这些规则怎么组合成流程是引擎里的连接关系。规则变成一种可以被注册、被引用、被改写的资源而不是一段拍在代码里的分支语句。2.3 上下文对象贯穿全程节点执行需要数据规则判断也需要数据这些数据都放在 Context 里。Context 本质是一个带日志功能的字典但我在设计它时多做了三件小事取代全局变量。整个流程执行期间节点之间传数据都通过 Context不搞隐式全局状态。谁都能往 Context 里写数据规则只负责读。记录执行轨迹。每个节点开始、结束、命中规则、结果数据全部追加到日志列表。上线后排查问题直接看一条流程的日志就能还原当时发生了什么。隔离每次请求。每个请求必须 new 一个上下文对象绝不能复用全局实例。这一点在后来的并发故障里救了我后面详细说。from dataclasses import dataclass, field from typing import Any, List dataclass class Context: data: dict field(default_factorydict) logs: List[str] field(default_factorylist) def get(self, key: str, default: Any None): return self.data.get(key, default) def set(self, key: str, value: Any): self.data[key] value def log(self, message: str): self.logs.append(message)2.4 调度器决定执行顺序用图描述而不是用嵌套很多人在自研编排引擎时容易陷入误区用函数递归表示流程A 方法里调用 B 方法B 方法里调用 C 方法。这种编码式编排虽然直观但流程一旦加了分支、循环、汇合调用关系立刻乱掉改起来还是要动代码。ruflo 从第一版开始就坚持用有向图来描述流程。节点是图中的点带规则的边是图中的弧。执行器从起点出发在当前节点完成后根据规则命中情况选择下一个节点。这个模型表达力足够强链路、分支、收敛都能表达而且图结构天然适合做可视化、静态分析和循环检测。2.5 节点和规则都要可注册、可扩展ruflo 没有做成一个大而全的框架而是定义了两套注册接口让业务方自己往里塞东西节点注册把一段处理逻辑包装成 Node 放进引擎。规则注册把条件函数放进规则表用字符串名字引用。业务方不需要继承任何基类只要写普通函数然后声明我是一个节点或我是一个规则即可。我不会引入复杂的 SPI 或插件机制因为对这种体量的工具来说普通函数 注册表已经够用了。3. 最小可运行版 ruflo从零实现一个规则流引擎代码永远是表达设计最好的方式。这一节我给出一版浓缩但能直接跑的 ruflo 核心实现约 150 行左右包含节点注册、规则边、循环检测和主执行循环。3.1 核心数据结构Node 和 Context 的定义上面已经写过这里补上 Edge 和 Enginefrom dataclasses import dataclass from typing import Any, Callable, Dict, List, Optional, Set # 边从一个节点到另一个节点带一个可选规则 dataclass class Edge: source: str target: str rule: Optional[Callable[[Context], bool]] None priority: int 0 class Engine: def __init__(self): self.nodes: Dict[str, Node] {} self.edges: Dict[str, List[Edge]] {} self.rule_table: Dict[str, Callable[[Context], bool]] {}3.2 引擎的注册与连接Engine 提供三个对外方法add_node注册节点connect连接两个节点并挂上规则register_rule把规则函数登记到规则表。def add_node(self, node: Node) - None: self.nodes[node.id] node self.edges.setdefault(node.id, []) def register_rule(self, name: str, fn: Callable[[Context], bool]) - None: self.rule_table[name] fn def connect( self, source: str, target: str, rule: Optional[Callable[[Context], bool]] None, priority: int 0, ) - None: self.edges[source].append(Edge(source, target, rule, priority))connect的rule参数可以直接传函数也可以传规则名字符串。直接传函数在写代码时比较顺手传字符串则方便以后改成配置化加载两种方式在底层都转成可调用的规则对象。3.3 主执行循环带分支的选择逻辑引擎的run(start, ctx)从指定起始节点开始执行节点然后根据规则优先级选择下一跳直到没有下一个节点可走。def run(self, start: str, ctx: Context) - None: if self._has_cycle(): raise RuntimeError(ruflo: flow contains a cycle) done: Set[str] set() pending: List[str] [start] while pending: node_id pending.pop(0) if node_id in done: continue node self.nodes.get(node_id) if node is None: raise KeyError(fruflo: unknown node: {node_id}) ctx.log(f[{node.id}] {node.name} 开始) if node.handler is not None: node.handler(ctx) ctx.log(f[{node.id}] {node.name} 完成) done.add(node_id) next_node self._choose_next(node_id, ctx) if next_node is not None and next_node not in done: pending.append(next_node) def _choose_next(self, source: str, ctx: Context) - Optional[str]: edges sorted(self.edges.get(source, []), keylambda e: e.priority, reverseTrue) for edge in edges: if edge.rule is None or edge.rule(ctx): ctx.log(f[{source}] 命中规则前往 {edge.target}) return edge.target return None这段逻辑值得重点解释。_choose_next把当前节点的所有出边按priority降序排列然后逐条评估规则只要遇到第一条规则结果为真的边就立即选定目标节点并返回。这其实就是业务里最常用的条件路由表模型适合绝大多数分支场景。done集合用来防止节点被重复执行。如果一个节点已经被执行过后续即使又从另一条路径指向它也会被跳过。这是简化处理适用于大多数汇聚但不需重复执行的流程。如果业务流程需要严格等待多条入边全部到达后再执行join 语义就不能用这个简单方案我在第 5 章会展开讲。3.4 循环检测不能在运行时才爆栈有向图里如果存在环执行会陷入死循环。所以我在run一进来就调用_has_cycle做一次全量检测。这里用拓扑排序边数少时间复杂度 O(NE)启动时跑一次完全没压力def _has_cycle(self) - bool: indegree: Dict[str, int] {nid: 0 for nid in self.nodes} for edges in self.edges.values(): for e in edges: if e.target in indegree: indegree[e.target] 1 queue [nid for nid, d in indegree.items() if d 0] visited 0 while queue: node_id queue.pop(0) visited 1 for e in self.edges.get(node_id, []): indegree[e.target] - 1 if indegree[e.target] 0: queue.append(e.target) return visited ! len(self.nodes)当图中存在环时环上的节点永远无法把入度降到 0最后visited会小于节点总数。这个检查必须在流程启动时做绝不能等执行到一半才发现死循环。3.5 一个可跑通的 Demo会员订单折扣流程空谈设计太虚上代码。下面这个例子模拟一个简单的订单折扣流开始节点检查用户信息然后根据用户等级分流VIP 走八折分支普通用户走九五折分支最后汇总计算最终价格。# 1. 初始化引擎和规则表 engine Engine() engine.register_rule(is_vip, lambda ctx: ctx.get(user_level) 1) engine.register_rule(not_vip, lambda ctx: ctx.get(user_level) ! 1) # 2. 注册节点 engine.add_node(Node(start, 开始校验, handlerlambda ctx: ctx.log(检查订单参数))) engine.add_node(Node(vip_discount, VIP 八折, handlerlambda ctx: ctx.set(final_price, ctx.get(price) * 0.8))) engine.add_node(Node(normal_discount, 普通九五折, handlerlambda ctx: ctx.set(final_price, ctx.get(price) * 0.95))) engine.add_node(Node(summary, 汇总输出, handlerlambda ctx: ctx.log(f最终价格: {ctx.get(final_price)}))) # 3. 连接节点带上规则和优先级 engine.connect(start, vip_discount, ruleengine.rule_table[is_vip], priority10) engine.connect(start, normal_discount, ruleengine.rule_table[not_vip], priority5) engine.connect(vip_discount, summary) engine.connect(normal_discount, summary) # 4. 执行 ctx Context({user_id: u_1001, user_level: 1, price: 100}) engine.run(start, ctx) for row in ctx.logs: print(row)执行结果[start] 开始校验 开始 [start] 开始校验 完成 [start] 命中规则前往 vip_discount [vip_discount] VIP 八折 开始 [vip_discount] VIP 八折 完成 [vip_discount] 命中规则前往 summary [summary] 汇总输出 开始 [summary] 汇总输出 完成如果我把user_level改成 2那么走的就是normal_discount分支最终价格是 95。整个流程的控制流完全由图的边和规则决定业务代码里看不到一个分支语句这就是规则流的核心价值。4. 真实业务接入把 ruflo 用在订单风控场景光会写 Hello World 没用我把 ruflo 接入订单风控链路后才真正体会到这套设计的价值。这一节用一个浓缩的订单风控场景讲清楚怎么把真实业务拆成节点和规则以及接入时要注意的工程化细节。4.1 风控规则流怎么拆订单风控不是单点判断而是一条多级链路。我当时的做法是把它拆成下面这个流程基础参数校验订单号、用户 ID、金额是否为空空则直接拒绝。黑名单检查查用户是否命中黑名单命中则直接拒绝。频次限制单个用户短时间内下单是否超频超频则拒绝。金额阈值判断大额订单再走人工审核。放行或转人工所有规则通过则自动放行金额超阈值转人工。对应到 ruflo 里每个动作是一个节点每个决策是一个规则。节点清单和规则清单可以完全对应到图中流程图长这样文字描述base_check基础校验 ├─ 规则: base_failed → reject └─ 规则: base_passed → blacklist_check blacklist_check黑名单检查 ├─ 规则: blacklist_hit → reject └─ 规则: blacklist_pass → frequency_check frequency_check频次检查 ├─ 规则: frequency_over → reject └─ 规则: frequency_pass → amount_check amount_check金额检查 ├─ 规则: amount_over → manual_review └─ 规则: amount_pass → pass实际落地时我不会在代码里一个个手工connect而是把流程定义放到配置里让引擎启动时加载。配置文件用 JSON大致长这样{ start: base_check, nodes: [ {id: base_check, handler: BaseCheckHandler}, {id: blacklist_check, handler: BlacklistCheckHandler}, {id: frequency_check, handler: FrequencyCheckHandler}, {id: amount_check, handler: AmountCheckHandler}, {id: manual_review, handler: ManualReviewHandler}, {id: reject, handler: RejectHandler}, {id: pass, handler: PassHandler} ], edges: [ {from: base_check, to: reject, rule: base_failed, priority: 20}, {from: base_check, to: blacklist_check, rule: base_passed, priority: 10}, {from: blacklist_check, to: reject, rule: blacklist_hit, priority: 20}, {from: blacklist_check, to: frequency_check, rule: blacklist_pass, priority: 10}, {from: frequency_check, to: reject, rule: frequency_over, priority: 20}, {from: frequency_check, to: amount_check, rule: frequency_pass, priority: 10}, {from: amount_check, to: manual_review, rule: amount_over, priority: 20}, {from: amount_check, to: pass, rule: amount_pass, priority: 10} ] }这里有一个非常关键的约定JSON 里只能写规则名不能写函数实现。所以引擎加载配置时需要维护一个规则名到规则函数的映射表。第 3 章的register_rule就是干这个的。有了这张表配置就纯粹是一份路由数据规则逻辑仍然留在代码里可测试、可单测。4.2 接入时的工程化处理把 ruflo 接入真实系统的过程中我踩了不少坑也沉淀下几个工程化要点。第一节点和规则的代码分层。节点 handler 负责执行动作比如查黑名单表、调频次服务规则函数只负责判断不做任何 IO。比如黑名单检查节点执行完后把结果写到 Context 里def check_blacklist(ctx: Context): user_id ctx.get(user_id) is_hit blacklist_service.is_in_blacklist(user_id) # 这是 IO ctx.set(blacklist_hit, is_hit)然后规则函数只读 Contextdef blacklist_hit(ctx: Context) - bool: return ctx.get(blacklist_hit, False)这个分层的好处是规则函数变成纯函数。我可以不跑完整流程只对规则做单元测试而且节点可以任意增删规则函数几乎不用改。第二执行轨迹要完整。Context 里的 logs 不能只记节点开始结束最好把关键结果也记进去。我在风控场景里会在日志里记录命中黑名单是/否频次3次/分钟最终动作转人工审核等关键信息。这样用户投诉时我把一份执行日志导出就能完整还原每一步为什么这样走。第三超时和熔断要放在节点层。规则流本身不做超时控制因为不是每个节点都需要超时。但调用第三方服务的节点必须自己设置超时和失败策略。失败后的处理有两种选择抛异常终止流程或者设置一个降级结果继续流程。我在 ruflo 里约定节点级异常默认终止流程但允许节点内部捕获业务异常后写入降级结果。比如频次服务超时就按未超频继续走保证订单流程不被第三方拖垮。4.3 集成时不能忽略的接口约定接入 ruflo 还有一个容易忽略的地方流程入口和出口的数据约定。上游调用方可能传的是订单对象、用户对象而不是原始字段。我在接入时做了一层薄薄的适配层把上游对象转换成 Context 的初始数据def execute_risk_control(order: Order) - str: ctx Context({ order_id: order.id, user_id: order.user_id, order_amount: order.amount, payment_method: order.payment_method, }) engine.run(base_check, ctx) final_action ctx.get(final_action, unknown) return final_actionfinal_action是流程终止节点写入的结果比如pass、reject、manual_review。这样对上游业务方来说调用方只看到一个返回字符串内部是规则流还是十几个 if-else完全不暴露。这也是规则流方案相比大引擎的优势对业务方的接口保持简单内部怎么做是维护者自己的事。5. 我在 ruflo 上踩过的坑和现在仍在改的部分任何工具都要经过真实流量的毒打才能成熟。ruflo 上线到现在我遇到了几个印象很深的坑每个都花了不少时间排查写出来供大家参考。5.1 循环依赖第一次跑就栈溢出第一版执行器我用的是递归写法执行完一个节点找到下个节点后直接递归调用。结果配置流程时不小心写了一个 A→B→A 的环运行后 RecursionError 直接爆栈。更麻烦的是流程是执行到一半才爆的栈前几个节点已经跑完了数据处于半更新状态。从那以后我意识到循环检测必须在引擎启动时做而不是等运行时才暴露。现在 ruflo 在run方法进入时会先对整个图做一次拓扑排序检测有环就直接抛异常。宁可启动时报错也不能让线上跑到一半爆栈。5.2 上下文线程安全的教训这个坑最隐蔽。早期版本的 Context 我图省事把data定义成模块级变量心想反正一个流程跑完就清零。结果上线后出现了一个特别诡异的现象A 订单的用户等级莫名其妙出现在了 B 订单的流程里。排查了半天最后发现问题出在线程模型上——多个请求共用同一个 Context 实例后写的数据覆盖了前面的数据。解决方式不复杂Context 必须在每次流程入口处创建绝不能被全局持有。execute_risk_control函数里ctx Context(...)这行就是血的教训换来的。每个请求有自己的上下文天然隔离百试百灵。5.3 重试与幂等节点失败三次的教训节点调用第三方接口偶发超时最开始我的策略是整个流程重跑。结果带来一个大问题重跑会重复执行已经成功的节点。有一次风控节点调用了发送通知短信的服务流程因为短信服务超时挂了重试后投资方收到了两条短信。这个问题的根源是流程重试必须区分已成功的节点和失败的节点。现在的做法是Context 里增加一个executed_nodes列表记录已经成功执行的节点 ID。重试时引擎从失败节点继续执行而不是从起点重新跑。这样既保证了流程可恢复又避免重复副作用。如果你的节点本身不是幂等的一定要在设计阶段就考虑这个机制。5.4 规则优先级的坑另一个让我印象深刻的坑来自规则顺序。早期版本我规定从出边列表里取第一个规则命中的边看起来没问题但工程实现时用的是普通的 list append而我后来为了调试方便在遍历时用了reversed()。结果规则注册顺序一变流程走向就变还特别难发现。从此我定了一条铁律路由规则必须显式声明优先级引擎按优先级降序排序再取第一条命中的规则。这样源码顺序和 list 顺序再也不会影响路由结果配置里看到的 priority 就是最终的执行依据。现在 JSON 配置里每条边都带priority从 5 到 20 不等一眼能看出来哪个规则优先。5.5 性能数据和优化方向ruflo 在单机上的性能我做过一轮简单压测。在 2C4G 容器里一个 200 节点的流程单次执行大约 8~20ms用 100 个并发同时跑CPU 没有明显尖刺。这个数据不作为绝对标准但至少说明纯 Python 实现这种轻量级编排引擎性能完全够用。真正的瓶颈几乎都集中在节点内部的 IO 调用上所以优化重点应该是节点层的超时、缓存和异步化而不是编排引擎本身。后续我计划做三件事。一是并行分支目前 ruflo 是一次走一条分支的模型如果两个节点互相没有依赖理想情况应该并行执行。我的方案是把并行做成一种特殊节点类型内部维护子分支的执行和汇总而不是把整个图调度器改成依赖计数模型。二是表达式引擎替换规则函数让规则也能变成配置里的一段表达式实现完全配置化不用改代码加规则。三是可视化流程编辑因为流程数据本身就是图结构画一个拖拽编辑器只是时间问题届时产品经理自己就能改流程。ruflo 从一个八百行 if-else 的替代品慢慢长成一个有完整设计的小引擎这个过程我自己收获很大。你要是也在被业务分支逻辑折磨建议别急着怼一个复杂的平台先用小工具把流程抽象出来你会发现代码会变得清爽很多。

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价