本文目录导读:

我将为您提供一个完整的工作流引擎Java实现案例,包含核心功能演示。
项目结构
workflow-demo/
├── src/main/java/com/workflow/
│ ├── model/
│ │ ├── WorkflowDefinition.java
│ │ ├── WorkflowInstance.java
│ │ ├── TaskInstance.java
│ │ └── WorkflowContext.java
│ ├── engine/
│ │ ├── WorkflowEngine.java
│ │ ├── WorkflowParser.java
│ │ └── TaskExecutor.java
│ ├── handler/
│ │ ├── TaskHandler.java
│ │ └── DefaultTaskHandler.java
│ ├── service/
│ │ └── WorkflowService.java
│ └── demo/
│ └── WorkflowDemo.java
核心模型类
WorkflowDefinition.java
package com.workflow.model;
import java.util.*;
import java.util.concurrent.ConcurrentHashMap;
/**
* 工作流定义
*/
public class WorkflowDefinition {
private String id;
private String name;
private String version;
private Map<String, Node> nodes = new LinkedHashMap<>();
private List<Transition> transitions = new ArrayList<>();
private String startNodeId;
private List<String> endNodeIds;
public WorkflowDefinition(String id, String name, String version) {
this.id = id;
this.name = name;
this.version = version;
}
public void addNode(Node node) {
nodes.put(node.getId(), node);
}
public void addTransition(Transition transition) {
transitions.add(transition);
}
// Getters and Setters
public String getId() { return id; }
public String getName() { return name; }
public String getVersion() { return version; }
public Map<String, Node> getNodes() { return nodes; }
public List<Transition> getTransitions() { return transitions; }
public String getStartNodeId() { return startNodeId; }
public void setStartNodeId(String startNodeId) { this.startNodeId = startNodeId; }
public List<String> getEndNodeIds() { return endNodeIds; }
public void setEndNodeIds(List<String> endNodeIds) { this.endNodeIds = endNodeIds; }
/**
* 节点定义
*/
public static class Node {
private String id;
private String name;
private NodeType type;
private Map<String, Object> properties = new HashMap<>();
public Node(String id, String name, NodeType type) {
this.id = id;
this.name = name;
this.type = type;
}
// Getters and Setters
public String getId() { return id; }
public String getName() { return name; }
public NodeType getType() { return type; }
public Map<String, Object> getProperties() { return properties; }
public void setProperty(String key, Object value) { properties.put(key, value); }
}
/**
* 转换定义(边的定义)
*/
public static class Transition {
private String fromNodeId;
private String toNodeId;
private String condition; // 条件表达式
public Transition(String fromNodeId, String toNodeId) {
this(fromNodeId, toNodeId, null);
}
public Transition(String fromNodeId, String toNodeId, String condition) {
this.fromNodeId = fromNodeId;
this.toNodeId = toNodeId;
this.condition = condition;
}
public String getFromNodeId() { return fromNodeId; }
public String getToNodeId() { return toNodeId; }
public String getCondition() { return condition; }
}
/**
* 节点类型
*/
public enum NodeType {
START, // 开始节点
TASK, // 任务节点
DECISION, // 决策节点
PARALLEL, // 并行任务
END // 结束节点
}
}
WorkflowInstance.java
package com.workflow.model;
import java.util.*;
import java.util.concurrent.ConcurrentHashMap;
/**
* 工作流实例
*/
public class WorkflowInstance {
private String id;
private String definitionId;
private WorkflowStatus status;
private String currentNodeId;
private Map<String, Object> variables = new ConcurrentHashMap<>();
private Map<String, TaskInstance> tasks = new ConcurrentHashMap<>();
private Date startTime;
private Date endTime;
private Map<String, NodeExecutionRecord> executionRecords = new LinkedHashMap<>();
public enum WorkflowStatus {
RUNNING, COMPLETED, FAILED, TERMINATED, PAUSED
}
public WorkflowInstance(String id, String definitionId) {
this.id = id;
this.definitionId = definitionId;
this.status = WorkflowStatus.RUNNING;
this.startTime = new Date();
}
public void addTask(TaskInstance task) {
tasks.put(task.getId(), task);
}
public void updateVariable(String key, Object value) {
variables.put(key, value);
}
public Object getVariable(String key) {
return variables.get(key);
}
public void addExecutionRecord(NodeExecutionRecord record) {
executionRecords.put(record.getNodeId() + "_" + record.getExecutionTime(), record);
}
public boolean hasActiveTasks() {
return tasks.values().stream()
.anyMatch(t -> t.getStatus() == TaskInstance.TaskStatus.RUNNING);
}
public int getTotalTaskCount() {
return tasks.size();
}
public int getCompletedTaskCount() {
return (int) tasks.values().stream()
.filter(t -> t.getStatus() == TaskInstance.TaskStatus.COMPLETED)
.count();
}
// Getters and Setters
public String getId() { return id; }
public String getDefinitionId() { return definitionId; }
public WorkflowStatus getStatus() { return status; }
public void setStatus(WorkflowStatus status) { this.status = status; }
public String getCurrentNodeId() { return currentNodeId; }
public void setCurrentNodeId(String currentNodeId) { this.currentNodeId = currentNodeId; }
public Map<String, Object> getVariables() { return variables; }
public Map<String, TaskInstance> getTasks() { return tasks; }
public Date getStartTime() { return startTime; }
public Date getEndTime() { return endTime; }
public void setEndTime(Date endTime) { this.endTime = endTime; }
public Map<String, NodeExecutionRecord> getExecutionRecords() { return executionRecords; }
/**
* 节点执行记录
*/
public static class NodeExecutionRecord {
private String nodeId;
private String nodeName;
private Date executionTime;
private NodeExecutionStatus status;
private String details;
public NodeExecutionRecord(String nodeId, String nodeName, NodeExecutionStatus status) {
this.nodeId = nodeId;
this.nodeName = nodeName;
this.executionTime = new Date();
this.status = status;
}
public String getNodeId() { return nodeId; }
public String getNodeName() { return nodeName; }
public Date getExecutionTime() { return executionTime; }
public NodeExecutionStatus getStatus() { return status; }
public String getDetails() { return details; }
public void setDetails(String details) { this.details = details; }
public enum NodeExecutionStatus {
PROCESSED, SKIPPED, FAILED
}
}
}
TaskInstance.java
package com.workflow.model;
import java.util.Date;
import java.util.HashMap;
import java.util.Map;
/**
* 任务实例
*/
public class TaskInstance {
private String id;
private String nodeId;
private String nodeName;
private String workflowInstanceId;
private TaskStatus status;
private String assignee;
private Date createdTime;
private Date claimTime;
private Date completedTime;
private String variables;
private Map<String, Object> taskData = new HashMap<>();
public enum TaskStatus {
CREATED, ASSIGNED, RUNNING, COMPLETED, CANCELLED
}
public TaskInstance(String id, String nodeId, String nodeName, String workflowInstanceId) {
this.id = id;
this.nodeId = nodeId;
this.nodeName = nodeName;
this.workflowInstanceId = workflowInstanceId;
this.status = TaskStatus.CREATED;
this.createdTime = new Date();
}
// Getters and Setters
public String getId() { return id; }
public String getNodeId() { return nodeId; }
public String getNodeName() { return nodeName; }
public String getWorkflowInstanceId() { return workflowInstanceId; }
public TaskStatus getStatus() { return status; }
public void setStatus(TaskStatus status) { this.status = status; }
public String getAssignee() { return assignee; }
public void setAssignee(String assignee) { this.assignee = assignee; }
public Date getCreatedTime() { return createdTime; }
public Date getClaimTime() { return claimTime; }
public void setClaimTime(Date claimTime) { this.claimTime = claimTime; }
public Date getCompletedTime() { return completedTime; }
public void setCompletedTime(Date completedTime) { this.completedTime = completedTime; }
public Map<String, Object> getTaskData() { return taskData; }
public void setTaskData(Map<String, Object> taskData) { this.taskData = taskData; }
}
工作流引擎
WorkflowEngine.java
package com.workflow.engine;
import com.workflow.handler.TaskHandler;
import com.workflow.model.*;
/**
* 工作流引擎(核心)
*/
public class WorkflowEngine {
private final WorkflowParser parser;
private final TaskExecutor taskExecutor;
private boolean running = true;
public WorkflowEngine(WorkflowParser parser, TaskExecutor taskExecutor) {
this.parser = parser;
this.taskExecutor = taskExecutor;
}
/**
* 启动工作流
*/
public WorkflowInstance startWorkflow(WorkflowDefinition definition, Map<String, Object> variables) {
if (definition == null) {
throw new IllegalArgumentException("工作流定义不能为空");
}
// 创建工作流实例
String instanceId = generateInstanceId(definition.getId());
WorkflowInstance instance = new WorkflowInstance(instanceId, definition.getId());
// 初始化流程变量
if (variables != null) {
variables.forEach(instance::updateVariable);
}
try {
System.out.println("=== 启动工作流: " + definition.getName() + " (版本: " + definition.getVersion() + ") ===");
// 执行开始节点链
String currentNodeId = definition.getStartNodeId();
instance.setCurrentNodeId(currentNodeId);
while (running && currentNodeId != null) {
currentNodeId = processNode(instance, definition, currentNodeId);
}
if (instance.getStatus() == WorkflowInstance.WorkflowStatus.RUNNING) {
// 检查是否所有节点都已完成
boolean allCompleted = checkAllNodesCompleted(definition);
if (allCompleted) {
completeWorkflow(instance);
}
}
} catch (Exception e) {
System.err.println("工作流执行失败: " + e.getMessage());
e.printStackTrace();
instance.setStatus(WorkflowInstance.WorkflowStatus.FAILED);
}
return instance;
}
/**
* 处理单个节点
*/
private String processNode(WorkflowInstance instance, WorkflowDefinition definition, String nodeId) {
WorkflowDefinition.Node node = definition.getNodes().get(nodeId);
if (node == null) {
throw new RuntimeException("节点不存在: " + nodeId);
}
System.out.println("\n▶ 执行节点: " + node.getName() + " (" + node.getType() + ")");
switch (node.getType()) {
case START:
return executeStartNode(instance, definition, node);
case TASK:
return executeTaskNode(instance, definition, node);
case DECISION:
return executeDecisionNode(instance, definition, node);
case PARALLEL:
return executeParallelNode(instance, definition, node);
case END:
return executeEndNode(instance, definition, node);
default:
throw new RuntimeException("未知节点类型: " + node.getType());
}
}
/**
* 执行开始节点
*/
private String executeStartNode(WorkflowInstance instance, WorkflowDefinition definition, WorkflowDefinition.Node node) {
instance.setCurrentNodeId(node.getId());
recordNodeExecution(instance, node, "开始节点,初始化流程");
// 找到下一个节点
return findNextNode(definition, node.getId());
}
/**
* 执行任务节点
*/
private String executeTaskNode(WorkflowInstance instance, WorkflowDefinition definition, WorkflowDefinition.Node node) {
// 创建任务实例
TaskInstance task = new TaskInstance(
generateTaskId(instance.getId()),
node.getId(),
node.getName(),
instance.getId()
);
// 设置任务数据
Object taskData = node.getProperties().get("taskData");
if (taskData instanceof Map) {
((Map<String, Object>) taskData).forEach(task.getTaskData()::put);
}
// 获取任务处理器
TaskHandler handler = taskExecutor.getHandler(node.getId());
if (handler == null) {
handler = new DefaultTaskHandler();
}
// 执行任务
try {
System.out.println(" 处理任务: " + task.getId());
handler.execute(task, instance);
task.setStatus(TaskInstance.TaskStatus.COMPLETED);
task.setCompletedTime(new Date());
instance.addTask(task);
recordNodeExecution(instance, node, "任务节点执行完成: " + node.getName());
} catch (Exception e) {
task.setStatus(TaskInstance.TaskStatus.CANCELLED);
instance.addTask(task);
recordNodeExecution(instance, node, "任务节点执行失败: " + e.getMessage());
throw e;
}
return findNextNode(definition, node.getId());
}
/**
* 执行决策节点
*/
private String executeDecisionNode(WorkflowInstance instance, WorkflowDefinition definition, WorkflowDefinition.Node node) {
// 获取决策结果
Object decisionValue = instance.getVariable("decision_result");
if (decisionValue == null) {
decisionValue = "approve"; // 默认决策
}
String decision = decisionValue.toString();
System.out.println(" 决策结果: " + decision);
// 根据决策找到下一个节点
List<WorkflowDefinition.Transition> transitions = definition.getTransitions().stream()
.filter(t -> t.getFromNodeId().equals(node.getId()))
.collect(Collectors.toList());
for (WorkflowDefinition.Transition transition : transitions) {
String condition = transition.getCondition();
if (condition == null || condition.equals(decision)) {
recordNodeExecution(instance, node, "决策节点执行,选择: " + decision);
return transition.getToNodeId();
}
}
throw new RuntimeException("决策节点没有找到匹配的分支");
}
/**
* 执行并行节点
*/
private String executeParallelNode(WorkflowInstance instance, WorkflowDefinition definition, WorkflowDefinition.Node node) {
// 找到所有并行的分支
List<String> parallelBranches = definition.getTransitions().stream()
.filter(t -> t.getFromNodeId().equals(node.getId()))
.map(WorkflowDefinition.Transition::getToNodeId)
.collect(Collectors.toList());
System.out.println(" 并行执行 " + parallelBranches.size() + " 个分支");
// 模拟并行执行(实际实现中会使用多线程)
for (String branchNodeId : parallelBranches) {
processNode(instance, definition, branchNodeId);
}
return findNextNode(definition, node.getId());
}
/**
* 执行结束节点
*/
private String executeEndNode(WorkflowInstance instance, WorkflowDefinition definition, WorkflowDefinition.Node node) {
instance.setCurrentNodeId(node.getId());
recordNodeExecution(instance, node, "流程结束");
completeWorkflow(instance);
return null; // 返回null表示流程结束
}
/**
* 查找下一个节点
*/
private String findNextNode(WorkflowDefinition definition, String currentNodeId) {
return definition.getTransitions().stream()
.filter(t -> t.getFromNodeId().equals(currentNodeId) && t.getCondition() == null)
.findFirst()
.map(WorkflowDefinition.Transition::getToNodeId)
.orElse(null);
}
/**
* 完成整个工作流
*/
private void completeWorkflow(WorkflowInstance instance) {
instance.setStatus(WorkflowInstance.WorkflowStatus.COMPLETED);
instance.setEndTime(new Date());
System.out.println("\n=== 工作流完成: " + instance.getId() + " ===");
System.out.println(" 总任务数: " + instance.getCompletedTaskCount() + "/" + instance.getTotalTaskCount());
System.out.println(" 执行记录数: " + instance.getExecutionRecords().size());
}
/**
* 检查所有节点是否完成
*/
private boolean checkAllNodesCompleted(WorkflowDefinition definition) {
return true; // 简化处理
}
/**
* 记录节点执行情况
*/
private void recordNodeExecution(WorkflowInstance instance, WorkflowDefinition.Node node, String details) {
WorkflowInstance.NodeExecutionRecord record =
new WorkflowInstance.NodeExecutionRecord(
node.getId(),
node.getName(),
WorkflowInstance.NodeExecutionRecord.NodeExecutionStatus.PROCESSED
);
record.setDetails(details);
instance.addExecutionRecord(record);
}
/**
* 生成实例ID
*/
private String generateInstanceId(String definitionId) {
return definitionId + "_" + System.currentTimeMillis();
}
/**
* 生成任务ID
*/
private String generateTaskId(String instanceId) {
return instanceId + "_task_" + System.nanoTime();
}
/**
* 停止引擎
*/
public void shutdown() {
this.running = false;
}
}
WorkflowParser.java
package com.workflow.engine;
import com.workflow.model.WorkflowDefinition;
import com.workflow.model.WorkflowDefinition.Node;
import com.workflow.model.WorkflowDefinition.NodeType;
import com.workflow.model.WorkflowDefinition.Transition;
import java.util.ArrayList;
import java.util.List;
/**
* 工作流解析器(用于构建工作流定义)
*/
public class WorkflowParser {
/**
* 基于JSON构建工作流(简化版本)
*/
public WorkflowDefinition parseJson(String json) {
// 这里简化处理,实际项目中会使用Jackson或Gson解析
WorkflowDefinition definition = new WorkflowDefinition("demo_workflow", "审批流程", "1.0");
// 创建节点
Node startNode = new Node("start", "开始", NodeType.START);
Node approvalNode = new Node("approval", "审批", NodeType.TASK);
Node reviewNode = new Node("review", "复核", NodeType.TASK);
Node decisionNode = new Node("decision", "决策", NodeType.DECISION);
Node endNode = new Node("end", "结束", NodeType.END);
// 设置审批节点属性
approvalNode.setProperty("taskData", new java.util.HashMap<>() {{
put("taskType", "approval");
put("assignee", "经理");
}});
reviewNode.setProperty("taskData", new java.util.HashMap<>() {{
put("taskType", "review");
put("assignee", "财务");
}});
// 添加节点
definition.addNode(startNode);
definition.addNode(approvalNode);
definition.addNode(reviewNode);
definition.addNode(decisionNode);
definition.addNode(endNode);
// 创建连接
definition.addTransition(new Transition("start", "approval"));
definition.addTransition(new Transition("approval", "review"));
definition.addTransition(new Transition("review", "decision"));
definition.addTransition(new Transition("decision", "end", "approve"));
definition.addTransition(new Transition("decision", "approval", "reject"));
// 设置开始和结束节点
definition.setStartNodeId("start");
definition.setEndNodeIds(new ArrayList<>(List.of("end")));
return definition;
}
/**
* DSL方式创建工作流
*/
public WorkflowDefinition buildDsl() {
WorkflowDefinition definition = new WorkflowDefinition("dsl_workflow", "DSL流程", "1.0");
// 简单链式工作流
addNodesForDsl(definition);
return definition;
}
private void addNodesForDsl(WorkflowDefinition definition) {
// 创建简单流程
Node startNode = new Node("start", "开始", NodeType.START);
Node processNode1 = new Node("process1", "处理任务1", NodeType.TASK);
Node processNode2 = new Node("process2", "处理任务2", NodeType.TASK);
Node endNode = new Node("end", "结束", NodeType.END);
processNode1.setProperty("taskData", new java.util.HashMap<>() {{
put("description", "第一步处理");
}});
processNode2.setProperty("taskData", new java.util.HashMap<>() {{
put("description", "第二步处理");
}});
definition.addNode(startNode);
definition.addNode(processNode1);
definition.addNode(processNode2);
definition.addNode(endNode);
definition.addTransition(new Transition("start", "process1"));
definition.addTransition(new Transition("process1", "process2"));
definition.addTransition(new Transition("process2", "end"));
definition.setStartNodeId("start");
definition.setEndNodeIds(new ArrayList<>(List.of("end")));
}
}
TaskExecutor.java
package com.workflow.engine;
import com.workflow.handler.TaskHandler;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
/**
* 任务执行器
*/
public class TaskExecutor {
private final Map<String, TaskHandler> handlers = new ConcurrentHashMap<>();
/**
* 注册任务处理器
*/
public void registerHandler(String nodeId, TaskHandler handler) {
handlers.put(nodeId, handler);
System.out.println("注册处理器: " + nodeId + " -> " + handler.getClass().getSimpleName());
}
/**
* 获取任务处理器
*/
public TaskHandler getHandler(String nodeId) {
return handlers.get(nodeId);
}
/**
* 移除处理器
*/
public void removeHandler(String nodeId) {
handlers.remove(nodeId);
}
/**
* 清空所有处理器
*/
public void clearHandlers() {
handlers.clear();
}
}
处理器实现
TaskHandler.java
package com.workflow.handler;
import com.workflow.model.TaskInstance;
import com.workflow.model.WorkflowInstance;
/**
* 任务处理器接口
*/
public interface TaskHandler {
/**
* 执行任务
* @param task 任务实例
* @param workflowInstance 工作流实例
*/
void execute(TaskInstance task, WorkflowInstance workflowInstance);
}
DefaultTaskHandler.java
package com.workflow.handler;
import com.workflow.model.TaskInstance;
import com.workflow.model.WorkflowInstance;
/**
* 默认任务处理器
*/
public class DefaultTaskHandler implements TaskHandler {
@Override
public void execute(TaskInstance task, WorkflowInstance workflowInstance) {
System.out.println(" 执行默认任务处理器");
System.out.println(" 任务信息: " + task.getNodeName());
// 模拟任务执行时间
try {
Thread.sleep(500);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
// 设置任务结果
task.getTaskData().put("result", "success");
}
}
服务层
WorkflowService.java
package com.workflow.service;
import com.workflow.engine.TaskExecutor;
import com.workflow.engine.WorkflowEngine;
import com.workflow.engine.WorkflowParser;
import com.workflow.handler.DefaultTaskHandler;
import com.workflow.handler.TaskHandler;
import com.workflow.model.WorkflowDefinition;
import com.workflow.model.WorkflowInstance;
import java.util.Map;
import java.util.UUID;
import java.util.concurrent.ConcurrentHashMap;
/**
* 工作流服务
*/
public class WorkflowService {
private final Map<String, WorkflowDefinition> definitions = new ConcurrentHashMap<>();
private final Map<String, WorkflowInstance> instances = new ConcurrentHashMap<>();
private final WorkflowEngine engine;
private final TaskExecutor taskExecutor;
private final WorkflowParser parser;
public WorkflowService() {
this.parser = new WorkflowParser();
this.taskExecutor = new TaskExecutor();
this.engine = new WorkflowEngine(parser, taskExecutor);
}
/**
* 部署工作流定义
*/
public void deployWorkflow(WorkflowDefinition definition, Map<String, TaskHandler> handlers) {
definitions.put(definition.getId(), definition);
// 注册任务处理器
if (handlers != null) {
handlers.forEach(taskExecutor::registerHandler);
}
System.out.println("✓ 工作流已部署: " + definition.getName() + " (ID: " + definition.getId() + ")");
}
/**
* 启动工作流实例
*/
public WorkflowInstance startWorkflow(String definitionId, Map<String, Object> variables) {
WorkflowDefinition definition = definitions.get(definitionId);
if (definition == null) {
throw new IllegalArgumentException("未找到工作流定义: " + definitionId);
}
// 生成实例ID
String instanceId = UUID.randomUUID().toString();
WorkflowInstance instance = new WorkflowInstance(instanceId, definitionId);
// 启动流程
instance = engine.startWorkflow(definition, variables);
instances.put(instance.getId(), instance);
return instance;
}
/**
* 注册任务处理器
*/
public void registerTaskHandler(String nodeId, TaskHandler handler) {
taskExecutor.registerHandler(nodeId, handler);
}
/**
* 获取工作流实例
*/
public WorkflowInstance getWorkflowInstance(String instanceId) {
return instances.get(instanceId);
}
/**
* 获取所有工作流实例
*/
public Map<String, WorkflowInstance> getAllInstances() {
return instances;
}
/**
* 获取所有定义
*/
public Map<String, WorkflowDefinition> getAllDefinitions() {
return definitions;
}
}
测试演示类
WorkflowDemo.java
package com.workflow.demo;
import com.workflow.engine.WorkflowParser;
import com.workflow.handler.TaskHandler;
import com.workflow.model.*;
import com.workflow.service.WorkflowService;
import java.util.HashMap;
import java.util.Map;
/**
* 工作流演示
*/
public class WorkflowDemo {
public static void main(String[] args) {
System.out.println("=== 工作流引擎演示 ===\n");
// 创建服务
WorkflowService service = new WorkflowService();
WorkflowParser parser = new WorkflowParser();
// 演示1: 审批流程
demoApprovalWorkflow(service, parser);
// 演示2: 简单线性流程
demoLinearWorkflow(service, parser);
// 演示3: 带自定义处理器的流程
demoCustomHandlerWorkflow(service, parser);
}
/**
* 演示审批流程
*/
private static void demoApprovalWorkflow(WorkflowService service, WorkflowParser parser) {
System.out.println("\n<<<<<< 演示1: 审批流程 >>>>>>");
// 构建审批流程
WorkflowDefinition approvalDefinition = parser.parseJson(null);
// 注册自定义处理器
Map<String, TaskHandler> handlers = new HashMap<>();
handlers.put("approval", new TaskHandler() {
@Override
public void execute(TaskInstance task, WorkflowInstance workflowInstance) {
System.out.println(" [审批节点] 审批中...");
// 模拟审批逻辑
Map<String, Object> data = workflowInstance.getVariables();
Object amount = data.get("amount");
boolean approved = amount != null && (Double) amount < 10000;
workflowInstance.updateVariable("decision_result", approved ? "approve" : "reject");
System.out.println(" [审批节点] 审批结果: " + (approved ? "通过" : "拒绝"));
}
});
handlers.put("review", new TaskHandler() {
@Override
public void execute(TaskInstance task, WorkflowInstance workflowInstance) {
System.out.println(" [复核节点] 复核中...");
// 模拟复核逻辑
workflowInstance.updateVariable("review_status", "completed");
System.out.println(" [复核节点] 复核完成");
}
});
// 部署工作流
service.deployWorkflow(approvalDefinition, handlers);
// 准备流程变量
Map<String, Object> variables = new HashMap<>();
variables.put("amount", 5000.0);
variables.put("applicant", "张三");
variables.put("purpose", "市场推广费用");
// 启动流程
WorkflowInstance instance = service.startWorkflow("demo_workflow", variables);
// 输出执行结果
printExecutionResult(instance);
}
/**
* 演示简单线性流程
*/
private static void demoLinearWorkflow(WorkflowService service, WorkflowParser parser) {
System.out.println("\n<<<<<< 演示2: 简单线性流程 >>>>>>");
// 构建DSL流程
WorkflowDefinition dslDefinition = parser.buildDsl();
// 部署
service.deployWorkflow(dslDefinition, null);
// 启动流程
WorkflowInstance instance = service.startWorkflow("dsl_workflow", null);
// 输出执行结果
printExecutionResult(instance);
}
/**
* 演示带自定义处理器的流程
*/
private static void demoCustomHandlerWorkflow(WorkflowService service, WorkflowParser parser) {
System.out.println("\n<<<<<< 演示3: 自定义处理器流程 >>>>>>");
// 创建自定义工作流
WorkflowDefinition customDefinition = new WorkflowDefinition("custom_workflow", "自定义流程", "1.0");
// 创建节点
WorkflowDefinition.Node startNode = new WorkflowDefinition.Node("start", "开始", WorkflowDefinition.NodeType.START);
WorkflowDefinition.Node dataProcessNode = new WorkflowDefinition.Node("dataProcess", "数据处理", WorkflowDefinition.NodeType.TASK);
WorkflowDefinition.Node reportNode = new WorkflowDefinition.Node("report", "生成报告", WorkflowDefinition.NodeType.TASK);
WorkflowDefinition.Node endNode = new WorkflowDefinition.Node("end", "结束", WorkflowDefinition.NodeType.END);
// 设置属性
dataProcessNode.setProperty("taskData", new HashMap<>() {{
put("description", "处理输入数据");
put("timeout", 1000);
}});
reportNode.setProperty("taskData", new HashMap<>() {{
put("description", "生成执行报告");
put("format", "PDF");
}});
// 添加节点
customDefinition.addNode(startNode);
customDefinition.addNode(dataProcessNode);
customDefinition.addNode(reportNode);
customDefinition.addNode(endNode);
// 添加连接
customDefinition.addTransition(new WorkflowDefinition.Transition("start", "dataProcess"));
customDefinition.addTransition(new WorkflowDefinition.Transition("dataProcess", "report"));
customDefinition.addTransition(new WorkflowDefinition.Transition("report", "end"));
// 设置开始和结束
customDefinition.setStartNodeId("start");
customDefinition.setEndNodeIds(new java.util.ArrayList<>(java.util.List.of("end")));
// 注册自定义处理器
Map<String, TaskHandler> handlers = new HashMap<>();
handlers.put("dataProcess", (task, instance) -> {
System.out.println(" [数据处理] 开始处理数据...");
try {
Thread.sleep(300);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
// 模拟数据处理
task.getTaskData().put("processedData", "处理完成的数据");
instance.updateVariable("dataStatus", "processed");
System.out.println(" [数据处理] 数据处理完成");
});
handlers.put("report", (task, instance) -> {
System.out.println(" [报告生成] 正在生成报告...");
try {
Thread.sleep(200);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
task.getTaskData().put("reportUrl", "/reports/report_" + instance.getId() + ".pdf");
System.out.println(" [报告生成] 报告已生成,可以下载");
});
// 部署并启动
service.deployWorkflow(customDefinition, handlers);
Map<String, Object> variables = new HashMap<>();
variables.put("inputData", "原始数据");
variables.put("priority", "high");
WorkflowInstance instance = service.startWorkflow("custom_workflow", variables);
// 输出执行结果
printExecutionResult(instance);
}
/**
* 打印执行结果
*/
private static void printExecutionResult(WorkflowInstance instance) {
System.out.println("\n----- 执行结果 -----");
System.out.println("实例ID: " + instance.getId());
System.out.println("状态: " + instance.getStatus());
System.out.println("开始时间: " + instance.getStartTime());
System.out.println("结束时间: " + instance.getEndTime());
System.out.println("总任务数: " + instance.getTotalTaskCount());
System.out.println("完成任务数: " + instance.getCompletedTaskCount());
System.out.println("\n变量信息:");
instance.getVariables().forEach((k, v) ->
System.out.println(" " + k + " = " + v)
);
System.out.println("\n执行轨迹:");
instance.getExecutionRecords().forEach((key, record) ->
System.out.println(" [" + record.getExecutionTime() + "] " +
record.getNodeName() + " (" + record.getStatus() + ")")
);
System.out.println("\n任务列表:");
instance.getTasks().forEach((id, task) ->
System.out.println(" Task: " + task.getNodeName() +
" - 状态: " + task.getStatus() +
" - 数据: " + task.getTaskData())
);
System.out.println("-------------------\n");
}
}
Maven依赖 (pom.xml)
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0
http://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>
<groupId>com.workflow</groupId>
<artifactId>workflow-demo</artifactId>
<version>1.0-SNAPSHOT</version>
<packaging>jar</packaging>
<properties>
<maven.compiler.source>1.8</maven.compiler.source>
<maven.compiler.target>1.8</maven.compiler.target>
<project.build