审批流分布式动态节点驳回

wen java案例 2

本文目录导读:

审批流分布式动态节点驳回

  1. 核心概念定义
  2. 动态节点驳回的三大技术难题
  3. 技术架构方案(最佳实践)
  4. 核心代码逻辑 (Java + Spring Boot 伪代码)
  5. 画图示意流程
  6. 重要建议与陷阱

这是一个非常专业的BPM(业务流程管理)技术问题,在分布式、微服务架构下实现“审批流动态节点驳回”,核心难点在于状态一致性动态节点实例化以及跨服务的事务回滚

下面从核心概念技术实现方案代码架构设计以及核心逻辑几个层面进行深度拆解。

核心概念定义

  • 动态节点: 不是预先在流程定义(BPMN)中固定好的,而是在运行时根据变量、脚本、外部服务返回的结果或人员池动态生成的节点。
    • 根据报销金额确定审批人(总经理 or 部门经理)。
    • 会签节点根据子部门数量动态生成多个审批子任务。
  • 分布式: 审批流引擎本身是一个微服务,而执行驳回操作可能涉及多个独立的微服务(如:订单服务、财务服务、客户服务)的状态回滚。
  • 驳回: 将流程从当前节点(或后置节点)回退到指定的历史节点(或流程发起人处),并且需要恢复该历史节点及其后续节点的执行上下文。

动态节点驳回的三大技术难题

  1. 动态节点实例的“复活”
    • 静态节点驳回,通常只要找到历史任务ID,重新激活即可。
    • 动态节点,原始审批人、候选人或审批策略可能已经发生了变化(A部门经理离职了,但流程要驳回给曾经审批过的A部门经理),此时需要动态重建或继承历史上下文。
  2. 状态补偿(Saga vs TCC)
    • 驳回不仅仅是流程引擎的节点跳转,往往意味着业务状态回滚,驳回订单审批,需要将订单状态从“审批中”变回“待提交”,同时可能需要释放库存锁定。
    • 难点:在分布式环境下,一旦驳回到了第3个节点,前面2个节点对应的业务微服务(订单、库存)已完成操作,如何保证回滚不出错?
  3. 多实例任务的并发处理

    如果是动态生成的“会签”节点(需要5个人审批,目前只有3人通过,2人未处理),驳回时不是简单地撤回已通过的,而是要终止未处理的,并重置所有状态。

技术架构方案(最佳实践)

推荐使用 事件溯源 + 状态机 + 分布式事务补偿 的方案。

存储模型:事件溯源的Token模型

不直接存储“当前节点”,而是存储一系列已执行的事件

  • Token(令牌): 流程实例中的移动指针。
  • Event(事件)TaskCreatedTaskCompletedTaskRejected
  • 核心优势: 对于动态节点驳回,我们不修改历史,而是追加一个 FlowJumpedEvent 事件,通过重算Token Map,来告诉引擎当前应该停在哪个节点。
-- 简化表结构
CREATE TABLE process_event (
    id BIGINT AUTO_INCREMENT,
    process_instance_id VARCHAR(64),
    event_type VARCHAR(50), -- TaskAssigned, FlowJumpOut, FlowJumpIn
    node_id VARCHAR(64), -- 动态节点ID
    payload JSON, -- 包含动态分配的信息(人员、角色、参数)
    timestamp BIGINT
);

实现流程:基于Activity/Flyway思想的CQRS

  • Command: 接受“驳回”请求,包含 目标节点ID是否强制继承当前上下文
  • Handler
    1. 反查事件流,找到从目标节点到当前节点的所有动态节点。
    2. 生成 FlowJumpedEvent 将Token重新指向目标节点。
    3. 发布 回滚命令 给每个被跳过的节点对应的业务服务。
  • Query: 直接基于事件流计算当前有效的Token位置。

核心代码逻辑 (Java + Spring Boot 伪代码)

假设我们有一个审批流引擎服务 ProcessEngineService

动态节点上下文定义

