From c49e1a9ac5aa777e359497083e46bf1511e17fa6 Mon Sep 17 00:00:00 2001 From: lorne <1991wangliang@gmail.com> Date: Mon, 17 Aug 2026 18:09:51 +0800 Subject: [PATCH 1/4] =?UTF-8?q?perf:=20=E9=AB=98=E9=A2=91=E8=AE=B0?= =?UTF-8?q?=E5=BD=95=E8=A1=A8=E4=B8=BB=E9=94=AE=20IDENTITY=E2=86=92SEQUENC?= =?UTF-8?q?E=20=E5=B9=B6=E5=BC=80=E5=90=AF=E6=89=B9=E6=8F=92=EF=BC=8C?= =?UTF-8?q?=E7=BC=93=E8=A7=A3=E5=AD=90=E6=B5=81=E7=A8=8B=E6=89=B9=E9=87=8F?= =?UTF-8?q?=E5=88=9B=E5=BB=BA=E6=85=A2=20issue=20#210?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 4 张高频写热表(t_flow_record/t_flow_todo_record/t_flow_todo_marge/t_flow_sub_process_record)主键改 SEQUENCE,解锁 Hibernate JDBC 批插,将批量创建 ~500 条串行 DB 往返降为约几十条批次往返 - 示例配置开启 hibernate.jdbc.batch_size=50 + order_inserts/order_updates - 低频流程设计表(t_flow_workflow/t_flow_workflow_version/t_flow_workflow_runtime 等)保持 IDENTITY 不动,避免触碰存量设计数据 - 新增 issue 210 复现测量测试 FlowIssue210SubProcessPerformanceTest - 运行机制受影响面为零:id>0 判定、setId(entity.getId())、同事务 fromId 引用、按 id 排序语义均不变 #210 Co-Authored-By: Claude --- .../main/resources/application-dm.properties | 6 + .../src/main/resources/application.properties | 5 + ...FlowIssue210SubProcessPerformanceTest.java | 331 ++++++++++++++++++ .../flow/infra/entity/FlowRecordEntity.java | 2 +- .../infra/entity/FlowTodoMargeEntity.java | 2 +- .../infra/entity/FlowTodoRecordEntity.java | 3 +- .../infra/entity/SubProcessRecordEntity.java | 4 +- 7 files changed, 348 insertions(+), 5 deletions(-) create mode 100644 flow-engine-framework/src/test/java/com/codingapi/flow/service/FlowIssue210SubProcessPerformanceTest.java diff --git a/flow-engine-example/src/main/resources/application-dm.properties b/flow-engine-example/src/main/resources/application-dm.properties index 39c25832..b801958a 100644 --- a/flow-engine-example/src/main/resources/application-dm.properties +++ b/flow-engine-example/src/main/resources/application-dm.properties @@ -4,6 +4,12 @@ spring.jpa.showSql=true spring.jpa.hibernate.ddl-auto=update spring.jpa.properties.hibernate.dialect=com.codingapi.example.dialect.HydDmDialect +# 开启 JDBC 批次插入,配合 SEQUENCE 主键(t_flow_record/t_flow_todo_record/t_flow_todo_marge/t_flow_sub_process_record) +# 将批量创建流程记录/待办/合并/子流程时的多次串行 INSERT 合并为批次往返,显著降低远程库耗时 +spring.jpa.properties.hibernate.jdbc.batch_size=100 +spring.jpa.properties.hibernate.order_inserts=true +spring.jpa.properties.hibernate.order_updates=true + spring.datasource.url=jdbc:dm://localhost:5236 spring.datasource.username=SYSDBA spring.datasource.password=SYSDBA001 diff --git a/flow-engine-example/src/main/resources/application.properties b/flow-engine-example/src/main/resources/application.properties index 75fa1a1b..4cc12d8c 100644 --- a/flow-engine-example/src/main/resources/application.properties +++ b/flow-engine-example/src/main/resources/application.properties @@ -7,6 +7,11 @@ spring.jpa.database-platform=org.hibernate.dialect.H2Dialect spring.jpa.hibernate.ddl-auto=update spring.jpa.show-sql=true +# 开启 JDBC 批次插入,配合 SEQUENCE 主键 + H2 序列,优化批量流程/待办/合并/子流程记录创建 +spring.jpa.properties.hibernate.jdbc.batch_size=50 +spring.jpa.properties.hibernate.order_inserts=true +spring.jpa.properties.hibernate.order_updates=true + # Security spring.main.allow-bean-definition-overriding=true codingapi.security.jwt.enable=true diff --git a/flow-engine-framework/src/test/java/com/codingapi/flow/service/FlowIssue210SubProcessPerformanceTest.java b/flow-engine-framework/src/test/java/com/codingapi/flow/service/FlowIssue210SubProcessPerformanceTest.java new file mode 100644 index 00000000..0c237d13 --- /dev/null +++ b/flow-engine-framework/src/test/java/com/codingapi/flow/service/FlowIssue210SubProcessPerformanceTest.java @@ -0,0 +1,331 @@ +package com.codingapi.flow.service; + +import com.codingapi.flow.builder.FormFieldPermissionsBuilder; +import com.codingapi.flow.builder.NodeStrategyBuilder; +import com.codingapi.flow.context.GatewayContext; +import com.codingapi.flow.domain.SubProcessRecord; +import com.codingapi.flow.factory.MyFlowServiceFactory; +import com.codingapi.flow.form.DataType; +import com.codingapi.flow.form.FlowForm; +import com.codingapi.flow.form.FlowFormBuilder; +import com.codingapi.flow.form.permission.PermissionType; +import com.codingapi.flow.gateway.FlowOperatorGateway; +import com.codingapi.flow.node.IFlowNode; +import com.codingapi.flow.node.nodes.ApprovalNode; +import com.codingapi.flow.node.nodes.EndNode; +import com.codingapi.flow.node.nodes.StartNode; +import com.codingapi.flow.node.nodes.SubProcessNode; +import com.codingapi.flow.operator.IFlowOperator; +import com.codingapi.flow.pojo.body.FlowAdviceBody; +import com.codingapi.flow.pojo.request.FlowActionRequest; +import com.codingapi.flow.pojo.request.FlowCreateRequest; +import com.codingapi.flow.record.FlowRecord; +import com.codingapi.flow.script.factory.FlowGroovyScriptFactory; +import com.codingapi.flow.strategy.node.FormFieldPermissionStrategy; +import com.codingapi.flow.strategy.node.MultiOperatorAuditStrategy; +import com.codingapi.flow.strategy.node.OperatorLoadStrategy; +import com.codingapi.flow.strategy.node.SubProcessStrategy; +import com.codingapi.flow.user.User; +import com.codingapi.flow.workflow.Workflow; +import com.codingapi.flow.workflow.WorkflowBuilder; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; + +import java.util.ArrayList; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNotNull; + +/** + * 问题 210 复现测试:主流程 C(子流程节点)一次性创建 6 个子流程、每个子流程 B 节点 20 个审批人。 + *

