Java流程引擎案例

wen java案例 2

本文目录导读:

Java流程引擎案例

  1. 整体架构设计
  2. 完整实现代码
  3. 测试案例
  4. 使用说明

我将为您提供一个完整的Java流程引擎案例,基于状态机和责任链模式实现一个灵活的审批流转引擎。

整体架构设计

1 核心组件

// 流程引擎核心接口
public interface WorkflowEngine {
    WorkflowInstance start(String flowKey, Map<String, Object> variables);
    void complete(String instanceId, String currentNodeId, String action, Map<String, Object> variables);
    WorkflowInstance getInstance(String instanceId);
    List<WorkflowTask> getTodoTasks(String assignee);
    void reject(String instanceId, String currentNodeId, String comment);
}

2 核心数据结构

// 流程定义
@Data
public class WorkflowDefinition {
    private String flowKey;          // 流程标识
    private String name;             // 流程名称
    private List<FlowNode> nodes;    // 节点集合
    private Map<String, List<Transition>> transitions; // 节点流转规则
}
// 流程节点
@Data
public class FlowNode {
    private String nodeId;           // 节点ID
    private String nodeType;         // START, APPROVAL, EXCLUSIVE_GATEWAY, END
    private String handlerBeanName;  // 节点处理器
    private Integer timeOut;         // 超时时间(小时)
    private Map<String, Object> properties; // 节点属性
}
// 流转路线
@Data
public class Transition {
    private String fromNodeId;       // 源节点
    private String toNodeId;         // 目标节点
    private String conditionExpression; // 条件表达式 (SpEL)
    private Boolean defaultFlow = false; // 默认路线
    private String action;           // 操作标识:agree/reject
}

完整实现代码

1 流程引擎核心实现

