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

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


The following commit(s) were added to 
refs/heads/xingtanzjr/align_by_device_distribution by this push:
     new 1bb7ff14c4 fix aligned reader
1bb7ff14c4 is described below

commit 1bb7ff14c4c450258ed6aca571f8f6cc408ecd67
Author: JackieTien97 <[email protected]>
AuthorDate: Fri May 27 21:36:17 2022 +0800

    fix aligned reader
---
 .../iotdb/db/mpp/plan/planner/LocalExecutionPlanner.java | 16 ++++++----------
 .../mpp/plan/planner/distribution/ExchangeNodeAdder.java |  7 +++++--
 .../db/mpp/plan/planner/distribution/SourceRewriter.java |  5 +++--
 .../mpp/plan/plan/distribution/AlignedByDeviceTest.java  |  5 ++---
 .../iotdb/tsfile/read/reader/page/AlignedPageReader.java | 10 ++++------
 5 files changed, 20 insertions(+), 23 deletions(-)

diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/LocalExecutionPlanner.java
 
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/LocalExecutionPlanner.java
index 1b7b78b44e..e2a3176515 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/LocalExecutionPlanner.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/LocalExecutionPlanner.java
@@ -231,9 +231,7 @@ public class LocalExecutionPlanner {
     return new SchemaDriver(root, context.getSinkHandle(), 
schemaDriverContext);
   }
 
-  /**
-   * This Visitor is responsible for transferring PlanNode Tree to Operator 
Tree
-   */
+  /** This Visitor is responsible for transferring PlanNode Tree to Operator 
Tree */
   private static class Visitor extends PlanVisitor<Operator, 
LocalExecutionPlanContext> {
 
     @Override
@@ -325,7 +323,7 @@ public class LocalExecutionPlanner {
                     descriptor.getAggregationType(), seriesDataType, 
ascending),
                 descriptor.getStep(),
                 Collections.singletonList(
-                    new InputLocation[]{new InputLocation(0, seriesIndex)})));
+                    new InputLocation[] {new InputLocation(0, seriesIndex)})));
       }
 
       AlignedSeriesAggregationScanOperator seriesAggregationScanOperator =
@@ -826,7 +824,6 @@ public class LocalExecutionPlanner {
           operatorContext, aggregators, children, ascending, 
node.getGroupByTimeParameter());
     }
 
