diff --git a/fe/fe-core/src/main/java/org/apache/doris/qe/NereidsCoordinator.java b/fe/fe-core/src/main/java/org/apache/doris/qe/NereidsCoordinator.java index 641d433c429754..b52ab3eb3fbbb6 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/qe/NereidsCoordinator.java +++ b/fe/fe-core/src/main/java/org/apache/doris/qe/NereidsCoordinator.java @@ -392,7 +392,7 @@ public Map getBeToInstancesNum() { for (MultiFragmentsPipelineTask beTasks : executionTask.getChildrenTasks().values()) { TNetworkAddress brpcAddress = beTasks.getBackend().getBrpcAddress(); String brpcAddrString = brpcAddress.hostname.concat(":").concat("" + brpcAddress.port); - result.put(brpcAddrString, beTasks.getChildrenTasks().size()); + result.put(brpcAddrString, beTasks.getInstanceNum()); } } return result; diff --git a/fe/fe-core/src/main/java/org/apache/doris/qe/runtime/AbstractRuntimeTask.java b/fe/fe-core/src/main/java/org/apache/doris/qe/runtime/AbstractRuntimeTask.java index 1059440792e529..7e4fb96d27cb78 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/qe/runtime/AbstractRuntimeTask.java +++ b/fe/fe-core/src/main/java/org/apache/doris/qe/runtime/AbstractRuntimeTask.java @@ -33,6 +33,12 @@ public void execute() throws Throwable { } } + public Integer getInstanceNum() { + return childrenTasks.allTasks().stream() + .mapToInt(Child::getInstanceNum) + .sum(); + } + public Map getChildrenTasks() { return childrenTasks.allTaskMap(); } diff --git a/fe/fe-core/src/main/java/org/apache/doris/qe/runtime/SingleFragmentPipelineTask.java b/fe/fe-core/src/main/java/org/apache/doris/qe/runtime/SingleFragmentPipelineTask.java index 2fba0654af9d67..1aa4c7f0478e50 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/qe/runtime/SingleFragmentPipelineTask.java +++ b/fe/fe-core/src/main/java/org/apache/doris/qe/runtime/SingleFragmentPipelineTask.java @@ -82,6 +82,10 @@ public int getFragmentId() { return fragmentId; } + public Integer getInstanceNum() { + return instanceIds.size(); + } + public List buildFragmentInstanceInfo() { TNetworkAddress address = new TNetworkAddress(backend.getHost(), backend.getBePort()); List infos = Lists.newArrayListWithCapacity(instanceIds.size());