@Service
public class WorkflowEngineImpl implements WorkflowEngine {
    @Autowired
    private WorkflowRegistry registry;           // 流程注册中心
    @Autowired
    private WorkflowInstanceRepository instanceRepo; // 实例仓库
    @Autowired
    private WorkflowTaskRepository taskRepo;     // 任务仓库
    @Autowired
    private ApplicationContext applicationContext; // Spring容器
    @Autowired
    private ExpressionEvaluator evaluator;       // 表达式求值器
    // SpEL表达式基于Spring
    private static final String ACTION_REJECT = "reject";
    @Override
    @Transactional(rollbackFor = Exception.class)
    public WorkflowInstance start(String flowKey, Map<String, Object> variables) {
        WorkflowDefinition definition = registry.getDefinition(flowKey);
        if (definition == null) {
            throw new IllegalArgumentException("流程不存在: " + flowKey);
        }
        // 创建流程实例
        WorkflowInstance instance = new WorkflowInstance();
        instance.setInstanceId(UUID.randomUUID().toString());
        instance.setFlowKey(flowKey);
        instance.setStatus(WorkflowStatus.RUNNING);
        instance.setVariables(variables);
        instance.setCurrentNodeId(getStartNode(definition).getNodeId());
        instance.setStartTime(new Date());
        instance = instanceRepo.save(instance);
        // 执行开始节点
        executeNode(instance, instance.getCurrentNodeId());
        return instance;
    }
    @Override
    @Transactional(rollbackFor = Exception.class)
    public void complete(String instanceId, String currentNodeId, String action, 
                        Map<String, Object> variables) {
        WorkflowInstance instance = getInstance(instanceId);
        if (!instance.getCurrentNodeId().equals(currentNodeId)) {
            throw new IllegalStateException("当前节点不在可操作状态");
        }
        // 合并变量
        if (variables != null) {
            instance.getVariables().putAll(variables);
        }
        // 查找节点绑定和流转
        WorkflowDefinition definition = registry.getDefinition(instance.getFlowKey());
        FloatNode node = findNode(definition, currentNodeId);
        // 校验是否有可执行的transition
        List<Transition> transitions = definition.getTransitions()
            .getOrDefault(currentNodeId, new ArrayList<>())
            .stream()
            .filter(t -> isTransitionExecutable(t, action, instance.getVariables()))
            .collect(Collectors.toList());
        if (transitions.isEmpty()) {
            throw new BusinessException("无满足条件的流转路线");
        }
        // 执行业务handler
        execNodeHandler(node, instance, action);
        // 执行流转
        for (Transition transition : transitions) {
            if (executeTransition(transition, instance)) {
                break;  // 找到匹配的路线并执行
            }
        }
        instanceRepo.save(instance);
    }
    @Override
    public void reject(String instanceId, String currentNodeId, String comment) {
        WorkflowInstance instance = getInstance(instanceId);
        WorkflowDefinition definition = registry.getDefinition(instance.getFlowKey());
        // 查找驳回目标(默认回去上一节点)
        FlowNode currentNode = findNode(definition, currentNodeId);
        FlowNode prevNode = findPreviousNode(definition, currentNodeId);
        // 回到上一节点
        instance.setCurrentNodeId(prevNode.getNodeId());
        instance.getVariables().put("rejectComment", comment);
        // 创建回退任务
        createTask(prevNode.getNodeId(), instance, "待审批");
        instanceRepo.save(instance);
    }
    @Override
    public List<WorkflowTask> getTodoTasks(String assignee) {
        return taskRepo.findByAssigneeAndStatus(assignee, TaskStatus.PENDING);
    }
    // 执行节点逻辑
    private void executeNode(WorkflowInstance instance, String nodeId) {
        WorkflowDefinition definition = registry.getDefinition(instance.getFlowKey());
        FlowNode node = findNode(definition, nodeId);
        switch (node.getNodeType()) {
            case "START":
                // 自动跳转到下一节点
                Transition transition = definition.getTransitions().get(nodeId).get(0);
                executeTransition(transition, instance);
                break;
            case "APPROVAL":
                // 创建审批任务
                createTask(nodeId, instance, "待审批");
                break;
            case "EXCLUSIVE_GATEWAY":
                // 条件网关自动处理
                List<Transition> transitions = definition.getTransitions().get(nodeId);
                for (Transition t : transitions) {
                    if (evaluator.evaluate(t.getConditionExpression(), 
                        instance.getVariables())) {
                        executeTransition(t, instance);
                        break;
                    }
                }
                break;
            case "END":
                instance.setStatus(WorkflowStatus.COMPLETED);
                instance.setEndTime(new Date());
                break;
        }
    }
    // 执行具体流转
    private boolean executeTransition(Transition transition, WorkflowInstance instance) {
        instance.setCurrentNodeId(transition.getToNodeId());
        instance.setLastTransitionTime(new Date());
        // 判断是否是结束节点
        WorkflowDefinition definition = registry.getDefinition(instance.getFlowKey());
        FlowNode targetNode = findNode(definition, transition.getToNodeId());
        if ("END".equals(targetNode.getNodeType())) {
            instance.setStatus(WorkflowStatus.COMPLETED);
            instance.setEndTime(new Date());
        } else {
            executeNode(instance, transition.getToNodeId());
        }
        instanceRepo.save(instance);
        return true;
    }
    // 创建任务
    private void createTask(String nodeId, WorkflowInstance instance, String status) {
        WorkflowTask task = new WorkflowTask();
        task.setTaskId(UUID.randomUUID().toString());
        task.setInstanceId(instance.getInstanceId());
        task.setNodeId(nodeId);
        task.setStatus(status);
        task.setAssignee(getHandlerForNode(nodeId, instance));
        taskRepo.save(task);
    }
    // 获取节点审批人
    private String getHandlerForNode(String nodeId, WorkflowInstance instance) {
        WorkflowDefinition definition = registry.getDefinition(instance.getFlowKey());
        FlowNode node = findNode(definition, nodeId);
        // 支持从变量中动态获取审批人
        Object assignee = instance.getVariables().get(node.getProperties().get("assignee"));
        if (assignee == null) {
            assignee = node.getProperties().get("defaultAssignee");
        }
        return assignee.toString();
    }
    // 判断某一transition是否可执行
    private boolean isTransitionExecutable(Transition transition, String action, 
                                         Map<String, Object> variables) {
        // 默认流转
        if (transition.getDefaultFlow()) {
            return true;
        }
        // action匹配
        String requiredAction = transition.getAction();
        if (requiredAction != null && !requiredAction.equals(action)) {
            return false;
        }
        // 条件表达式
        String condition = transition.getConditionExpression();
        if (condition != null && !evaluator.evaluate(condition, variables)) {
            return false;
        }
        return true;
    }
    // 查找节点
    private FlowNode findNode(WorkflowDefinition definition, String nodeId) {
        return definition.getNodes().stream()
            .filter(n -> n.getNodeId().equals(nodeId))
            .findFirst()
            .orElseThrow(() -> new IllegalArgumentException("节点不存在: " + nodeId));
    }
    // 查找开始节点
    private FlowNode getStartNode(WorkflowDefinition definition) {
        return definition.getNodes().stream()
            .filter(n -> "START".equals(n.getNodeType()))
            .findFirst()
            .orElseThrow(() -> new IllegalArgumentException("流程缺少开始节点"));
    }
    // 查找前序节点
    private FloatNode findPreviousNode(WorkflowDefinition definition, String nodeId) {
        // 遍历所有transitions找到能到达currentNodeId的源节点
        return definition.getTransitions().entrySet().stream()
            .flatMap(entry -> entry.getValue().stream()
                .filter(t -> t.getToNodeId().equals(nodeId))
                .map(t -> findNode(definition, entry.getKey())))
            .findFirst()
            .orElseThrow(() -> new IllegalArgumentException("未找到前序节点"));
    }
    // 调用节点处理器
    private void execNodeHandler(FlowNode node, WorkflowInstance instance, String action) {
        if (node.getHandlerBeanName() != null) {
            try {
                FlowNodeHandler handler = (FlowNodeHandler) 
                    applicationContext.getBean(node.getHandlerBeanName());
                handler.handle(instance, action);
            } catch (Exception e) {
                throw new BusinessException("节点处理器执行失败", e);
            }
        }
    }
    // 保存任务状态
    public void completeTask(String taskId, String action, String comment) {
        WorkflowTask task = taskRepo.findById(taskId).orElseThrow();
        task.setStatus(TaskStatus.COMPLETED);
        task.setComment(comment);
        task.setCompleteTime(new Date());
        taskRepo.save(task);
    }
}