-
     @Override
     public Operator visitSlidingWindowAggregation(
         SlidingWindowAggregationNode node, LocalExecutionPlanContext context) {
@@ -945,11 +942,11 @@ public class LocalExecutionPlanner {
       List<InputLocation[]> inputLocationList = new ArrayList<>();
       for (int i = 0; i < inputLocationParts.get(0).size(); i++) {
         if (inputColumnNames.size() == 1) {
-          inputLocationList.add(new 
InputLocation[]{inputLocationParts.get(0).get(i)});
+          inputLocationList.add(new InputLocation[] 
{inputLocationParts.get(0).get(i)});
         } else {
           inputLocationList.add(
-              new InputLocation[]{
-                  inputLocationParts.get(0).get(i), 
inputLocationParts.get(1).get(i)
+              new InputLocation[] {
+                inputLocationParts.get(0).get(i), 
inputLocationParts.get(1).get(i)
               });
         }
       }
@@ -1309,8 +1306,7 @@ public class LocalExecutionPlanner {
 
   private static class InstanceHolder {
 
-    private InstanceHolder() {
-    }
+    private InstanceHolder() {}
 
     private static final LocalExecutionPlanner INSTANCE = new 
LocalExecutionPlanner();
   }
diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/distribution/ExchangeNodeAdder.java
 
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/distribution/ExchangeNodeAdder.java
index fdeaab5e7a..f85931231b 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/distribution/ExchangeNodeAdder.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/distribution/ExchangeNodeAdder.java
@@ -276,8 +276,11 @@ public class ExchangeNodeAdder extends 
PlanVisitor<PlanNode, NodeGroupContext> {
             .collect(
                 Collectors.groupingBy(
                     child -> {
-                      TRegionReplicaSet region = 
context.getNodeDistribution(child.getPlanNodeId()).region;
-                      if (region == null && 
context.getNodeDistribution(child.getPlanNodeId()).type == 
NodeDistributionType.SAME_WITH_ALL_CHILDREN) {
+                      TRegionReplicaSet region =
+                          
context.getNodeDistribution(child.getPlanNodeId()).region;
+                      if (region == null
+                          && 
context.getNodeDistribution(child.getPlanNodeId()).type
+                              == NodeDistributionType.SAME_WITH_ALL_CHILDREN) {
                         return 
calculateSchemaRegionByChildren(child.getChildren(), context);
                       }
                       return region;
diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/distribution/SourceRewriter.java
 
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/distribution/SourceRewriter.java
index 64a348ad76..903e60d481 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/distribution/SourceRewriter.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/distribution/SourceRewriter.java
@@ -28,7 +28,6 @@ import 
org.apache.iotdb.db.mpp.common.schematree.PathPatternTree;
 import org.apache.iotdb.db.mpp.plan.analyze.Analysis;
 import org.apache.iotdb.db.mpp.plan.expression.Expression;
 import org.apache.iotdb.db.mpp.plan.planner.plan.node.PlanNode;
-import org.apache.iotdb.db.mpp.plan.planner.plan.node.PlanNodeUtil;
 import org.apache.iotdb.db.mpp.plan.planner.plan.node.SimplePlanNodeRewriter;
 import 
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.read.CountSchemaMergeNode;
 import 
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.read.SchemaFetchMergeNode;
@@ -153,7 +152,9 @@ public class SourceRewriter extends 
SimplePlanNodeRewriter<DistributionPlanConte
     private PlanNode buildPlanNodeInRegion(
         PlanNode root, TRegionReplicaSet regionReplicaSet, MPPQueryContext 
context) {
       List<PlanNode> children =
-          root.getChildren().stream().map(child -> 
buildPlanNodeInRegion(child, regionReplicaSet, 
context)).collect(Collectors.toList());
+          root.getChildren().stream()
+              .map(child -> buildPlanNodeInRegion(child, regionReplicaSet, 
context))
+              .collect(Collectors.toList());
       PlanNode newRoot = root.cloneWithChildren(children);
       newRoot.setPlanNodeId(context.getQueryId().genPlanNodeId());
       if (newRoot instanceof SourceNode) {
diff --git 
a/server/src/test/java/org/apache/iotdb/db/mpp/plan/plan/distribution/AlignedByDeviceTest.java
 
b/server/src/test/java/org/apache/iotdb/db/mpp/plan/plan/distribution/AlignedByDeviceTest.java
index f10817ab8a..af64fc39bb 100644
--- 
a/server/src/test/java/org/apache/iotdb/db/mpp/plan/plan/distribution/AlignedByDeviceTest.java
+++ 
b/server/src/test/java/org/apache/iotdb/db/mpp/plan/plan/distribution/AlignedByDeviceTest.java
@@ -20,6 +20,7 @@
 package org.apache.iotdb.db.mpp.plan.plan.distribution;
 
 import org.apache.iotdb.db.mpp.plan.planner.plan.LogicalQueryPlan;
+
 import org.junit.Test;
 
 import java.util.List;
@@ -27,9 +28,7 @@ import java.util.List;
 public class AlignedByDeviceTest {
 
   @Test
-  public void test1Device1Region() {
-
-  }
+  public void test1Device1Region() {}
 
   private LogicalQueryPlan constructLogicalPlan(List<String> series) {
     return null;
diff --git 
a/tsfile/src/main/java/org/apache/iotdb/tsfile/read/reader/page/AlignedPageReader.java
 
b/tsfile/src/main/java/org/apache/iotdb/tsfile/read/reader/page/AlignedPageReader.java
index 5dc9a466a6..e7578ed630 100644
--- 
a/tsfile/src/main/java/org/apache/iotdb/tsfile/read/reader/page/AlignedPageReader.java
+++ 
b/tsfile/src/main/java/org/apache/iotdb/tsfile/read/reader/page/AlignedPageReader.java
@@ -37,7 +37,6 @@ import java.io.IOException;
 import java.nio.ByteBuffer;
 import java.util.ArrayList;
 import java.util.List;
-import java.util.stream.Collectors;
 
 public class AlignedPageReader implements IPageReader, IAlignedPageReader {
 
@@ -47,6 +46,8 @@ public class AlignedPageReader implements IPageReader, 
IAlignedPageReader {
   private Filter filter;
   private boolean isModified;
 
+  private final TsBlockBuilder builder;
+
   public AlignedPageReader(
       PageHeader timePageHeader,
       ByteBuffer timePageData,
@@ -75,6 +76,7 @@ public class AlignedPageReader implements IPageReader, 
IAlignedPageReader {
     }
     this.filter = filter;
     this.valueCount = valuePageReaderList.size();
+    this.builder = new TsBlockBuilder(valueDataTypeList);
   }
 
   @Override
@@ -108,11 +110,7 @@ public class AlignedPageReader implements IPageReader, 
IAlignedPageReader {
   @Override
   public TsBlock getAllSatisfiedData() throws IOException {
     // TODO change from the row-based style to column-based style
-    TsBlockBuilder builder =
-        new TsBlockBuilder(
-            valuePageReaderList.stream()
-                .map(ValuePageReader::getDataType)
-                .collect(Collectors.toList()));
+    builder.reset();
     int timeIndex = -1;
     while (timePageReader.hasNextTime()) {
       long timestamp = timePageReader.nextTime();

Reply via email to