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;
}
}