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

Reply via email to