本文目录导读:

我将为您提供一个完整的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", "预算不足");
这个流程引擎提供了完整的审批流转功能,可根据实际需求进行扩展和调整。