From 046e42400e0cf01baa0a5d81d26c585ac56a6f40 Mon Sep 17 00:00:00 2001 From: Edvard Fonsell Date: Mon, 8 Jun 2026 00:07:49 +0300 Subject: [PATCH] do executor recovery in chunks to avoid #718 --- .../engine/internal/dao/WorkflowInstanceDao.java | 13 +++++++++---- 1 file changed, 9 insertions(+), 4 deletions(-) diff --git a/nflow-engine/src/main/java/io/nflow/engine/internal/dao/WorkflowInstanceDao.java b/nflow-engine/src/main/java/io/nflow/engine/internal/dao/WorkflowInstanceDao.java index e1631e475..943a3233f 100644 --- a/nflow-engine/src/main/java/io/nflow/engine/internal/dao/WorkflowInstanceDao.java +++ b/nflow-engine/src/main/java/io/nflow/engine/internal/dao/WorkflowInstanceDao.java @@ -399,11 +399,16 @@ public void recoverWorkflowInstancesFromDeadNodes() { } WorkflowInstanceAction.Builder builder = new WorkflowInstanceAction.Builder().setExecutionStart(now()).setExecutionEnd(now()) .setType(recovery).setStateText("Recovered"); - for (InstanceInfo instance : getRecoverableWorkflowInstances(recoverableExecutorIds)) { - WorkflowInstanceAction action = builder.setState(instance.state()).setWorkflowInstanceId(instance.id()).build(); - recoverWorkflowInstance(instance.id(), instance.executorId(), action); + List recoverableExecutorIdList = new ArrayList<>(recoverableExecutorIds); + for (int from = 0; from < recoverableExecutorIdList.size(); from += 100) { + List recoverableExecutorIdChunk = recoverableExecutorIdList.subList(from, + min(from + 100, recoverableExecutorIdList.size())); + for (InstanceInfo instance : getRecoverableWorkflowInstances(recoverableExecutorIdChunk)) { + WorkflowInstanceAction action = builder.setState(instance.state()).setWorkflowInstanceId(instance.id()).build(); + recoverWorkflowInstance(instance.id(), instance.executorId(), action); + } + recoverableExecutorIdChunk.forEach(executorInfo::markRecovered); } - recoverableExecutorIds.forEach(executorInfo::markRecovered); } private List getRecoverableWorkflowInstances(Collection executorsIds) {