// 动态节点的运行时上下文
@Data
public class DynamicNodeContext {
    private String nodeId;          //  "level1_approval"
    private String assignee;        // 具体的审批人
    private String dynamicSource;   // 产生这个节点的来源(如:规则引擎返回的部门ID)
    private String ruleExpression;  // 产生动态的规则表达式
    private boolean isRollback;     // 当前是否为回滚状态
    private Map<String, Object> originalVariables; // 回滚时需要恢复的变量
}

核心驳回服务 (支持动态节点补偿)

@Service
public class DynamicNodeRejectService {
    @Autowired
    private EventStore eventStore;
    @Autowired
    private TaskInstanceService taskInstanceService;
    @Autowired
    private SagaManager sagaManager; // 分布式事务补偿管理器
    public ProcessResult rejectFlow(String processInstanceId, String targetNodeId) {
        // 1. 调取事件流,重建当前流程的Token状态
        ProcessToken token = eventStore.rebuildToken(processInstanceId);
        // 2. 计算跳转路径 (当前节点 -> ... -> 目标节点)
        //    注意:这些节点可能是动态生成的,需要从Token的历史映射中查找
        List<String> skippedNodeIds = token.calculateSkippedNodes(targetNodeId);
        // 3. 创建Saga事务 (分布式补偿)
        Saga saga = sagaManager.begin(processInstanceId);
        try {
            // 4. 对每个被跳过的动态节点执行"回滚"动作
            for (String nodeId : skippedNodeIds) {
                DynamicNodeContext ctx = token.getDynamicNodeContext(nodeId);
                if (ctx != null) {
                    // 调用各业务微服务进行状态回滚
                    //  对于"订单锁定库存"节点,这里触发库存解锁
                    saga.addCompensationStep(
                        () -> notifyBusinessService(nodeId, "ROLLBACK", ctx)
                    );
                }
            }
            // 5. 终止当前所有活跃的任务实例 (包括动态生成的多实例)
            List<TaskInstance> activeTasks = taskInstanceService
                .getActiveTasksByProcessId(processInstanceId);
            for (TaskInstance task : activeTasks) {
                taskInstanceService.forceTerminate(task.getId());
            }
            // 6. 创建"回退"事件并写入事件存储
            //    关键点:将动态节点上下文恢复到目标节点
            eventStore.appendEvent(ProcessEvent.builder()
                .processInstanceId(processInstanceId)
                .eventType("FLOW_JUMP_TO")
                .targetNodeId(targetNodeId)
                .dynamicContextSnapshots(token.getDynamicContextsUpTo(targetNodeId))
                .build());
            // 7. 重建目标节点的任务 (如果是动态节点,需要重新解析规则)
            DynamicNodeContext targetCtx = token.getDynamicNodeContext(targetNodeId);
            if (targetCtx != null && targetCtx.getRuleExpression() != null) {
                // 重新评估动态规则(部门经理变了,要用最新的人员池)
                String newAssignee = ruleService.evaluate(
                    targetCtx.getRuleExpression(),
                    targetCtx.getOriginalVariables()
                );
                taskInstanceService.createDynamicTask(
                    targetNodeId, newAssignee, processInstanceId
                );
            }
            saga.commit();
            return ProcessResult.success();
        } catch (Exception e) {
            // 如果补偿失败,Saga会自动调用前面成功步骤的逆向操作
            saga.compensate(); 
            return ProcessResult.fail("驳回失败,已触发补偿,错误:" + e.getMessage());
        }
    }
    // 补偿业务逻辑 (分布式调用)
    private void notifyBusinessService(String nodeId, String action, DynamicNodeContext ctx) {
        // 通过消息队列或Feign调用业务微服务
        businessClient.rollbackNodeAction(
            new RollbackRequest(nodeId, ctx.getOriginalVariables())
        );
    }
}

关键数据处理:Token重建

这是处理动态节点的核心。

