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