2 流程注册中心

@Component
public class WorkflowRegistry {
    private final Map<String, WorkflowDefinition> definitions = new ConcurrentHashMap<>();
    public void register(WorkflowDefinition definition) {
        definitions.put(definition.getFlowKey(), definition);
        validateDefinition(definition);
    }
    public WorkflowDefinition getDefinition(String flowKey) {
        return definitions.get(flowKey);
    }
    public void validateDefinition(WorkflowDefinition definition) {
        // 至少有一个开始节点和一个结束节点
        long startCount = definition.getNodes().stream()
            .filter(n -> "START".equals(n.getNodeType())).count();
        long endCount = definition.getNodes().stream()
            .filter(n -> "END".equals(n.getNodeType())).count();
        if (startCount != 1) {
            throw new IllegalArgumentException("必须有且仅有一个开始节点");
        }
        if (endCount < 1) {
            throw new IllegalArgumentException("至少需要一个结束节点");
        }
        // 必须有从START的流向
        if (!definition.getTransitions().containsKey("START")) {
            throw new IllegalArgumentException("START节点必须要有流出路径");
        }
    }
}

3 表达式求值器

@Component
public class SpelExpressionEvaluator implements ExpressionEvaluator {
    private final ExpressionParser parser = new SpelExpressionParser();
    @Override
    public boolean evaluate(String expression, Map<String, Object> variables) {
        if (expression == null || expression.trim().isEmpty()) {
            return true;
        }
        EvaluationContext context = new StandardEvaluationContext();
        context.setVariables(new VariableMap<>(variables));
        Expression exp = parser.parseExpression(expression);
        return Boolean.TRUE.equals(exp.getValue(context, Boolean.class));
    }
}

4 节点处理器接口

