This is an automated email from the ASF dual-hosted git repository.
jackietien pushed a commit to branch mpp-ty
in repository https://gitbox.apache.org/repos/asf/iotdb.git
The following commit(s) were added to refs/heads/mpp-ty by this push:
new 5dc88c3 add filter operator implementation and visitor in
LocalExecutionPlanner
5dc88c3 is described below
commit 5dc88c377ac61b2c4b98ac620fb2e50d19a54655
Author: JackieTien97 <[email protected]>
AuthorDate: Sat Mar 19 21:53:33 2022 +0800
add filter operator implementation and visitor in LocalExecutionPlanner
---
.../org/apache/iotdb/db/mpp/common/TsBlock.java | 25 +++++++++++++++++
.../iotdb/db/mpp/execution/InstanceContext.java | 27 ++++++++++++++++++
.../iotdb/db/mpp/operator/OperatorContext.java | 4 +++
.../db/mpp/operator/process/LimitOperator.java | 32 +++++++++++++++++++---
.../db/mpp/sql/planner/LocalExecutionPlanner.java | 8 +++++-
5 files changed, 91 insertions(+), 5 deletions(-)
diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/common/TsBlock.java
b/server/src/main/java/org/apache/iotdb/db/mpp/common/TsBlock.java
index aa40205..7ea49ca 100644
--- a/server/src/main/java/org/apache/iotdb/db/mpp/common/TsBlock.java
+++ b/server/src/main/java/org/apache/iotdb/db/mpp/common/TsBlock.java
@@ -20,6 +20,8 @@ package org.apache.iotdb.db.mpp.common;
import org.apache.iotdb.tsfile.read.common.RowRecord;
+import static java.lang.String.format;
+
/**
* Intermediate result for most of ExecOperators. The Tablet contains data
from one or more columns
* and constructs them as a row based view The columns can be series,
aggregation result for one
@@ -33,6 +35,8 @@ public class TsBlock {
// Describe the column info
private TsBlockMetadata metadata;
+ private int count;
+
public boolean hasNext() {
return false;
}
@@ -45,4 +49,25 @@ public class TsBlock {
public TsBlockMetadata getMetadata() {
return metadata;
}
+
+ public int getCount() {
+ return count;
+ }
+
+ /**
+ * TODO has not been implemented yet
+ *
+ * @param positionOffset start offset
+ * @param length slice length
+ * @return view of current TsBlock start from positionOffset to
positionOffset + length
+ */
+ public TsBlock getRegion(int positionOffset, int length) {
+ if (positionOffset < 0 || length < 0 || positionOffset + length > count) {
+ throw new IndexOutOfBoundsException(
+ format(
+ "Invalid position %s and length %s in page with %s positions",
+ positionOffset, length, count));
+ }
+ return this;
+ }
}
diff --git
a/server/src/main/java/org/apache/iotdb/db/mpp/execution/InstanceContext.java
b/server/src/main/java/org/apache/iotdb/db/mpp/execution/InstanceContext.java
index 3a5f99d..2299a1c 100644
---
a/server/src/main/java/org/apache/iotdb/db/mpp/execution/InstanceContext.java
+++
b/server/src/main/java/org/apache/iotdb/db/mpp/execution/InstanceContext.java
@@ -19,11 +19,22 @@
package org.apache.iotdb.db.mpp.execution;
import org.apache.iotdb.db.mpp.common.InstanceId;
+import org.apache.iotdb.db.mpp.operator.OperatorContext;
+import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanNodeId;
+
+import java.util.ArrayList;
+import java.util.List;
+
+import static com.google.common.base.Preconditions.checkArgument;
public class InstanceContext {
private InstanceId id;
+ // TODO if we split one fragment instance into multiple pipelines to run, we
need to replace it
+ // with CopyOnWriteArrayList or some other thread safe data structure
+ private final List<OperatorContext> operatorContexts = new ArrayList<>();
+
private final long createNanos = System.nanoTime();
// private final GcMonitor gcMonitor;
@@ -37,4 +48,20 @@ public class InstanceContext {
public InstanceContext(InstanceId id) {
this.id = id;
}
+
+ public OperatorContext addOperatorContext(
+ int operatorId, PlanNodeId planNodeId, String operatorType) {
+ checkArgument(operatorId >= 0, "operatorId is negative");
+
+ for (OperatorContext operatorContext : operatorContexts) {
+ checkArgument(
+ operatorId != operatorContext.getOperatorId(),
+ "A context already exists for operatorId %s",
+ operatorId);
+ }
+
+ OperatorContext operatorContext = new OperatorContext(operatorId,
planNodeId, operatorType);
+ operatorContexts.add(operatorContext);
+ return operatorContext;
+ }
}
diff --git
a/server/src/main/java/org/apache/iotdb/db/mpp/operator/OperatorContext.java
b/server/src/main/java/org/apache/iotdb/db/mpp/operator/OperatorContext.java
index 472fa05..feedc13 100644
--- a/server/src/main/java/org/apache/iotdb/db/mpp/operator/OperatorContext.java
+++ b/server/src/main/java/org/apache/iotdb/db/mpp/operator/OperatorContext.java
@@ -36,4 +36,8 @@ public class OperatorContext {
this.planNodeId = planNodeId;
this.operatorType = operatorType;
}
+
+ public int getOperatorId() {
+ return operatorId;
+ }
}
diff --git
a/server/src/main/java/org/apache/iotdb/db/mpp/operator/process/LimitOperator.java
b/server/src/main/java/org/apache/iotdb/db/mpp/operator/process/LimitOperator.java
index 0efd325..1ebca2c 100644
---
a/server/src/main/java/org/apache/iotdb/db/mpp/operator/process/LimitOperator.java
+++
b/server/src/main/java/org/apache/iotdb/db/mpp/operator/process/LimitOperator.java
@@ -19,29 +19,53 @@
package org.apache.iotdb.db.mpp.operator.process;
import org.apache.iotdb.db.mpp.common.TsBlock;
+import org.apache.iotdb.db.mpp.operator.Operator;
import org.apache.iotdb.db.mpp.operator.OperatorContext;
import com.google.common.util.concurrent.ListenableFuture;
+import static com.google.common.base.Preconditions.checkArgument;
+import static java.util.Objects.requireNonNull;
+
public class LimitOperator implements ProcessOperator {
+
+ private final OperatorContext operatorContext;
+ private long remainingLimit;
+ private final Operator child;
+
+ public LimitOperator(OperatorContext operatorContext, long limit, Operator
child) {
+ this.operatorContext = requireNonNull(operatorContext, "operatorContext is
null");
+ checkArgument(limit >= 0, "limit must be at least zero");
+ this.remainingLimit = limit;
+ this.child = requireNonNull(child, "child operator is null");
+ }
+
@Override
public OperatorContext getOperatorContext() {
- return null;
+ return operatorContext;
}
@Override
public ListenableFuture<Void> isBlocked() {
- return ProcessOperator.super.isBlocked();
+ return child.isBlocked();
}
@Override
public TsBlock next() {
- return null;
+ TsBlock block = child.next();
+ TsBlock res = block;
+ if (block.getCount() <= remainingLimit) {
+ remainingLimit -= block.getCount();
+ } else {
+ res = block.getRegion(0, (int) remainingLimit);
+ remainingLimit = 0;
+ }
+ return res;
}
@Override
public boolean hasNext() {
- return false;
+ return child.hasNext();
}
@Override
diff --git
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/LocalExecutionPlanner.java
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/LocalExecutionPlanner.java
index a4dd7539..069783a 100644
---
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/LocalExecutionPlanner.java
+++
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/LocalExecutionPlanner.java
@@ -20,6 +20,7 @@ package org.apache.iotdb.db.mpp.sql.planner;
import org.apache.iotdb.db.mpp.execution.InstanceContext;
import org.apache.iotdb.db.mpp.operator.Operator;
+import org.apache.iotdb.db.mpp.operator.process.LimitOperator;
import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanNode;
import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanVisitor;
import org.apache.iotdb.db.mpp.sql.planner.plan.node.process.*;
@@ -86,7 +87,12 @@ public class LocalExecutionPlanner {
@Override
public Operator visitLimit(LimitNode node, LocalExecutionPlanContext
context) {
- return super.visitLimit(node, context);
+ Operator child = node.getChild().accept(this, context);
+ return new LimitOperator(
+ context.taskContext.addOperatorContext(
+ context.getNextOperatorId(), node.getId(),
LimitOperator.class.getSimpleName()),
+ node.getLimit(),
+ child);
}
@Override