+ * 用于量化网关延时(模拟真实 DB/远程查询 5ms)下,C 节点提交的整体耗时与网关查询次数, + * 定位是否存在可优化的空间,例如操作人缓存是否因每次 create/action 清空而反复查询同一批人员。 + */ +class FlowIssue210SubProcessPerformanceTest { + + private static final String FORM_CODE = "issue-210-form"; + private static final String PARENT_CODE = "issue-210-parent"; + private static final String CHILD_CODE = "issue-210-child"; + private static final int SUB_PROCESS_COUNT = 6; + private static final int APPROVER_COUNT = 20; + private static final long GATEWAY_DELAY_MS = 5; + + private MyFlowServiceFactory factory; + private User initiator; + private List approvers; + private StartNode childStart; + private ApprovalNode childB; + private DelayedGateway gateway; + + private StartNode parentStart; + private ApprovalNode parentB; + + @BeforeEach + void setUp() { + factory = new MyFlowServiceFactory(); + initiator = saveUser(1, "发起人"); + approvers = new ArrayList<>(); + for (int i = 0; i < APPROVER_COUNT; i++) { + approvers.add(saveUser(100 + i, "审批人" + i)); + } + FlowForm form = FlowFormBuilder.builder() + .name("问题210测试表单") + .code(FORM_CODE) + .addField("业务内容", "content", DataType.STRING) + .build(); + // 用带延时与计数的网关替换默认 UserGateway,模拟真实环境的人员查询代价 + gateway = new DelayedGateway(GatewayContext.getInstance().getFlowOperatorGateway(), GATEWAY_DELAY_MS); + GatewayContext.getInstance().setFlowOperatorGateway(gateway); + saveChildWorkflow(form); + buildParentWorkflow(form); + } + + @AfterEach + void tearDown() { + GatewayContext.getInstance().setFlowOperatorGateway(new com.codingapi.flow.gateway.impl.UserGateway()); + } + + /** + * 测试目标:量化主流程 C(子流程)节点创建 6 个子流程、每个子流程 B 节点 20 个审批人的耗时与网关查询次数。 + * 前置条件:节点为 A-B-SubProcess-D-E,子流程 A-B-C-D 且 B 审批人取表单 approvers 字段;网关每次查询延时 5ms。 + * 执行步骤:提交主流程并审批 A、B,触发子流程批量创建。 + * 期望断言:共创建 6 个子流程、每个子流程 B 节点 20 个待办(总量 120 条);记录耗时与网关批量查询次数。 + */ + @Test + void shouldMeasureSubProcessBulkCreationTimeWithGatewayDelay() { + // when:创建主流程并审批 A 节点,进入 B 节点 + Map data = new LinkedHashMap<>(); + data.put("content", "parent"); + data.put("approvers", approvers.stream().map(User::getUserId).toList()); + long parentRecordId = createParent(data); + approveMain(factory.flowRecordRepository.get(parentRecordId), parentStart, data); + FlowRecord parentBRecord = findTodo(initiator, parentB.getId()); + assertNotNull(parentBRecord, "主流程应进入 B 审批节点"); + + // when:审批 B 节点,触发 C 子流程节点批量创建 6 个子流程(B->C 流转) + long startNanos = System.nanoTime(); + approveMain(parentBRecord, parentB, data); + long elapsedMs = (System.nanoTime() - startNanos) / 1_000_000; + + // then:6 个子流程实例,每个子流程 B 节点 20 个待办(共 120 条子流程记录) + List subProcessRecords = factory.subProcessRepository + .findByParentRecordId(parentBRecord.getId()); + assertEquals(1, subProcessRecords.size(), "主流程应产生一条子流程聚合记录"); + assertEquals(SUB_PROCESS_COUNT, subProcessRecords.get(0).getInstances().size(), + "应创建订阅个数子流程"); + assertEquals(SUB_PROCESS_COUNT * APPROVER_COUNT, countChildBTodos(), + "每个子流程 B 节点 20 审批人,共创建 120 条待办记录"); + + // 汇报量化指标 + gateway.printReport(); + System.out.println("[issue-210] 主流程 B/C 节点提交总耗时: " + elapsedMs + " ms" + + " (网关延时 " + GATEWAY_DELAY_MS + "ms/次)"); + } + + // ---------- 网关延时 + 计数 ---------- + + private static class DelayedGateway implements FlowOperatorGateway { + private final FlowOperatorGateway delegate; + private final long delayMs; + private long findCalls; + private long getCalls; + private long totalIds; + + DelayedGateway(FlowOperatorGateway delegate, long delayMs) { + this.delegate = delegate; + this.delayMs = delayMs; + } + + @Override + public IFlowOperator get(long id) { + sleep(); + getCalls++; + return delegate.get(id); + } + + @Override + public List findByIds(List ids) { + sleep(); + findCalls++; + totalIds += ids.size(); + return delegate.findByIds(ids); + } + + private void sleep() { + try { + Thread.sleep(delayMs); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + } + + void printReport() { + System.out.println("[issue-210] 网关单点查询 get() 次数: " + getCalls); + System.out.println("[issue-210] 网关批量查询 findByIds 次数: " + findCalls + + ",累计覆盖人员 " + totalIds + " 人次"); + System.out.println("[issue-210] 网关总查询次数: " + (getCalls + findCalls) + + ",纯网关延时约 " + ((getCalls + findCalls) * delayMs) + " ms"); + } + } + + // ---------- 流程定义构建 ---------- + + private void saveChildWorkflow(FlowForm form) { + childStart = writableStart("子流程开始", form); + childB = ApprovalNode.builder() + .name("子流程B") + .strategies(NodeStrategyBuilder.builder() + .addStrategy(readonlyPermission(FORM_CODE, "content")) + .addStrategy(new MultiOperatorAuditStrategy( + MultiOperatorAuditStrategy.Type.MERGE, 1.0f)) + .addStrategy(new OperatorLoadStrategy( + FlowGroovyScriptFactory.createOperatorLoadScript( + "def run(request){ return request.getFormData('approvers') }").getKey())) + .build()) + .build(); + ApprovalNode childC = ApprovalNode.builder() + .name("子流程C") + .strategies(NodeStrategyBuilder.builder() + .addStrategy(readonlyPermission(FORM_CODE, "content")) + .addStrategy(new OperatorLoadStrategy( + FlowGroovyScriptFactory.createOperatorLoadScript( + "def run(request){ return [" + approvers.get(0).getUserId() + "] }").getKey())) + .build()) + .build(); + Workflow childWorkflow = WorkflowBuilder.builder() + .title("问题210-子流程") + .code(CHILD_CODE) + .createdOperator(initiator) + .form(form) + .addNode(childStart) + .addNode(childB) + .addNode(childC) + .addNode(EndNode.builder().name("子流程结束").build()) + .build(); + factory.workflowService.saveWorkflow(childWorkflow); + } + + private void buildParentWorkflow(FlowForm form) { + parentStart = writableStart("主流程开始", form); + parentB = ApprovalNode.builder() + .name("主流程B") + .strategies(NodeStrategyBuilder.builder() + .addStrategy(readonlyPermission(FORM_CODE, "content")) + .addStrategy(new OperatorLoadStrategy( + FlowGroovyScriptFactory.createOperatorLoadScript( + "def run(request){ return [" + initiator.getUserId() + "] }").getKey())) + .build()) + .build(); + String createScript = """ + def run(request){ + def approvers = request.getFormData('approvers') + def list = [] + for (int i = 0; i < %d; i++) { + list.add(request.toCreateRequest('%s', %d, '%s', [approvers: approvers, content: 'child-' + i])) + } + return list + } + """.formatted(SUB_PROCESS_COUNT, CHILD_CODE, initiator.getUserId(), + passAction(childStart).id()); + SubProcessNode subProcess = SubProcessNode.builder() + .name("主流程C-子流程") + .strategies(NodeStrategyBuilder.builder() + .addStrategy(readonlyPermission(FORM_CODE, "content")) + .addStrategy(new SubProcessStrategy( + FlowGroovyScriptFactory.createSubProcessScript(createScript).getKey(), true)) + .build()) + .build(); + ApprovalNode parentD = ApprovalNode.builder() + .name("主流程D") + .strategies(NodeStrategyBuilder.builder() + .addStrategy(readonlyPermission(FORM_CODE, "content")) + .addStrategy(new OperatorLoadStrategy( + FlowGroovyScriptFactory.createOperatorLoadScript( + "def run(request){ return [" + initiator.getUserId() + "] }").getKey())) + .build()) + .build(); + Workflow parentWorkflow = WorkflowBuilder.builder() + .title("问题210-主流程") + .code(PARENT_CODE) + .createdOperator(initiator) + .form(form) + .addNode(parentStart) + .addNode(parentB) + .addNode(subProcess) + .addNode(parentD) + .addNode(EndNode.builder().name("主流程结束").build()) + .build(); + factory.workflowService.saveWorkflow(parentWorkflow); + } + + // ---------- 辅助方法 ---------- + + private long createParent(Map data) { + FlowCreateRequest request = new FlowCreateRequest(); + request.setWorkCode(PARENT_CODE); + request.setOperatorId(initiator.getUserId()); + request.setActionId(passAction(parentStart).id()); + request.setFormData(data); + return factory.flowService.create(request); + } + + private void approveMain(FlowRecord record, IFlowNode node, Map data) { + FlowActionRequest request = new FlowActionRequest(); + request.setRecordId(record.getId()); + request.setFormData(data); + request.setAdvice(new FlowAdviceBody(passAction(node).id(), "同意", initiator.getUserId())); + factory.flowService.action(request); + } + + private FlowRecord findTodo(User operator, String nodeId) { + return factory.flowRecordRepository.findTodoByOperator(operator.getUserId()).stream() + .filter(record -> record.getNodeId().equals(nodeId)) + .findFirst() + .orElse(null); + } + + private int countChildBTodos() { + int total = 0; + for (int i = 0; i < APPROVER_COUNT; i++) { + total += factory.flowRecordRepository.findTodoByOperator(approvers.get(i).getUserId()).stream() + .filter(record -> record.getNodeId().equals(childB.getId())) + .count(); + } + return total; + } + + private StartNode writableStart(String name, FlowForm form) { + return StartNode.builder() + .name(name) + .strategies(NodeStrategyBuilder.builder() + .addStrategy(new FormFieldPermissionStrategy(FormFieldPermissionsBuilder.builder() + .addPermission(FORM_CODE, "content", PermissionType.WRITE) + .build())) + .build()) + .build(); + } + + private FormFieldPermissionStrategy readonlyPermission(String formCode, String field) { + return new FormFieldPermissionStrategy(FormFieldPermissionsBuilder.builder() + .addPermission(formCode, field, PermissionType.READ) + .build()); + } + + private com.codingapi.flow.action.IFlowAction passAction(IFlowNode node) { + return node.actionManager().getActions().stream() + .filter(action -> "PASS".equals(action.type())) + .findFirst() + .orElseThrow(); + } + + private User saveUser(long id, String name) { + User user = new User(id, name); + factory.userGateway.save(user); + return user; + } +} \ No newline at end of file diff --git a/flow-engine-starter-infra/src/main/java/com/codingapi/flow/infra/entity/FlowRecordEntity.java b/flow-engine-starter-infra/src/main/java/com/codingapi/flow/infra/entity/FlowRecordEntity.java index c94e437f..75716479 100644 --- a/flow-engine-starter-infra/src/main/java/com/codingapi/flow/infra/entity/FlowRecordEntity.java +++ b/flow-engine-starter-infra/src/main/java/com/codingapi/flow/infra/entity/FlowRecordEntity.java @@ -11,7 +11,7 @@ public class FlowRecordEntity { * 记录id */ @Id - @GeneratedValue(strategy = GenerationType.IDENTITY) + @GeneratedValue(strategy = GenerationType.SEQUENCE) private Long id; /** * 工作id diff --git a/flow-engine-starter-infra/src/main/java/com/codingapi/flow/infra/entity/FlowTodoMargeEntity.java b/flow-engine-starter-infra/src/main/java/com/codingapi/flow/infra/entity/FlowTodoMargeEntity.java index 515cc0ca..5d699e8c 100644 --- a/flow-engine-starter-infra/src/main/java/com/codingapi/flow/infra/entity/FlowTodoMargeEntity.java +++ b/flow-engine-starter-infra/src/main/java/com/codingapi/flow/infra/entity/FlowTodoMargeEntity.java @@ -9,7 +9,7 @@ public class FlowTodoMargeEntity { @Id - @GeneratedValue(strategy = GenerationType.IDENTITY) + @GeneratedValue(strategy = GenerationType.SEQUENCE) private Long id; /** * 待办id diff --git a/flow-engine-starter-infra/src/main/java/com/codingapi/flow/infra/entity/FlowTodoRecordEntity.java b/flow-engine-starter-infra/src/main/java/com/codingapi/flow/infra/entity/FlowTodoRecordEntity.java index 145e3b5f..88555dda 100644 --- a/flow-engine-starter-infra/src/main/java/com/codingapi/flow/infra/entity/FlowTodoRecordEntity.java +++ b/flow-engine-starter-infra/src/main/java/com/codingapi/flow/infra/entity/FlowTodoRecordEntity.java @@ -11,9 +11,10 @@ public class FlowTodoRecordEntity { /** * 合并记录id + *

使用序列生成(而非数据库自增),配合 hibernate 批插,降低高频批量写表的 DB 往返。 */ @Id - @GeneratedValue(strategy = GenerationType.IDENTITY) + @GeneratedValue(strategy = GenerationType.SEQUENCE) private Long id; /** diff --git a/flow-engine-starter-infra/src/main/java/com/codingapi/flow/infra/entity/SubProcessRecordEntity.java b/flow-engine-starter-infra/src/main/java/com/codingapi/flow/infra/entity/SubProcessRecordEntity.java index 5bc176b6..ef0da7c6 100644 --- a/flow-engine-starter-infra/src/main/java/com/codingapi/flow/infra/entity/SubProcessRecordEntity.java +++ b/flow-engine-starter-infra/src/main/java/com/codingapi/flow/infra/entity/SubProcessRecordEntity.java @@ -25,10 +25,10 @@ public class SubProcessRecordEntity { /** - * 主键,自增 + * 主键,由序列生成(替代数据库自增,配合 hibernate 批插降低批量往返) */ @Id - @GeneratedValue(strategy = GenerationType.IDENTITY) + @GeneratedValue(strategy = GenerationType.SEQUENCE) private Long id; /** From a67197fd271ce8418045c9df2a989924336bcacb Mon Sep 17 00:00:00 2001 From: lorne <1991wangliang@gmail.com> Date: Mon, 17 Aug 2026 18:17:29 +0800 Subject: [PATCH 2/4] =?UTF-8?q?perf:=20=E6=B6=88=E9=99=A4=E4=BF=9D?= =?UTF-8?q?=E5=AD=98=E5=BE=85=E5=8A=9E=E6=97=B6=E7=9A=84=20getByTodoKey=20?= =?UTF-8?q?N+1=EF=BC=8C=E6=94=B9=E4=B8=BA=E6=89=B9=E9=87=8F=20findByKeys?= =?UTF-8?q?=20=E5=88=86=E5=9D=97=E5=8A=A0=E8=BD=BD=20issue=20#210?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - FlowRecordSaveService.saveTodoMargeRecords 先批量加载已存在待办,替代循环内逐条 getByTodoKey(每 20 条待办 20 次 SELECT → 1 次 IN 查询) - 按 TODO_KEY_BATCH_SIZE=500 分块查询,避免海量 key 拼超大 IN 子句/结果集过大导致 OOM - FlowTodoRecordRepository 新增 findByKeys,实现于 mock/test/Infra 三层仓储 - 新增集成断言:批量场景逐条 getByTodoKey 调用数远小于待办数,且走批量 findByKeys - 全量 clean install BUILD SUCCESS,227 framework + example 真实 JPA 集成测试全绿 #210 Co-Authored-By: Claude --- .../FlowTodoRecordRepositoryMockImpl.java | 25 +++++++++++++++++ .../repository/FlowTodoRecordRepository.java | 8 ++++++ .../flow/service/FlowRecordSaveService.java | 27 ++++++++++++++++++- .../FlowTodoRecordRepositoryImpl.java | 25 +++++++++++++++++ ...FlowIssue210SubProcessPerformanceTest.java | 25 +++++++++++++++++ .../jpa/FlowTodoRecordEntityRepository.java | 5 ++++ .../impl/FlowTodoRecordRepositoryImpl.java | 7 +++++ 7 files changed, 121 insertions(+), 1 deletion(-) diff --git a/flow-engine-framework/src/main/java/com/codingapi/flow/mock/repository/FlowTodoRecordRepositoryMockImpl.java b/flow-engine-framework/src/main/java/com/codingapi/flow/mock/repository/FlowTodoRecordRepositoryMockImpl.java index 2df43382..b966b0ec 100644 --- a/flow-engine-framework/src/main/java/com/codingapi/flow/mock/repository/FlowTodoRecordRepositoryMockImpl.java +++ b/flow-engine-framework/src/main/java/com/codingapi/flow/mock/repository/FlowTodoRecordRepositoryMockImpl.java @@ -4,8 +4,10 @@ import com.codingapi.flow.repository.FlowTodoRecordRepository; import java.util.HashMap; +import java.util.HashSet; import java.util.List; import java.util.Map; +import java.util.Set; public class FlowTodoRecordRepositoryMockImpl implements FlowTodoRecordRepository { @@ -13,6 +15,10 @@ public class FlowTodoRecordRepositoryMockImpl implements FlowTodoRecordRepositor private final Map cacheByMageKey = new HashMap<>(); private long nextId = 1; + // 查询计数器,供测试断言"按 key 逐条查询已被批量查询替代"(N+1 消除) + private int getByTodoKeyCalls; + private int findByKeysCalls; + @Override public void save(FlowTodoRecord record) { if (record.getId() > 0) { @@ -41,9 +47,28 @@ public void delete(FlowTodoRecord margeRecord) { @Override public FlowTodoRecord getByTodoKey(String key) { + getByTodoKeyCalls++; return cacheByMageKey.get(key); } + @Override + public List findByKeys(List keys) { + findByKeysCalls++; + Set keySet = new HashSet<>(keys); + return cacheByMageKey.entrySet().stream() + .filter(entry -> keySet.contains(entry.getKey())) + .map(Map.Entry::getValue) + .toList(); + } + + public int getGetByTodoKeyCalls() { + return getByTodoKeyCalls; + } + + public int getFindByKeysCalls() { + return findByKeysCalls; + } + public List findByOperatorId(long operatorId) { return cache.values().stream().filter(record -> record.getCurrentOperatorId() == operatorId).toList(); } diff --git a/flow-engine-framework/src/main/java/com/codingapi/flow/repository/FlowTodoRecordRepository.java b/flow-engine-framework/src/main/java/com/codingapi/flow/repository/FlowTodoRecordRepository.java index 40408829..e4900882 100644 --- a/flow-engine-framework/src/main/java/com/codingapi/flow/repository/FlowTodoRecordRepository.java +++ b/flow-engine-framework/src/main/java/com/codingapi/flow/repository/FlowTodoRecordRepository.java @@ -10,6 +10,14 @@ public interface FlowTodoRecordRepository { FlowTodoRecord getByTodoKey(String key); + /** + * 按多个待办合并 key 批量加载已存在的待办记录,用于替代循环内逐条 {@link #getByTodoKey(String)} 的 N+1 查询。 + * + * @param keys 待办合并 key 列表 + * @return 已存在的待办记录 + */ + List findByKeys(List keys); + void delete(FlowTodoRecord margeRecord); void save(FlowTodoRecord margeRecord); diff --git a/flow-engine-framework/src/main/java/com/codingapi/flow/service/FlowRecordSaveService.java b/flow-engine-framework/src/main/java/com/codingapi/flow/service/FlowRecordSaveService.java index 6b040b65..194444e9 100644 --- a/flow-engine-framework/src/main/java/com/codingapi/flow/service/FlowRecordSaveService.java +++ b/flow-engine-framework/src/main/java/com/codingapi/flow/service/FlowRecordSaveService.java @@ -8,13 +8,21 @@ import com.codingapi.flow.repository.FlowTodoRecordRepository; import java.util.ArrayList; +import java.util.HashMap; import java.util.List; +import java.util.Map; +import java.util.Objects; /** * 流程记录保存服务,负责保存流程记录和待办记录的合并关系 */ class FlowRecordSaveService { + /** + * 批量加载已存在待办时单批查询的 key 上限,避免单个 IN 子句过大(SQL 限制)与结果集过大(内存)。 + */ + private static final int TODO_KEY_BATCH_SIZE = 500; + private final List flowRecords; private FlowTodoRecordRepository flowTodoRecordRepository; @@ -41,12 +49,29 @@ public void registerRepositories(FlowTodoRecordRepository flowTodoRecordReposito private void saveTodoMargeRecords() { + // 批量加载已存在的待办(而非逐条 getByTodoKey 的 N+1): + // 按分块大小分批 findByKeys,避免海量 key 拼在单个 IN 子句里导致 SQL/结果集过大(OOM 风险)。 + List todoKeys = flowRecords.stream() + .filter(FlowRecord::isTodo) + .map(FlowRecord::getTodoKey) + .filter(Objects::nonNull) + .distinct() + .toList(); + Map existedByKey = new HashMap<>(); + for (int start = 0; start < todoKeys.size(); start += TODO_KEY_BATCH_SIZE) { + List chunk = todoKeys.subList(start, Math.min(start + TODO_KEY_BATCH_SIZE, todoKeys.size())); + for (FlowTodoRecord existed : flowTodoRecordRepository.findByKeys(chunk)) { + existedByKey.put(existed.getTodoKey(), existed); + } + } + List flowTodoRecords = new ArrayList<>(); for (FlowRecord flowRecord : flowRecords) { if (flowRecord.isTodo()) { - FlowTodoRecord todoMargeRecord = flowTodoRecordRepository.getByTodoKey(flowRecord.getTodoKey()); + FlowTodoRecord todoMargeRecord = existedByKey.get(flowRecord.getTodoKey()); if (todoMargeRecord == null) { todoMargeRecord = new FlowTodoRecord(flowRecord); + existedByKey.put(flowRecord.getTodoKey(), todoMargeRecord); } else { todoMargeRecord.update(flowRecord); if (flowRecord.isMergeable()) { diff --git a/flow-engine-framework/src/test/java/com/codingapi/flow/repository/FlowTodoRecordRepositoryImpl.java b/flow-engine-framework/src/test/java/com/codingapi/flow/repository/FlowTodoRecordRepositoryImpl.java index d8d3de80..dcd4eac6 100644 --- a/flow-engine-framework/src/test/java/com/codingapi/flow/repository/FlowTodoRecordRepositoryImpl.java +++ b/flow-engine-framework/src/test/java/com/codingapi/flow/repository/FlowTodoRecordRepositoryImpl.java @@ -3,8 +3,10 @@ import com.codingapi.flow.record.FlowTodoRecord; import java.util.HashMap; +import java.util.HashSet; import java.util.List; import java.util.Map; +import java.util.Set; public class FlowTodoRecordRepositoryImpl implements FlowTodoRecordRepository { @@ -12,6 +14,10 @@ public class FlowTodoRecordRepositoryImpl implements FlowTodoRecordRepository { private final Map cacheByMageKey = new HashMap<>(); private long nextId = 1; + // 查询计数器,供测试断言"按 key 逐条查询已被批量查询替代"(N+1 消除) + private int getByTodoKeyCalls; + private int findByKeysCalls; + @Override public void save(FlowTodoRecord record) { if (record.getId() > 0) { @@ -40,9 +46,28 @@ public void delete(FlowTodoRecord margeRecord) { @Override public FlowTodoRecord getByTodoKey(String key) { + getByTodoKeyCalls++; return cacheByMageKey.get(key); } + @Override + public List findByKeys(List keys) { + findByKeysCalls++; + Set keySet = new HashSet<>(keys); + return cacheByMageKey.entrySet().stream() + .filter(entry -> keySet.contains(entry.getKey())) + .map(Map.Entry::getValue) + .toList(); + } + + public int getGetByTodoKeyCalls() { + return getByTodoKeyCalls; + } + + public int getFindByKeysCalls() { + return findByKeysCalls; + } + public List findByOperatorId(long operatorId) { return cache.values().stream().filter(record -> record.getCurrentOperatorId() == operatorId).toList(); } diff --git a/flow-engine-framework/src/test/java/com/codingapi/flow/service/FlowIssue210SubProcessPerformanceTest.java b/flow-engine-framework/src/test/java/com/codingapi/flow/service/FlowIssue210SubProcessPerformanceTest.java index 0c237d13..497d918d 100644 --- a/flow-engine-framework/src/test/java/com/codingapi/flow/service/FlowIssue210SubProcessPerformanceTest.java +++ b/flow-engine-framework/src/test/java/com/codingapi/flow/service/FlowIssue210SubProcessPerformanceTest.java @@ -39,6 +39,7 @@ import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertTrue; /** * 问题 210 复现测试:主流程 C(子流程节点)一次性创建 6 个子流程、每个子流程 B 节点 20 个审批人。 @@ -127,6 +128,30 @@ void shouldMeasureSubProcessBulkCreationTimeWithGatewayDelay() { + " (网关延时 " + GATEWAY_DELAY_MS + "ms/次)"); } + /** + * 测试目标:验证保存待办时不再逐条 getByTodoKey(N+1),而是改用批量 findByKeys。 + * 前置条件:网关不带延时(本类每次测试都新建 factory),强制执行 bulk 场景。 + * 执行步骤:完成与 {@code shouldMeasureSubProcessBulkCreationTimeWithGatewayDelay} 相同的 + * 6 子流程 × 20 审批人(120 条待办)创建。 + * 期望断言:逐条 getByTodoKey 的调用次数远小于待办数量(仅来自节点完成清理路径), + * 且调用过批量 findByKeys。 + */ + @Test + void shouldBatchLoadExistingTodosInsteadOfPerRecordGet() { + Map data = new LinkedHashMap<>(); + data.put("content", "parent"); + data.put("approvers", approvers.stream().map(User::getUserId).toList()); + long parentRecordId = createParent(data); + approveMain(factory.flowRecordRepository.get(parentRecordId), parentStart, data); + FlowRecord parentBRecord = findTodo(initiator, parentB.getId()); + approveMain(parentBRecord, parentB, data); + + assertTrue(factory.flowTodoRecordRepository.getGetByTodoKeyCalls() < APPROVER_COUNT, + "批量创建待办时应以批量 findByKeys 替代逐条 getByTodoKey,避免 N+1"); + assertTrue(factory.flowTodoRecordRepository.getFindByKeysCalls() >= 1, + "应调用批量 findByKeys 加载已存在待办"); + } + // ---------- 网关延时 + 计数 ---------- private static class DelayedGateway implements FlowOperatorGateway { diff --git a/flow-engine-starter-infra/src/main/java/com/codingapi/flow/infra/jpa/FlowTodoRecordEntityRepository.java b/flow-engine-starter-infra/src/main/java/com/codingapi/flow/infra/jpa/FlowTodoRecordEntityRepository.java index 48221d6b..b6153b47 100644 --- a/flow-engine-starter-infra/src/main/java/com/codingapi/flow/infra/jpa/FlowTodoRecordEntityRepository.java +++ b/flow-engine-starter-infra/src/main/java/com/codingapi/flow/infra/jpa/FlowTodoRecordEntityRepository.java @@ -6,10 +6,15 @@ import org.springframework.data.domain.PageRequest; import org.springframework.data.jpa.repository.Query; +import java.util.List; + public interface FlowTodoRecordEntityRepository extends FastRepository { FlowTodoRecordEntity getByTodoKey(String todoKey); + @Query("from FlowTodoRecordEntity r where r.todoKey in ?1") + List findByKeys(List keys); + @Query("from FlowTodoRecordEntity r where r.currentOperatorId = ?1") Page findTodoRecordPage(long currentOperatorId, PageRequest pageRequest); diff --git a/flow-engine-starter-infra/src/main/java/com/codingapi/flow/infra/repository/impl/FlowTodoRecordRepositoryImpl.java b/flow-engine-starter-infra/src/main/java/com/codingapi/flow/infra/repository/impl/FlowTodoRecordRepositoryImpl.java index dea2aa94..b49fa3b5 100644 --- a/flow-engine-starter-infra/src/main/java/com/codingapi/flow/infra/repository/impl/FlowTodoRecordRepositoryImpl.java +++ b/flow-engine-starter-infra/src/main/java/com/codingapi/flow/infra/repository/impl/FlowTodoRecordRepositoryImpl.java @@ -26,6 +26,13 @@ public FlowTodoRecord getByTodoKey(String key) { return FlowTodoRecordConvertor.convert(flowTodoRecordEntityRepository.getByTodoKey(key)); } + @Override + public List findByKeys(List keys) { + return flowTodoRecordEntityRepository.findByKeys(keys).stream() + .map(FlowTodoRecordConvertor::convert) + .toList(); + } + @Override public void delete(FlowTodoRecord margeRecord) { flowTodoRecordEntityRepository.deleteById(margeRecord.getId()); From db1b2c5630a202ffba660bfef600d28ec7e2d3e6 Mon Sep 17 00:00:00 2001 From: lorne <1991wangliang@gmail.com> Date: Mon, 17 Aug 2026 18:22:51 +0800 Subject: [PATCH 3/4] =?UTF-8?q?chore:=20=E6=9B=B4=E6=96=B0=20flow-frontend?= =?UTF-8?q?=20=E5=AD=90=E6=A8=A1=E5=9D=97=EF=BC=880.2.0=20=E7=89=88?= =?UTF-8?q?=E6=9C=AC=E5=8D=87=E7=BA=A7=20+=20=E6=8E=A5=E5=8F=A3=E8=AF=B7?= =?UTF-8?q?=E6=B1=82=E8=B6=85=E6=97=B6=E6=97=B6=E9=97=B4=E6=94=AF=E6=8C=81?= =?UTF-8?q?=E5=8A=A8=E6=80=81=E9=85=8D=E7=BD=AE=EF=BC=89issue=20#211?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 接口请求超时时间可通过 localStorage 动态配置,未配置时回退默认值 - 发布 @coding-flow/*@0.2.0,去除各调用点 10000ms 硬编码 Co-Authored-By: Claude --- flow-frontend | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/flow-frontend b/flow-frontend index 4b204551..27b49956 160000 --- a/flow-frontend +++ b/flow-frontend @@ -1 +1 @@ -Subproject commit 4b204551b494be4817351bd8f9d9f6fb080c428d +Subproject commit 27b49956dbeedc86d25bb6e5dc07adca800a462c From b51dd2d809ccbeb5b35548f33fe7b89173d26e24 Mon Sep 17 00:00:00 2001 From: lorne <1991wangliang@gmail.com> Date: Mon, 17 Aug 2026 19:46:44 +0800 Subject: [PATCH 4/4] Update flow-frontend to latest main --- flow-frontend | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/flow-frontend b/flow-frontend index 27b49956..88f44a61 160000 --- a/flow-frontend +++ b/flow-frontend @@ -1 +1 @@ -Subproject commit 27b49956dbeedc86d25bb6e5dc07adca800a462c +Subproject commit 88f44a61e0078eba6dcea8e6c1f4fa31af32bed9