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

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

commit e374570ecf5be27059b4ba433993d4cf39aeee2f
Author: Jinrui.Zhang <[email protected]>
AuthorDate: Sun Mar 20 10:28:06 2022 +0800

    complete SourceRewriter
---
 .../org/apache/iotdb/db/mpp/common/Analysis.java   | 36 ++++++++-
 .../ThriftSinkNode.java => common/DataRegion.java} | 46 +++++-------
 .../{Analysis.java => DataRegionTimeSlice.java}    |  9 ++-
 .../common/{Analysis.java => SchemaRegion.java}    | 13 +++-
 .../mpp/sql/planner/plan/DistributionPlanner.java  | 87 +++++++++++++++++++---
 .../db/mpp/sql/planner/plan/node/PlanNode.java     |  4 +
 .../{PlanNode.java => SimplePlanNodeRewriter.java} | 40 +++++-----
 .../planner/plan/node/process/AggregateNode.java   | 10 +++
 .../planner/plan/node/process/DeviceMergeNode.java | 10 +++
 .../sql/planner/plan/node/process/FillNode.java    | 10 +++
 .../sql/planner/plan/node/process/FilterNode.java  | 10 +++
 .../planner/plan/node/process/FilterNullNode.java  | 10 +++
 .../plan/node/process/GroupByLevelNode.java        | 10 +++
 .../sql/planner/plan/node/process/LimitNode.java   | 10 +++
 .../sql/planner/plan/node/process/OffsetNode.java  | 10 +++
 .../sql/planner/plan/node/process/SortNode.java    | 10 +++
 .../planner/plan/node/process/TimeJoinNode.java    | 10 +++
 .../sql/planner/plan/node/sink/CsvSinkNode.java    | 10 +++
 .../planner/plan/node/sink/FragmentSinkNode.java   | 10 +++
 .../sql/planner/plan/node/sink/ThriftSinkNode.java | 10 +++
 .../planner/plan/node/source/CsvSourceNode.java    | 10 +++
 .../plan/node/source/SeriesAggregateScanNode.java  | 10 +++
 .../planner/plan/node/source/SeriesScanNode.java   | 30 ++++++++
 23 files changed, 351 insertions(+), 64 deletions(-)

diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/common/Analysis.java 
b/server/src/main/java/org/apache/iotdb/db/mpp/common/Analysis.java
index 178402a..08fbc45 100644
--- a/server/src/main/java/org/apache/iotdb/db/mpp/common/Analysis.java
+++ b/server/src/main/java/org/apache/iotdb/db/mpp/common/Analysis.java
@@ -18,5 +18,39 @@
  */
 package org.apache.iotdb.db.mpp.common;
 
+import org.apache.iotdb.db.metadata.path.PartialPath;
+import org.apache.iotdb.tsfile.read.filter.basic.Filter;
+
+import java.util.*;
+
 /** Analysis used for planning a query. TODO: This class may need to store 
more info for a query. */
-public class Analysis {}
+public class Analysis {
+    // Description for each series. Such as dataType, existence
+
+    // Data distribution info for each series. Series -> [VSG, VSG]
+
+    // Map<PartialPath, List<FullPath>> Used to remove asterisk
+
+    // Statement
+    private String statement;
+
+    // DataPartitionInfo
+    private Map<String, Map<DataRegionTimeSlice, List<DataRegion>>> 
dataPartitionInfo;
+
+    // SchemaPartitionInfo
+    private Map<String, List<SchemaRegion>> schemaPartitionInfo;
+
+
+    public Set<DataRegion> getPartitionInfo(PartialPath seriesPath, Filter 
timefilter) {
+        if (timefilter == null) {
+            //TODO: (xingtanzjr) we need to have a method to get the 
deviceGroup by device
+            String deviceGroup = seriesPath.getDevice();
+            Set<DataRegion> result = new HashSet<>();
+            
this.dataPartitionInfo.get(deviceGroup).values().forEach(result::addAll);
+            return result;
+        } else {
+            //TODO: (xingtanzjr) complete this branch
+            return null;
+        }
+    }
+}
diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/sink/ThriftSinkNode.java
 b/server/src/main/java/org/apache/iotdb/db/mpp/common/DataRegion.java
similarity index 56%
copy from 
server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/sink/ThriftSinkNode.java
copy to server/src/main/java/org/apache/iotdb/db/mpp/common/DataRegion.java
index bb343f4..a46afdf 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/sink/ThriftSinkNode.java
+++ b/server/src/main/java/org/apache/iotdb/db/mpp/common/DataRegion.java
@@ -16,33 +16,25 @@
  * specific language governing permissions and limitations
  * under the License.
  */