// Token 重建逻辑 (从事件流)
public ProcessToken rebuildToken(String processInstanceId) {
    List<ProcessEvent> events = eventStore.getEventsByProcessId(processInstanceId);
    ProcessToken token = new ProcessToken();
    Map<String, DynamicNodeContext> contextMap = new HashMap<>();
    for (ProcessEvent event : events) {
        switch (event.getEventType()) {
            case "TASK_CREATED":
                // 解析动态生成的节点上下文
                DynamicNodeContext ctx = parseDynamicContext(event);
                contextMap.put(event.getNodeId(), ctx);
                break;
            case "TASK_COMPLETED":
                // 记录完成状态
                token.addCompletedNode(event.getNodeId());
                break;
            case "FLOW_JUMP_TO":
                // 如果是回退事件,截断后续所有节点状态
                List<String> activeList = token.getActiveNodes();
                activeList.removeIf(n -> n.compareTo(event.getTargetNodeId()) > 0);
                // 同时恢复动态上下文
                Map<String, DynamicNodeContext> snapshot = event.getDynamicContextSnapshots();
                contextMap.putAll(snapshot);
                break;
        }
    }
    token.setDynamicContexts(contextMap);
    return token;
}

画图示意流程

场景: 动态节点A (经理审批) -> 动态节点B (总监审批) -> 动态节点C (VP审批),现在需要从C驳回至A。

graph TD
    subgraph 流程状态
        A[动态节点A: 经理审批<br>Token: 完成] --> B
        B[动态节点B: 总监审批<br>Token: 完成] --> C
        C[动态节点C: VP审批<br>Token: 正在审批] -- 驳回操作触发 --> RejectHandler
    end
    subgraph 核心处理逻辑 RejectHandler
        direction LR
        R1[1. 读取事件流<br>重建Token状态: {A:completed, B:completed, C:active}]
        R2[2. 计算跳过节点: [B, C]]
        R3[3. 发起Saga补偿]
        R4[4. 发送补偿命令]
        R5[5. 终止C节点的任务]
        R6[6. 追加FLOW_JUMP_TO事件]
        R7[7. 根据A节点的Rule Expression<br>重新创建动态任务]
    end
    subgraph 分布式补偿
        S1[订单服务:<br>回退库存锁定]
        S2[财务服务:<br>取消预算占用]
        S3[客户服务:<br>撤回通知]
    end
    RejectHandler --> R1 --> R2 --> R3 --> R4
    R4 --> S1
    R4 --> S2
    R4 --> S3
    R3 --> R5 --> R6 --> R7
    R7 --> Done[流程回到A节点<br>重新等待经理审批]

重要建议与陷阱

  1. 避免强一致,拥抱最终一致
    • 不要试图用分布式事务(XA)锁住所有微服务,使用Saga或TCC模式,并引入消息表本地消息表来保证补偿操作一定被执行。
  2. 动态节点规则的幂等性与重新评估
    • 驳回后,节点A的审批人可能已经不是原来的那个人(原经理离职了)。
    • 建议策略:提供“保留历史审批人”和“重新按照当前规则分配”两种模式,系统应该允许用户在发起驳回时选择。
  3. 分支与聚合模型
    • 如果流程中有“并行分支(Parallel Gateway)”,动态节点驳回会变得极其复杂,需要同时考虑所有分支的Token状态。
    • 简化方案:对于复杂的并行动态驳回,建议设计为“全局驳回到发起人”“仅驳回到前一个并行网关”,避免陷入令牌匹配地狱。
  4. 性能考虑
    • 大量的事件重演(Event Sourcing)可能在超长流程中变慢,建议引入快照(Snapshot)机制,每N个事件存储一份Token状态的完整快照。

“审批流分布式动态节点驳回” 本质上是一个 状态恢复(State Recovery) 问题,最佳实践是:

  1. 采用 事件溯源 动态存储每个节点的分配上下文。
  2. 通过 Saga 补偿 处理跨微服务的业务回滚。
  3. 使用 Token 重算 机制来精准决定流程应该跳转到哪个动态节点。

这样设计,即使在高度动态、分布式的微服务环境下,也能保证审批回滚的逻辑正确性和数据一致性。

抱歉,评论功能暂时关闭!