public interface FlowNodeHandler {
    void handle(WorkflowInstance instance, String action);
}
// 示例:请假审批处理器
@Component("leaveApprovalHandler")
public class LeaveApprovalHandler implements FlowNodeHandler {
    @Override
    public void handle(WorkflowInstance instance, String action) {
        Map<String, Object> variables = instance.getVariables();
        Integer days = (Integer) variables.get("days");
        String applicant = (String) variables.get("applicant");
        // 三天以上自动通知上级
        if (days > 3) {
            sendEmail(applicant, "您的请假申请超过3天,需要部门经理审批", null);
        }
    }
    private void sendEmail(String to, String subject, String content) {
        // 邮件发送实现
    }
}

5 仓库接口

public interface WorkflowInstanceRepository extends JpaRepository<WorkflowInstance, String> {
    List<WorkflowInstance> findByStatus(WorkflowStatus status);
    Optional<WorkflowInstance> findByBusinessKey(String businessKey);
}
public interface WorkflowTaskRepository extends JpaRepository<WorkflowTask, String> {
    List<WorkflowTask> findByAssigneeAndStatus(String assignee, TaskStatus status);
    List<WorkflowTask> findByInstanceId(String instanceId);
}

6 实体类

@Entity
@Data
public class WorkflowInstance {
    @Id
    private String instanceId;
    private String flowKey;
    private String businessKey;
    private WorkflowStatus status;
    private String currentNodeId;
    @Convert(converter = MapToStringConverter.class)
    private Map<String, Object> variables;
    private Date startTime;
    private Date endTime;
    private Date lastTransitionTime;
    @Version
    private Long version;
}
@Entity
@Data
public class WorkflowTask {
    @Id
    private String taskId;
    private String instanceId;
    private String nodeId;
    private String assignee;
    private TaskStatus status;
    private String comment;
    private Date createTime;
    private Date completeTime;
}

测试案例

1 测试流程配置

@Configuration
public class WorkflowConfig {
    @Bean
    public WorkflowRegistry workflowRegistry() {
        WorkflowRegistry registry = new WorkflowRegistry();
        // 请假流程
        WorkflowDefinition leaveFlow = buildLeaveFlow();
        registry.register(leaveFlow);
        return registry;
    }
    private WorkflowDefinition buildLeaveFlow() {
        WorkflowDefinition def = new WorkflowDefinition();
        def.setFlowKey("leave_process");
        def.setName("请假审批流程");
        // 节点
        List<FlowNode> nodes = new ArrayList<>();
        FlowNode start = new FlowNode();
        start.setNodeId("start");
        start.setNodeType("START");
        nodes.add(start);
        FlowNode apply = new FlowNode();
        apply.setNodeId("apply");
        apply.setNodeType("APPROVAL");
        apply.setHandlerBeanName("leaveApplyHandler");
        nodes.add(apply);
        FlowNode managerReview = new FlowNode();
        managerReview.setNodeId("manager_review");
        managerReview.setNodeType("APPROVAL");
        managerReview.setHandlerBeanName("leaveApprovalHandler");
        nodes.add(managerReview);
        FlowNode hrReview = new FlowNode();
        hrReview.setNodeId("hr_review");
        hrReview.setNodeType("APPROVAL");
        hrReview.setHandlerBeanName("hrApprovalHandler");
        nodes.add(hrReview);
        FlowNode end = new FlowNode();
        end.setNodeId("end");
        end.setNodeType("END");
        nodes.add(end);
        def.setNodes(nodes);
        // 流转规则
        Map<String, List<Transition>> transitions = new HashMap<>();
        Transition start2Apply = new Transition();
        start2Apply.setFromNodeId("start");
        start2Apply.setToNodeId("apply");
        transitions.put("start", Collections.singletonList(start2Apply));
        Transition apply2Manager = new Transition();
        apply2Manager.setFromNodeId("apply");
        apply2Manager.setToNodeId("manager_review");
        apply2Manager.setAction("submit");
        transitions.put("apply", Collections.singletonList(apply2Manager));
        // 条件分支
        Transition manager2Hr = new Transition();
        manager2Hr.setFromNodeId("manager_review");
        manager2Hr.setToNodeId("hr_review");
        manager2Hr.setConditionExpression("#days >= 3");
        Transition manager2End = new Transition();
        manager2End.setFromNodeId("manager_review");
        manager2End.setToNodeId("end");
        manager2End.setConditionExpression("#days < 3");
        manager2End.setDefaultFlow(true);
        List<Transition> managerTransitions = new ArrayList<>();
        managerTransitions.add(manager2Hr);
        managerTransitions.add(manager2End);
        transitions.put("manager_review", managerTransitions);
        Transition hr2End = new Transition();
        hr2End.setFromNodeId("hr_review");
        hr2End.setToNodeId("end");
        transitions.put("hr_review", Collections.singletonList(hr2End));
        def.setTransitions(transitions);
        return def;
    }
}