-package org.apache.iotdb.db.mpp.sql.planner.plan.node.sink;
 
-import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanNode;
-import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanNodeId;
+package org.apache.iotdb.db.mpp.common;
 
-import java.util.List;
-
-/** not implemented in current IoTDB yet */
-public class ThriftSinkNode extends SinkNode {
-
-  public ThriftSinkNode(PlanNodeId id) {
-    super(id);
-  }
-
-  @Override
-  public List<PlanNode> getChildren() {
-    return null;
-  }
-
-  @Override
-  public List<String> getOutputColumnNames() {
-    return null;
-  }
-
-  @Override
-  public void close() throws Exception {}
-
-  @Override
-  public void send() {}
+/**
+ * This class is used to represent the data partition info including the 
DataRegionId and physical node IP address
+ */
+//TODO: (xingtanzjr) This class should be substituted with the class defined 
in Consensus level
+public class DataRegion {
+    private Integer dataRegionId;
+    private String endpoint;
+
+    public int hashCode() {
+        return dataRegionId.hashCode();
+    }
+
+    public boolean equals(Object obj) {
+        if (obj instanceof DataRegion) {
+            return this.dataRegionId.equals(((DataRegion)obj).dataRegionId);
+        }
+        return false;
+    }
 }
diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/common/Analysis.java 
b/server/src/main/java/org/apache/iotdb/db/mpp/common/DataRegionTimeSlice.java
similarity index 79%
copy from server/src/main/java/org/apache/iotdb/db/mpp/common/Analysis.java
copy to 
server/src/main/java/org/apache/iotdb/db/mpp/common/DataRegionTimeSlice.java
index 178402a..52f48cc 100644
--- a/server/src/main/java/org/apache/iotdb/db/mpp/common/Analysis.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/common/DataRegionTimeSlice.java
@@ -7,7 +7,7 @@
  * "License"); you may not use this file except in compliance
  * with the License.  You may obtain a copy of the License at
  *
- *      http://www.apache.org/licenses/LICENSE-2.0
+ *     http://www.apache.org/licenses/LICENSE-2.0
  *
  * Unless required by applicable law or agreed to in writing,
  * software distributed under the License is distributed on an
@@ -16,7 +16,10 @@
  * specific language governing permissions and limitations
  * under the License.
  */
+
 package org.apache.iotdb.db.mpp.common;
 
-/** Analysis used for planning a query. TODO: This class may need to store 
more info for a query. */
-public class Analysis {}
+//TODO: (xingtanzjr) This class should be substituted with the class defined 
in Consensus level
+public class DataRegionTimeSlice {
+    long startTimestamp;
+}
diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/common/Analysis.java 
b/server/src/main/java/org/apache/iotdb/db/mpp/common/SchemaRegion.java
similarity index 68%
copy from server/src/main/java/org/apache/iotdb/db/mpp/common/Analysis.java
copy to server/src/main/java/org/apache/iotdb/db/mpp/common/SchemaRegion.java
index 178402a..8352f10 100644
--- a/server/src/main/java/org/apache/iotdb/db/mpp/common/Analysis.java
+++ b/server/src/main/java/org/apache/iotdb/db/mpp/common/SchemaRegion.java
@@ -7,7 +7,7 @@
  * "License"); you may not use this file except in compliance
  * with the License.  You may obtain a copy of the License at
  *
- *      http://www.apache.org/licenses/LICENSE-2.0
+ *     http://www.apache.org/licenses/LICENSE-2.0
  *
  * Unless required by applicable law or agreed to in writing,
  * software distributed under the License is distributed on an
@@ -16,7 +16,14 @@
  * specific language governing permissions and limitations
  * under the License.
  */
+
 package org.apache.iotdb.db.mpp.common;
 
-/** Analysis used for planning a query. TODO: This class may need to store 
more info for a query. */
-public class Analysis {}
+/**
+ * This class is used to represent the schema partition info including the 
DataRegionId and physical node IP address
+ */
+//TODO: (xingtanzjr) This class should be substituted with the class defined 
in Consensus level
+public class SchemaRegion {
+    private Integer DataRegionId;
+    private String endpoint;
+}
diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/DistributionPlanner.java
 
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/DistributionPlanner.java
index fc428cf..d9fa64c 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/DistributionPlanner.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/DistributionPlanner.java
@@ -19,17 +19,86 @@
 package org.apache.iotdb.db.mpp.sql.planner.plan;
 
 import org.apache.iotdb.db.mpp.common.Analysis;
