This is an automated email from the ASF dual-hosted git repository.

qiaojialin pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iotdb.git


The following commit(s) were added to refs/heads/master by this push:
     new 8a8e727a20 [IOTDB-3370] Fix count nodes number (#6178)
8a8e727a20 is described below

commit 8a8e727a2071118d05a4f4d0727f27082cd69d7a
Author: ZhangHongYin <[email protected]>
AuthorDate: Tue Jun 7 12:06:22 2022 +0800

    [IOTDB-3370] Fix count nodes number (#6178)
---
 .../schema/NodeManageMemoryMergeOperator.java      |  7 ++--
 .../operator/schema/NodePathsConvertOperator.java  |  5 +--
 .../operator/schema/NodePathsCountOperator.java    | 38 ++++++++++++++++------
 3 files changed, 31 insertions(+), 19 deletions(-)

diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/schema/NodeManageMemoryMergeOperator.java
 
b/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/schema/NodeManageMemoryMergeOperator.java
index 5532092483..e6e15b2e3f 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/schema/NodeManageMemoryMergeOperator.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/schema/NodeManageMemoryMergeOperator.java
@@ -36,16 +36,14 @@ import static java.util.Objects.requireNonNull;
 
 public class NodeManageMemoryMergeOperator implements ProcessOperator {
   private final OperatorContext operatorContext;
-  private Set<String> data;
+  private final Set<String> data;
   private final Operator child;
-  private boolean isFinished;
 
   public NodeManageMemoryMergeOperator(
       OperatorContext operatorContext, Set<String> data, Operator child) {
     this.operatorContext = requireNonNull(operatorContext, "operatorContext is 
null");
     this.data = data;
     this.child = requireNonNull(child, "child operator is null");
-    isFinished = false;
   }
 
   @Override
@@ -60,7 +58,6 @@ public class NodeManageMemoryMergeOperator implements 
ProcessOperator {
 
   @Override
   public TsBlock next() {
-    isFinished = true;
     TsBlock block = child.next();
     TsBlockBuilder tsBlockBuilder =
         new 
TsBlockBuilder(HeaderConstant.showChildPathsHeader.getRespDataTypes());
@@ -90,6 +87,6 @@ public class NodeManageMemoryMergeOperator implements 
ProcessOperator {
 
   @Override
   public boolean isFinished() {
-    return isFinished || child.isFinished();
+    return child.isFinished();
   }
 }
diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/schema/NodePathsConvertOperator.java
 
b/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/schema/NodePathsConvertOperator.java
index 3ad0619b7b..bb1b50228d 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/schema/NodePathsConvertOperator.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/schema/NodePathsConvertOperator.java
@@ -41,12 +41,10 @@ public class NodePathsConvertOperator implements 
ProcessOperator {
 
   private final OperatorContext operatorContext;
   private final Operator child;
-  private boolean isFinished;
 
   public NodePathsConvertOperator(OperatorContext operatorContext, Operator 
child) {
     this.operatorContext = requireNonNull(operatorContext, "operatorContext is 
null");
     this.child = requireNonNull(child, "child operator is null");
-    isFinished = false;
   }
 
   @Override
@@ -61,7 +59,6 @@ public class NodePathsConvertOperator implements 
ProcessOperator {
 
   @Override
   public TsBlock next() {
-    isFinished = true;
     TsBlock block = child.next();
     TsBlockBuilder tsBlockBuilder =
         new 
TsBlockBuilder(HeaderConstant.showChildNodesHeader.getRespDataTypes());
@@ -95,6 +92,6 @@ public class NodePathsConvertOperator implements 
ProcessOperator {
 
   @Override
   public boolean isFinished() {
-    return isFinished || child.isFinished();
+    return child.isFinished();
   }
 }
diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/schema/NodePathsCountOperator.java
 
b/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/schema/NodePathsCountOperator.java
index 6276c4d5eb..3c0614589d 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/schema/NodePathsCountOperator.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/schema/NodePathsCountOperator.java
@@ -28,6 +28,9 @@ import 
org.apache.iotdb.tsfile.read.common.block.TsBlockBuilder;
 
 import com.google.common.util.concurrent.ListenableFuture;
 
+import java.util.HashSet;
+import java.util.Set;
+
 import static java.util.Objects.requireNonNull;
 
 public class NodePathsCountOperator implements ProcessOperator {
@@ -35,11 +38,13 @@ public class NodePathsCountOperator implements 
ProcessOperator {
   private final OperatorContext operatorContext;
   private final Operator child;
   private boolean isFinished;
+  private final Set<String> nodePaths;
 
   public NodePathsCountOperator(OperatorContext operatorContext, Operator 
child) {
     this.operatorContext = requireNonNull(operatorContext, "operatorContext is 
null");
     this.child = requireNonNull(child, "child operator is null");
-    isFinished = false;
+    this.isFinished = false;
+    this.nodePaths = new HashSet<>();
   }
 
   @Override
@@ -47,27 +52,40 @@ public class NodePathsCountOperator implements 
ProcessOperator {
     return operatorContext;
   }
 
-  @Override
-  public ListenableFuture<Void> isBlocked() {
-    return child.isBlocked();
-  }
-
   @Override
   public TsBlock next() {
     isFinished = true;
-    TsBlock block = child.next();
     TsBlockBuilder tsBlockBuilder =
         new TsBlockBuilder(HeaderConstant.countNodesHeader.getRespDataTypes());
 
     tsBlockBuilder.getTimeColumnBuilder().writeLong(0L);
-    tsBlockBuilder.getColumnBuilder(0).writeInt(block.getPositionCount());
+    tsBlockBuilder.getColumnBuilder(0).writeInt(nodePaths.size());
     tsBlockBuilder.declarePosition();
     return tsBlockBuilder.build();
   }
 
   @Override
   public boolean hasNext() {
-    return child.hasNext();
+    return !isFinished;
+  }
+
+  @Override
+  public ListenableFuture<Void> isBlocked() {
+    ListenableFuture<Void> blocked = child.isBlocked();
+    while (child.hasNext() && blocked.isDone()) {
+      TsBlock tsBlock = child.next();
+      if (null != tsBlock && !tsBlock.isEmpty()) {
+        for (int i = 0; i < tsBlock.getPositionCount(); i++) {
+          String path = tsBlock.getColumn(0).getBinary(i).toString();
+          nodePaths.add(path);
+        }
+      }
+      blocked = child.isBlocked();
+    }
+    if (!blocked.isDone()) {
+      return blocked;
+    }
+    return NOT_BLOCKED;
   }
 
   @Override
@@ -77,6 +95,6 @@ public class NodePathsCountOperator implements 
ProcessOperator {
 
   @Override
   public boolean isFinished() {
-    return isFinished || child.isFinished();
+    return isFinished;
   }
 }

Reply via email to