2 测试代码

@SpringBootTest
public class WorkflowEngineTest {
    @Autowired
    private WorkflowEngine engine;
    @Test
    public void testLeaveWorkflow() {
        // 启动流程
        Map<String, Object> variables = new HashMap<>();
        variables.put("applicant", "张三");
        variables.put("days", 5);
        WorkflowInstance instance = engine.start("leave_process", variables);
        String instanceId = instance.getInstanceId();
        // 企业审批节点完成
        engine.complete(instanceId, "apply", "submit", null);
        // 校验当前节点
        WorkflowInstance updated = engine.getInstance(instanceId);
        assertEquals("manager_review", updated.getCurrentNodeId());
        // 经理审批(同意)
        engine.complete(instanceId, "manager_review", "agree", null);
        // 校验进入HR审批
        updated = engine.getInstance(instanceId);
        assertEquals("hr_review", updated.getCurrentNodeId());
        // HR审批
        engine.complete(instanceId, "hr_review", "approve", null);
        // 校验流程完成
        updated = engine.getInstance(instanceId);
        assertEquals(WorkflowStatus.COMPLETED, updated.getStatus());
    }
    @Test
    public void testRejectWorkflow() {
        Map<String, Object> variables = Map.of("applicant", "李四", "days", 2);
        WorkflowInstance instance = engine.start("leave_process", variables);
        // 经理驳回
        engine.reject(instance.getInstanceId(), "manager_review", "请假理由不充分");
        WorkflowInstance updated = engine.getInstance(instance.getInstanceId());
        assertEquals("apply", updated.getCurrentNodeId());
        assertNotNull(updated.getVariables().get("rejectComment"));
    }
    @Test
    public void testConditionalFlow() {
        // 请假1天,应该直接结束
        Map<String, Object> variables = Map.of("applicant", "王五", "days", 1);
        WorkflowInstance instance = engine.start("leave_process", variables);
        engine.complete(instance.getInstanceId(), "apply", "submit", null);
        WorkflowInstance updated = engine.getInstance(instance.getInstanceId());
        // 等待manager和HR节点处理
        // 由于条件, 经理审批后直接到结束
        engine.complete(instance.getInstanceId(), "manager_review", "agree", null);
        updated = engine.getInstance(instance.getInstanceId());
        assertEquals(WorkflowStatus.COMPLETED, updated.getStatus());
    }
}

使用说明

1 核心特点

  • 动态决策:基于SpEL表达式实现复杂的流转逻辑
  • 可扩展性:通过FlowNodeHandler接口扩展业务节点
  • 条件分支:支持条件网关实现多分支流转
  • 驳回机制:支持任务驳回、重新提交
  • 事务管理:整个流程操作原子性保证

2 优化建议

  • 增加流程版本管理
  • 添加流程跟踪(审计日志)
  • 增加超时提醒、自动任务
  • 提供图形化流程设计器
  • 增加流程催办、代理功能

3 使用示例

// 获取待办任务
List<WorkflowTask> tasks = engine.getTodoTasks("张三");
// 提交申请
Map<String, Object> variables = new HashMap<>();
variables.put("amount", 5000);
engine.complete(instanceId, "apply", "submit", variables);
// 审批通过
engine.complete(instanceId, "manager_review", "agree", null);
// 审批拒绝
engine.reject(instanceId, "manager_review", "预算不足");

这个流程引擎提供了完整的审批流转功能,可根据实际需求进行扩展和调整。

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