+import org.apache.iotdb.db.mpp.common.DataRegion;
+import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanNode;
+import org.apache.iotdb.db.mpp.sql.planner.plan.node.SimplePlanNodeRewriter;
+import org.apache.iotdb.db.mpp.sql.planner.plan.node.process.TimeJoinNode;
+import 
org.apache.iotdb.db.mpp.sql.planner.plan.node.source.SeriesAggregateScanNode;
+import org.apache.iotdb.db.mpp.sql.planner.plan.node.source.SeriesScanNode;
+
+import java.util.*;
 
 public class DistributionPlanner {
-  private Analysis analysis;
-  private LogicalQueryPlan logicalPlan;
+    private Analysis analysis;
+    private LogicalQueryPlan logicalPlan;
+
+    public DistributionPlanner(Analysis analysis, LogicalQueryPlan 
logicalPlan) {
+        this.analysis = analysis;
+        this.logicalPlan = logicalPlan;
+    }
+
+    public DistributedQueryPlan planFragments() {
+        return null;
+    }
+
+    private class SourceRewriter extends 
SimplePlanNodeRewriter<DistributionPlanContext> {
+        public PlanNode visitTimeJoin(TimeJoinNode node, 
DistributionPlanContext context) {
+            TimeJoinNode root = (TimeJoinNode) node.clone();
+
+            // Step 1: Get all source nodes. For the node which is not source, 
add it as the child of current TimeJoinNode
+            List<SeriesScanNode> sources = new ArrayList<>();
+            for (PlanNode child : node.getChildren()) {
+                if (child instanceof SeriesScanNode) {
+                    // If the child is SeriesScanNode, we need to check 
whether this node should be seperated into several splits.
+                    SeriesScanNode handle = (SeriesScanNode) child;
+                    Set<DataRegion> dataDistribution = 
analysis.getPartitionInfo(handle.getSeriesPath(), handle.getTimeFilter());
+                    // If the size of dataDistribution is m, this 
SeriesScanNode should be seperated into m SeriesScanNode.
+                    for (DataRegion dataRegion : dataDistribution) {
+                        SeriesScanNode split = (SeriesScanNode) handle.clone();
+                        split.setDataRegion(dataRegion);
+                        sources.add(split);
+                    }
+                } else if (child instanceof SeriesAggregateScanNode) {
+                    //TODO: (xingtanzjr) We should do the same thing for 
SeriesAggregateScanNode. Consider to make SeriesAggregateScanNode
+                    // and SeriesScanNode to derived from the same parent 
Class because they have similar process logic in many scenarios
+                } else {
+                    // In a general logical query plan, the children of 
TimeJoinNode should only be SeriesScanNode or SeriesAggregateScanNode
+                    // So this branch should not be touched.
+                    root.addChild(generateDistributedPlan(child, context));
+                }
+            }
+
+            // Step 2: For the source nodes, group them by the DataRegion.
+            Map<DataRegion, List<SeriesScanNode>> sourceGroup = new 
HashMap<>();
+            sources.forEach(source -> {
+                List<SeriesScanNode> group = 
sourceGroup.containsKey(source.getDataRegion()) ?
+                        sourceGroup.get(source.getDataRegion()) : new 
ArrayList<>();
+                group.add(source);
+                sourceGroup.put(source.getDataRegion(), group);
+            });
+
+            // Step 3: For the source nodes which belong to same data region, 
add a TimeJoinNode for them and make the
+            // new TimeJoinNode as the child of current TimeJoinNode
+            sourceGroup.forEach((dataRegion, seriesScanNodes) -> {
+                if (seriesScanNodes.size() == 1) {
+                    root.addChild(seriesScanNodes.get(0));
+                } else {
+                    // We clone a TimeJoinNode from root to make the params to 
be consistent
+                    TimeJoinNode parentOfGroup = (TimeJoinNode) root.clone();
+                    seriesScanNodes.forEach(parentOfGroup::addChild);
+                    root.addChild(parentOfGroup);
+                }
+            });
+
+            return root;
+        }
+
+        public PlanNode generateDistributedPlan(PlanNode node, 
DistributionPlanContext context) {
+            return node.accept(this, context);
+        }
+    }
 
-  public DistributionPlanner(Analysis analysis, LogicalQueryPlan logicalPlan) {
-    this.analysis = analysis;
-    this.logicalPlan = logicalPlan;
-  }
+    private class DistributionPlanContext {
 
-  public DistributedQueryPlan planFragments() {
-    return null;
-  }
+    }
 }
diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/PlanNode.java
 
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/PlanNode.java
index 3a3975f..e65af59 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/PlanNode.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/PlanNode.java
@@ -39,6 +39,10 @@ public abstract class PlanNode {
 
   public abstract List<PlanNode> getChildren();
 
+  public abstract PlanNode clone();
+
+  public abstract PlanNode cloneWithChildren(List<PlanNode> children);
+
   public abstract List<String> getOutputColumnNames();
 
   public <R, C> R accept(PlanVisitor<R, C> visitor, C context) {
diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/PlanNode.java
 
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/SimplePlanNodeRewriter.java
similarity index 52%
copy from 
server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/PlanNode.java
copy to 
server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/SimplePlanNodeRewriter.java
index 3a3975f..639e805 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/PlanNode.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/SimplePlanNodeRewriter.java
@@ -7,7 +7,7 @@
  * "License"); you may not use this file except in compliance
  * with the License.  You may obtain a copy of the License at
  *
- *      http://www.apache.org/licenses/LICENSE-2.0
+ *     http://www.apache.org/licenses/LICENSE-2.0
  *
  * Unless required by applicable law or agreed to in writing,
  * software distributed under the License is distributed on an
@@ -16,32 +16,30 @@
  * specific language governing permissions and limitations
  * under the License.
  */
+
 package org.apache.iotdb.db.mpp.sql.planner.plan.node;
 
 import java.util.List;
 
-import static java.util.Objects.requireNonNull;
-
-/** The base class of query executable operators, which is used to compose 
logical query plan. */
-// TODO: consider how to restrict the children type for each type of 
ExecOperator
-public abstract class PlanNode {
-
-  private PlanNodeId id;
-
-  protected PlanNode(PlanNodeId id) {
-    requireNonNull(id, "id is null");
-    this.id = id;
-  }
+import static com.google.common.base.Verify.verifyNotNull;
+import static com.google.common.collect.ImmutableList.toImmutableList;
 
-  public PlanNodeId getId() {
-    return id;
-  }
+public class SimplePlanNodeRewriter<C> extends PlanVisitor<PlanNode, C>{
+    @Override
+    public PlanNode visitPlan(PlanNode node, C context) {
+        return defaultRewrite(node, context);
+    }
 
-  public abstract List<PlanNode> getChildren();
+    public PlanNode defaultRewrite(PlanNode node, C context) {
+        List<PlanNode> children = node.getChildren().stream()
+                .map(child -> rewrite(child, context))
+                .collect(toImmutableList());
 
-  public abstract List<String> getOutputColumnNames();
+        return node.cloneWithChildren(children);
+    }
 
-  public <R, C> R accept(PlanVisitor<R, C> visitor, C context) {
-    return visitor.visitPlan(this, context);
-  }
+    public PlanNode rewrite(PlanNode node, C userContext)
+    {
+        return node.accept(this, userContext);
+    }
 }
diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/AggregateNode.java
 
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/AggregateNode.java
index 9382fbf..7d5705e 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/AggregateNode.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/AggregateNode.java
@@ -63,6 +63,16 @@ public class AggregateNode extends ProcessNode {
   }
 
   @Override
+  public PlanNode clone() {
+    return null;
+  }
+
+  @Override
+  public PlanNode cloneWithChildren(List<PlanNode> children) {
+    return null;
+  }
+
+  @Override
   public List<String> getOutputColumnNames() {
     return columnNames;
   }
diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/DeviceMergeNode.java
 
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/DeviceMergeNode.java
index 54a0cf8..5073d15 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/DeviceMergeNode.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/DeviceMergeNode.java
@@ -64,6 +64,16 @@ public class DeviceMergeNode extends ProcessNode {
   }
 
   @Override
+  public PlanNode clone() {
+    return null;
+  }
+
+  @Override
+  public PlanNode cloneWithChildren(List<PlanNode> children) {
+    return null;
+  }
+
+  @Override
   public List<String> getOutputColumnNames() {
     return columnNames;
   }
diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/FillNode.java
 
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/FillNode.java
index af88895..7d44870 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/FillNode.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/FillNode.java
@@ -45,6 +45,16 @@ public class FillNode extends ProcessNode {
   }
 
   @Override
+  public PlanNode clone() {
+    return null;
+  }
+
+  @Override
+  public PlanNode cloneWithChildren(List<PlanNode> children) {
+    return null;
+  }
+
+  @Override
   public List<String> getOutputColumnNames() {
     return child.getOutputColumnNames();
   }
diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/FilterNode.java
 
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/FilterNode.java
index 4504890..053589d 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/FilterNode.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/FilterNode.java
@@ -47,6 +47,16 @@ public class FilterNode extends ProcessNode {
   }
 
   @Override
+  public PlanNode clone() {
+    return null;
+  }
+
+  @Override
+  public PlanNode cloneWithChildren(List<PlanNode> children) {
+    return null;
+  }
+
+  @Override
   public List<String> getOutputColumnNames() {
     return child.getOutputColumnNames();
   }
diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/FilterNullNode.java
 
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/FilterNullNode.java
index fbf8cc9..9d339ef 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/FilterNullNode.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/FilterNullNode.java
@@ -54,6 +54,16 @@ public class FilterNullNode extends ProcessNode {
   }
 
   @Override
+  public PlanNode clone() {
+    return null;
+  }
+
+  @Override
+  public PlanNode cloneWithChildren(List<PlanNode> children) {
+    return null;
+  }
+
+  @Override
   public List<String> getOutputColumnNames() {
     return child.getOutputColumnNames();
   }
diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/GroupByLevelNode.java
 
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/GroupByLevelNode.java
index 7dcca76..6f3716e 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/GroupByLevelNode.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/GroupByLevelNode.java
@@ -56,6 +56,16 @@ public class GroupByLevelNode extends ProcessNode {
   }
 
   @Override
+  public PlanNode clone() {
+    return null;
+  }
+
+  @Override
+  public PlanNode cloneWithChildren(List<PlanNode> children) {
+    return null;
+  }
+
+  @Override
   public List<String> getOutputColumnNames() {
     return columnNames;
   }
diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/LimitNode.java
 
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/LimitNode.java
index 8fdc9b4..5c5b1bc 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/LimitNode.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/LimitNode.java
@@ -45,6 +45,16 @@ public class LimitNode extends ProcessNode {
   }
 
   @Override
+  public PlanNode clone() {
+    return null;
+  }
+
+  @Override
+  public PlanNode cloneWithChildren(List<PlanNode> children) {
+    return null;
+  }
+
+  @Override
   public List<String> getOutputColumnNames() {
     return child.getOutputColumnNames();
   }
diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/OffsetNode.java
 
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/OffsetNode.java
index 2e3fc78..bde203c 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/OffsetNode.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/OffsetNode.java
@@ -46,6 +46,16 @@ public class OffsetNode extends ProcessNode {
   }
 
   @Override
+  public PlanNode clone() {
+    return null;
+  }
+
+  @Override
+  public PlanNode cloneWithChildren(List<PlanNode> children) {
+    return null;
+  }
+
+  @Override
   public List<String> getOutputColumnNames() {
     return null;
   }
diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/SortNode.java
 
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/SortNode.java
index 19464a2..b0752f0 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/SortNode.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/SortNode.java
@@ -52,6 +52,16 @@ public class SortNode extends ProcessNode {
   }
 
   @Override
+  public PlanNode clone() {
+    return null;
+  }
+
+  @Override
+  public PlanNode cloneWithChildren(List<PlanNode> children) {
+    return null;
+  }
+
+  @Override
   public List<String> getOutputColumnNames() {
     return child.getOutputColumnNames();
   }
diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/TimeJoinNode.java
 
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/TimeJoinNode.java
index 74fac14..4886055 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/TimeJoinNode.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/TimeJoinNode.java
@@ -63,6 +63,16 @@ public class TimeJoinNode extends ProcessNode {
   }
 
   @Override
+  public PlanNode clone() {
+    return null;
+  }
+
+  @Override
+  public PlanNode cloneWithChildren(List<PlanNode> children) {
+    return null;
+  }
+
+  @Override
   public List<String> getOutputColumnNames() {
     return children.stream()
         .flatMap(child -> child.getOutputColumnNames().stream())
diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/sink/CsvSinkNode.java
 
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/sink/CsvSinkNode.java
index 2fddbdd..a379c33 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/sink/CsvSinkNode.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/sink/CsvSinkNode.java
@@ -34,6 +34,16 @@ public class CsvSinkNode extends SinkNode {
   }
 
   @Override
+  public PlanNode clone() {
+    return null;
+  }
+
+  @Override
+  public PlanNode cloneWithChildren(List<PlanNode> children) {
+    return null;
+  }
+
+  @Override
   public List<String> getOutputColumnNames() {
     return null;
   }
diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/sink/FragmentSinkNode.java
 
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/sink/FragmentSinkNode.java
index da30223..6151916 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/sink/FragmentSinkNode.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/sink/FragmentSinkNode.java
@@ -34,6 +34,16 @@ public class FragmentSinkNode extends SinkNode {
   }
 
   @Override
+  public PlanNode clone() {
+    return null;
+  }
+
+  @Override
+  public PlanNode cloneWithChildren(List<PlanNode> children) {
+    return null;
+  }
+
+  @Override
   public List<String> getOutputColumnNames() {
     return null;
   }
diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/sink/ThriftSinkNode.java
 
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/sink/ThriftSinkNode.java
index bb343f4..4781e7e 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/sink/ThriftSinkNode.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/sink/ThriftSinkNode.java
@@ -36,6 +36,16 @@ public class ThriftSinkNode extends SinkNode {
   }
 
   @Override
+  public PlanNode clone() {
+    return null;
+  }
+
+  @Override
+  public PlanNode cloneWithChildren(List<PlanNode> children) {
+    return null;
+  }
+
+  @Override
   public List<String> getOutputColumnNames() {
     return null;
   }
diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/source/CsvSourceNode.java
 
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/source/CsvSourceNode.java
index 612e930..797611e 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/source/CsvSourceNode.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/source/CsvSourceNode.java
@@ -36,6 +36,16 @@ public class CsvSourceNode extends SourceNode {
   }
 
   @Override
+  public PlanNode clone() {
+    return null;
+  }
+
+  @Override
+  public PlanNode cloneWithChildren(List<PlanNode> children) {
+    return null;
+  }
+
+  @Override
   public List<String> getOutputColumnNames() {
     return null;
   }
diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/source/SeriesAggregateScanNode.java
 
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/source/SeriesAggregateScanNode.java
index ee01ac1..d071fd9 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/source/SeriesAggregateScanNode.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/source/SeriesAggregateScanNode.java
@@ -69,6 +69,16 @@ public class SeriesAggregateScanNode extends SourceNode {
   }
 
   @Override
+  public PlanNode clone() {
+    return null;
+  }
+
+  @Override
+  public PlanNode cloneWithChildren(List<PlanNode> children) {
+    return null;
+  }
+
+  @Override
   public List<String> getOutputColumnNames() {
     return ImmutableList.of(columnName);
   }
diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/source/SeriesScanNode.java
 
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/source/SeriesScanNode.java
index 8858074..939e2d4 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/source/SeriesScanNode.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/source/SeriesScanNode.java
@@ -19,6 +19,7 @@
 package org.apache.iotdb.db.mpp.sql.planner.plan.node.source;
 
 import org.apache.iotdb.db.metadata.path.PartialPath;
+import org.apache.iotdb.db.mpp.common.DataRegion;
 import org.apache.iotdb.db.mpp.common.OrderBy;
 import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanNode;
 import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanNodeId;
@@ -60,6 +61,9 @@ public class SeriesScanNode extends SourceNode {
 
   private String columnName;
 
+  // The id of DataRegion where the node will run
+  private DataRegion dataRegion;
+
   public SeriesScanNode(PlanNodeId id, PartialPath seriesPath) {
     super(id);
     this.seriesPath = seriesPath;
@@ -97,6 +101,16 @@ public class SeriesScanNode extends SourceNode {
   }
 
   @Override
+  public PlanNode clone() {
+    return null;
+  }
+
+  @Override
+  public PlanNode cloneWithChildren(List<PlanNode> children) {
+    return null;
+  }
+
+  @Override
   public List<String> getOutputColumnNames() {
     return ImmutableList.of(columnName);
   }
@@ -105,4 +119,20 @@ public class SeriesScanNode extends SourceNode {
   public <R, C> R accept(PlanVisitor<R, C> visitor, C context) {
     return visitor.visitSeriesScan(this, context);
   }
+
+  public PartialPath getSeriesPath() {
+    return seriesPath;
+  }
+
+  public Filter getTimeFilter() {
+    return timeFilter;
+  }
+
+  public void setDataRegion(DataRegion dataRegion) {
+    this.dataRegion = dataRegion;
+  }
+
+  public DataRegion getDataRegion() {
+    return dataRegion;
+  }
 }

Reply via email to