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

jackietien pushed a commit to branch ty-mpp
in repository https://gitbox.apache.org/repos/asf/iotdb.git

commit 10b8ec3d6ee1b1f95298ee2ed752c41739ae1c18
Author: JackieTien97 <[email protected]>
AuthorDate: Sat Apr 9 11:58:36 2022 +0800

    Construct ExchangeOperator in LocalExecutionPlanner and move 
initQueryDataSource method from SourceOperator to DataSourceOperator
---
 .../statemachine/DataRegionStateMachine.java       |  7 ++-
 .../org/apache/iotdb/db/mpp/common/DataRegion.java |  4 +-
 .../apache/iotdb/db/mpp/execution/DataDriver.java  |  4 +-
 .../iotdb/db/mpp/execution/DataDriverContext.java  |  8 +--
 .../db/mpp/execution/FragmentInstanceInfo.java     |  4 +-
 ...SourceOperator.java => DataSourceOperator.java} |  8 +--
 ...gateScanOperator.java => ExchangeOperator.java} | 61 ++++++++++++++++------
 .../source/SeriesAggregateScanOperator.java        |  6 +--
 .../db/mpp/operator/source/SeriesScanOperator.java |  7 ++-
 .../db/mpp/operator/source/SourceOperator.java     |  3 --
 .../db/mpp/sql/planner/LocalExecutionPlanner.java  | 41 +++++++++------
 .../iotdb/db/mpp/execution/DataDriverTest.java     |  8 ++-
 .../iotdb/db/mpp/operator/LimitOperatorTest.java   |  8 ++-
 .../db/mpp/operator/SeriesScanOperatorTest.java    |  4 +-
 .../db/mpp/operator/TimeJoinOperatorTest.java      |  8 ++-
 15 files changed, 118 insertions(+), 63 deletions(-)

diff --git 
a/server/src/main/java/org/apache/iotdb/db/consensus/statemachine/DataRegionStateMachine.java
 
b/server/src/main/java/org/apache/iotdb/db/consensus/statemachine/DataRegionStateMachine.java
index 06312cd13f..883bf457fa 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/consensus/statemachine/DataRegionStateMachine.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/consensus/statemachine/DataRegionStateMachine.java
@@ -21,6 +21,7 @@ package org.apache.iotdb.db.consensus.statemachine;
 
 import org.apache.iotdb.consensus.common.DataSet;
 import org.apache.iotdb.db.mpp.common.DataRegion;
+import org.apache.iotdb.db.mpp.execution.FragmentInstanceManager;
 import org.apache.iotdb.db.mpp.sql.planner.plan.FragmentInstance;
 import org.apache.iotdb.rpc.TSStatusCode;
 import org.apache.iotdb.service.rpc.thrift.TSStatus;
@@ -32,6 +33,9 @@ public class DataRegionStateMachine extends BaseStateMachine {
 
   private static final Logger logger = 
LoggerFactory.getLogger(DataRegionStateMachine.class);
 
+  private static final FragmentInstanceManager INSTANCE_MANAGER =
+      FragmentInstanceManager.getInstance();
+
   private final DataRegion region;
 
   public DataRegionStateMachine(DataRegion region) {
@@ -52,7 +56,8 @@ public class DataRegionStateMachine extends BaseStateMachine {
 
   @Override
   protected DataSet read(FragmentInstance fragmentInstance) {
-    logger.info("Execute read plan in DataRegionStateMachine");
+    // TODO here we need to change after VSG being renamed to DataRegion
+    //    return 
INSTANCE_MANAGER.execDataQueryFragmentInstance(fragmentInstance, region);
     return null;
   }
 }
diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/common/DataRegion.java 
b/server/src/main/java/org/apache/iotdb/db/mpp/common/DataRegion.java
index 81c9665ba3..d9454e5271 100644
--- a/server/src/main/java/org/apache/iotdb/db/mpp/common/DataRegion.java
+++ b/server/src/main/java/org/apache/iotdb/db/mpp/common/DataRegion.java
@@ -25,8 +25,8 @@ package org.apache.iotdb.db.mpp.common;
  */
 // TODO: (xingtanzjr) This class should be substituted with the class defined 
in Consensus level
 public class DataRegion {
-  private Integer dataRegionId;
-  private String endpoint;
+  private final Integer dataRegionId;
+  private final String endpoint;
 
   public DataRegion(Integer dataRegionId, String endpoint) {
     this.dataRegionId = dataRegionId;
diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/execution/DataDriver.java 
b/server/src/main/java/org/apache/iotdb/db/mpp/execution/DataDriver.java
index 5a1168fac3..9996381e4f 100644
--- a/server/src/main/java/org/apache/iotdb/db/mpp/execution/DataDriver.java
+++ b/server/src/main/java/org/apache/iotdb/db/mpp/execution/DataDriver.java
@@ -27,7 +27,7 @@ import org.apache.iotdb.db.metadata.path.PartialPath;
 import org.apache.iotdb.db.mpp.buffer.ISinkHandle;
 import org.apache.iotdb.db.mpp.common.FragmentInstanceId;
 import org.apache.iotdb.db.mpp.operator.Operator;
-import org.apache.iotdb.db.mpp.operator.source.SourceOperator;
+import org.apache.iotdb.db.mpp.operator.source.DataSourceOperator;
 import org.apache.iotdb.db.query.control.FileReaderManager;
 import org.apache.iotdb.tsfile.read.common.block.TsBlock;
 
@@ -174,7 +174,7 @@ public class DataDriver implements Driver {
    * we should change all the blocked lock operation into tryLock
    */
   private void initialize() throws QueryProcessException {
-    List<SourceOperator> sourceOperators = driverContext.getSourceOperators();
+    List<DataSourceOperator> sourceOperators = 
driverContext.getSourceOperators();
     if (sourceOperators != null && !sourceOperators.isEmpty()) {
       QueryDataSource dataSource = initQueryDataSourceCache();
       sourceOperators.forEach(
diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/execution/DataDriverContext.java 
b/server/src/main/java/org/apache/iotdb/db/mpp/execution/DataDriverContext.java
index 52113e5586..c1bc46f5bd 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/execution/DataDriverContext.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/execution/DataDriverContext.java
@@ -20,7 +20,7 @@ package org.apache.iotdb.db.mpp.execution;
 
 import org.apache.iotdb.db.engine.storagegroup.VirtualStorageGroupProcessor;
 import org.apache.iotdb.db.metadata.path.PartialPath;
-import org.apache.iotdb.db.mpp.operator.source.SourceOperator;
+import org.apache.iotdb.db.mpp.operator.source.DataSourceOperator;
 import org.apache.iotdb.tsfile.read.filter.basic.Filter;
 
 import java.util.List;
@@ -29,14 +29,14 @@ public class DataDriverContext extends DriverContext {
   private final List<PartialPath> paths;
   private final Filter timeFilter;
   private final VirtualStorageGroupProcessor dataRegion;
-  private final List<SourceOperator> sourceOperators;
+  private final List<DataSourceOperator> sourceOperators;
 
   public DataDriverContext(
       FragmentInstanceContext fragmentInstanceContext,
       List<PartialPath> paths,
       Filter timeFilter,
       VirtualStorageGroupProcessor dataRegion,
-      List<SourceOperator> sourceOperators) {
+      List<DataSourceOperator> sourceOperators) {
     super(fragmentInstanceContext);
     this.paths = paths;
     this.timeFilter = timeFilter;
@@ -56,7 +56,7 @@ public class DataDriverContext extends DriverContext {
     return dataRegion;
   }
 
-  public List<SourceOperator> getSourceOperators() {
+  public List<DataSourceOperator> getSourceOperators() {
     return sourceOperators;
   }
 }
diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/execution/FragmentInstanceInfo.java
 
b/server/src/main/java/org/apache/iotdb/db/mpp/execution/FragmentInstanceInfo.java
index 4f02f8465d..d93f276045 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/execution/FragmentInstanceInfo.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/execution/FragmentInstanceInfo.java
@@ -18,7 +18,9 @@
  */
 package org.apache.iotdb.db.mpp.execution;
 
-public class FragmentInstanceInfo {
+import org.apache.iotdb.consensus.common.DataSet;
+
+public class FragmentInstanceInfo implements DataSet {
 
   private final FragmentInstanceState state;
 
diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/operator/source/SourceOperator.java
 
b/server/src/main/java/org/apache/iotdb/db/mpp/operator/source/DataSourceOperator.java
similarity index 79%
copy from 
server/src/main/java/org/apache/iotdb/db/mpp/operator/source/SourceOperator.java
copy to 
server/src/main/java/org/apache/iotdb/db/mpp/operator/source/DataSourceOperator.java
index 2e81002f7a..7bc2c46733 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/operator/source/SourceOperator.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/operator/source/DataSourceOperator.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
@@ -19,12 +19,8 @@
 package org.apache.iotdb.db.mpp.operator.source;
 
 import org.apache.iotdb.db.engine.querycontext.QueryDataSource;
-import org.apache.iotdb.db.mpp.operator.Operator;
-import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanNodeId;
 
-public interface SourceOperator extends Operator {
-
-  PlanNodeId getSourceId();
+public interface DataSourceOperator extends SourceOperator {
 
   void initQueryDataSource(QueryDataSource dataSource);
 }
diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/operator/source/SeriesAggregateScanOperator.java
 
b/server/src/main/java/org/apache/iotdb/db/mpp/operator/source/ExchangeOperator.java
similarity index 52%
copy from 
server/src/main/java/org/apache/iotdb/db/mpp/operator/source/SeriesAggregateScanOperator.java
copy to 
server/src/main/java/org/apache/iotdb/db/mpp/operator/source/ExchangeOperator.java
index 07aa5a4c2f..7bdfc8e745 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/operator/source/SeriesAggregateScanOperator.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/operator/source/ExchangeOperator.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
@@ -18,49 +18,76 @@
  */
 package org.apache.iotdb.db.mpp.operator.source;
 
-import org.apache.iotdb.db.engine.querycontext.QueryDataSource;
+import org.apache.iotdb.db.mpp.buffer.ISourceHandle;
 import org.apache.iotdb.db.mpp.operator.OperatorContext;
 import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanNodeId;
 import org.apache.iotdb.tsfile.read.common.block.TsBlock;
 
 import com.google.common.util.concurrent.ListenableFuture;
 
-public class SeriesAggregateScanOperator implements SourceOperator {
-  @Override
-  public OperatorContext getOperatorContext() {
-    return null;
+import java.io.IOException;
+
+public class ExchangeOperator implements SourceOperator {
+
+  private final OperatorContext operatorContext;
+
+  private final ISourceHandle sourceHandle;
+
+  private final PlanNodeId sourceId;
+
+  private ListenableFuture<Void> isBlocked = NOT_BLOCKED;
+
+  public ExchangeOperator(
+      OperatorContext operatorContext, ISourceHandle sourceHandle, PlanNodeId 
sourceId) {
+    this.operatorContext = operatorContext;
+    this.sourceHandle = sourceHandle;
+    this.sourceId = sourceId;
   }
 
   @Override
-  public ListenableFuture<Void> isBlocked() {
-    return SourceOperator.super.isBlocked();
+  public OperatorContext getOperatorContext() {
+    return operatorContext;
   }
 
   @Override
   public TsBlock next() {
-    return null;
+    try {
+      return sourceHandle.receive();
+    } catch (IOException e) {
+      throw new RuntimeException(
+          "Error happened while reading from source handle " + sourceHandle, 
e);
+    }
   }
 
   @Override
   public boolean hasNext() {
-    return false;
+    return sourceHandle.isFinished();
   }
 
   @Override
-  public void close() throws Exception {
-    SourceOperator.super.close();
+  public boolean isFinished() throws IOException {
+    return sourceHandle.isFinished();
   }
 
   @Override
-  public boolean isFinished() {
-    return false;
+  public PlanNodeId getSourceId() {
+    return sourceId;
   }
 
   @Override
-  public PlanNodeId getSourceId() {
-    return null;
+  public ListenableFuture<Void> isBlocked() {
+    // Avoid registering a new callback in the source handle when one is 
already pending
+    if (isBlocked.isDone()) {
+      isBlocked = sourceHandle.isBlocked();
+      if (isBlocked.isDone()) {
+        isBlocked = NOT_BLOCKED;
+      }
+    }
+    return isBlocked;
   }
 
   @Override
-  public void initQueryDataSource(QueryDataSource dataSource) {}
+  public void close() throws Exception {
+    sourceHandle.close();
+  }
 }
diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/operator/source/SeriesAggregateScanOperator.java
 
b/server/src/main/java/org/apache/iotdb/db/mpp/operator/source/SeriesAggregateScanOperator.java
index 07aa5a4c2f..c4b1015529 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/operator/source/SeriesAggregateScanOperator.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/operator/source/SeriesAggregateScanOperator.java
@@ -25,7 +25,7 @@ import org.apache.iotdb.tsfile.read.common.block.TsBlock;
 
 import com.google.common.util.concurrent.ListenableFuture;
 
-public class SeriesAggregateScanOperator implements SourceOperator {
+public class SeriesAggregateScanOperator implements DataSourceOperator {
   @Override
   public OperatorContext getOperatorContext() {
     return null;
@@ -33,7 +33,7 @@ public class SeriesAggregateScanOperator implements 
SourceOperator {
 
   @Override
   public ListenableFuture<Void> isBlocked() {
-    return SourceOperator.super.isBlocked();
+    return DataSourceOperator.super.isBlocked();
   }
 
   @Override
@@ -48,7 +48,7 @@ public class SeriesAggregateScanOperator implements 
SourceOperator {
 
   @Override
   public void close() throws Exception {
-    SourceOperator.super.close();
+    DataSourceOperator.super.close();
   }
 
   @Override
diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/operator/source/SeriesScanOperator.java
 
b/server/src/main/java/org/apache/iotdb/db/mpp/operator/source/SeriesScanOperator.java
index e979710c3f..8e58e36b61 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/operator/source/SeriesScanOperator.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/operator/source/SeriesScanOperator.java
@@ -29,15 +29,17 @@ import org.apache.iotdb.tsfile.read.filter.basic.Filter;
 import java.io.IOException;
 import java.util.Set;
 
-public class SeriesScanOperator implements SourceOperator {
+public class SeriesScanOperator implements DataSourceOperator {
 
   private final OperatorContext operatorContext;
   private final SeriesScanUtil seriesScanUtil;
+  private final PlanNodeId sourceId;
   private TsBlock tsBlock;
   private boolean hasCachedTsBlock = false;
   private boolean finished = false;
 
   public SeriesScanOperator(
+      PlanNodeId sourceId,
       PartialPath seriesPath,
       Set<String> allSensors,
       TSDataType dataType,
@@ -45,6 +47,7 @@ public class SeriesScanOperator implements SourceOperator {
       Filter timeFilter,
       Filter valueFilter,
       boolean ascending) {
+    this.sourceId = sourceId;
     this.operatorContext = context;
     this.seriesScanUtil =
         new SeriesScanUtil(
@@ -140,7 +143,7 @@ public class SeriesScanOperator implements SourceOperator {
 
   @Override
   public PlanNodeId getSourceId() {
-    return null;
+    return sourceId;
   }
 
   @Override
diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/operator/source/SourceOperator.java
 
b/server/src/main/java/org/apache/iotdb/db/mpp/operator/source/SourceOperator.java
index 2e81002f7a..8454fd685b 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/operator/source/SourceOperator.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/operator/source/SourceOperator.java
@@ -18,13 +18,10 @@
  */
 package org.apache.iotdb.db.mpp.operator.source;
 
-import org.apache.iotdb.db.engine.querycontext.QueryDataSource;
 import org.apache.iotdb.db.mpp.operator.Operator;
 import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanNodeId;
 
 public interface SourceOperator extends Operator {
 
   PlanNodeId getSourceId();
-
-  void initQueryDataSource(QueryDataSource dataSource);
 }
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 1af476b7a2..b9c514795f 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
@@ -25,6 +25,7 @@ import org.apache.iotdb.db.metadata.schemaregion.SchemaRegion;
 import org.apache.iotdb.db.mpp.buffer.DataBlockManager;
 import org.apache.iotdb.db.mpp.buffer.DataBlockService;
 import org.apache.iotdb.db.mpp.buffer.ISinkHandle;
+import org.apache.iotdb.db.mpp.buffer.ISourceHandle;
 import org.apache.iotdb.db.mpp.common.FragmentInstanceId;
 import org.apache.iotdb.db.mpp.common.filter.QueryFilter;
 import org.apache.iotdb.db.mpp.execution.DataDriver;
@@ -36,8 +37,9 @@ import org.apache.iotdb.db.mpp.operator.Operator;
 import org.apache.iotdb.db.mpp.operator.OperatorContext;
 import org.apache.iotdb.db.mpp.operator.process.LimitOperator;
 import org.apache.iotdb.db.mpp.operator.process.TimeJoinOperator;
+import org.apache.iotdb.db.mpp.operator.source.DataSourceOperator;
+import org.apache.iotdb.db.mpp.operator.source.ExchangeOperator;
 import org.apache.iotdb.db.mpp.operator.source.SeriesScanOperator;
-import org.apache.iotdb.db.mpp.operator.source.SourceOperator;
 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.AggregateNode;
@@ -55,7 +57,6 @@ import 
org.apache.iotdb.db.mpp.sql.planner.plan.node.sink.FragmentSinkNode;
 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 org.apache.iotdb.db.mpp.sql.statement.component.OrderBy;
-import org.apache.iotdb.mpp.rpc.thrift.TFragmentInstanceId;
 import org.apache.iotdb.tsfile.read.filter.basic.Filter;
 
 import java.io.IOException;
@@ -132,6 +133,7 @@ public class LocalExecutionPlanner {
 
       SeriesScanOperator seriesScanOperator =
           new SeriesScanOperator(
+              node.getId(),
               seriesPath,
               node.getAllSensors(),
               seriesPath.getSeriesType(),
@@ -221,9 +223,24 @@ public class LocalExecutionPlanner {
 
     @Override
     public Operator visitExchange(ExchangeNode node, LocalExecutionPlanContext 
context) {
+      OperatorContext operatorContext =
+          context.instanceContext.addOperatorContext(
+              context.getNextOperatorId(), node.getId(), 
SeriesScanOperator.class.getSimpleName());
+      FragmentInstanceId localInstanceId = context.instanceContext.getId();
+      FragmentInstanceId remoteInstanceId = node.getUpstreamInstanceId();
+      Endpoint source = node.getUpstreamEndpoint();
 
-      // TODO(jackie tien) create SourceHandle here
-      return super.visitExchange(node, context);
+      try {
+        ISourceHandle sourceHandle =
+            DATA_BLOCK_MANAGER.createSourceHandle(
+                localInstanceId.toThrift(),
+                node.getId().getId(),
+                source.getIp(),
+                remoteInstanceId.toThrift());
+        return new ExchangeOperator(operatorContext, sourceHandle, 
node.getUpstreamPlanNodeId());
+      } catch (IOException e) {
+        throw new RuntimeException("Error happened while creating source 
handle", e);
+      }
     }
 
     @Override
@@ -235,15 +252,9 @@ public class LocalExecutionPlanner {
       try {
         ISinkHandle sinkHandle =
             DATA_BLOCK_MANAGER.createSinkHandle(
-                new TFragmentInstanceId(
-                    localInstanceId.getQueryId().getId(),
-                    String.valueOf(localInstanceId.getFragmentId().getId()),
-                    localInstanceId.getInstanceId()),
+                localInstanceId.toThrift(),
                 target.getIp(),
-                new TFragmentInstanceId(
-                    targetInstanceId.getQueryId().getId(),
-                    String.valueOf(targetInstanceId.getFragmentId().getId()),
-                    targetInstanceId.getInstanceId()),
+                targetInstanceId.toThrift(),
                 node.getDownStreamPlanNodeId().getId());
         context.setSinkHandle(sinkHandle);
         return child;
@@ -263,7 +274,7 @@ public class LocalExecutionPlanner {
   private static class LocalExecutionPlanContext {
     private final FragmentInstanceContext instanceContext;
     private final List<PartialPath> paths;
-    private final List<SourceOperator> sourceOperators;
+    private final List<DataSourceOperator> sourceOperators;
     private ISinkHandle sinkHandle;
 
     private int nextOperatorId = 0;
@@ -282,7 +293,7 @@ public class LocalExecutionPlanner {
       return paths;
     }
 
-    public List<SourceOperator> getSourceOperators() {
+    public List<DataSourceOperator> getSourceOperators() {
       return sourceOperators;
     }
 
@@ -290,7 +301,7 @@ public class LocalExecutionPlanner {
       paths.add(path);
     }
 
-    public void addSourceOperator(SourceOperator sourceOperator) {
+    public void addSourceOperator(DataSourceOperator sourceOperator) {
       sourceOperators.add(sourceOperator);
     }
 
diff --git 
a/server/src/test/java/org/apache/iotdb/db/mpp/execution/DataDriverTest.java 
b/server/src/test/java/org/apache/iotdb/db/mpp/execution/DataDriverTest.java
index 9a5c16c7a6..70d3236850 100644
--- a/server/src/test/java/org/apache/iotdb/db/mpp/execution/DataDriverTest.java
+++ b/server/src/test/java/org/apache/iotdb/db/mpp/execution/DataDriverTest.java
@@ -97,16 +97,19 @@ public class DataDriverTest {
       FragmentInstanceContext fragmentInstanceContext =
           new FragmentInstanceContext(
               new FragmentInstanceId(new PlanFragmentId(queryId, 0), 
"stub-instance"), state);
+      PlanNodeId planNodeId1 = new PlanNodeId("1");
       fragmentInstanceContext.addOperatorContext(
-          1, new PlanNodeId("1"), SeriesScanOperator.class.getSimpleName());
+          1, planNodeId1, SeriesScanOperator.class.getSimpleName());
+      PlanNodeId planNodeId2 = new PlanNodeId("2");
       fragmentInstanceContext.addOperatorContext(
-          2, new PlanNodeId("2"), SeriesScanOperator.class.getSimpleName());
+          2, planNodeId2, SeriesScanOperator.class.getSimpleName());
       fragmentInstanceContext.addOperatorContext(
           3, new PlanNodeId("3"), TimeJoinOperator.class.getSimpleName());
       fragmentInstanceContext.addOperatorContext(
           4, new PlanNodeId("4"), LimitOperator.class.getSimpleName());
       SeriesScanOperator seriesScanOperator1 =
           new SeriesScanOperator(
+              planNodeId1,
               measurementPath1,
               allSensors,
               TSDataType.INT32,
@@ -119,6 +122,7 @@ public class DataDriverTest {
           new MeasurementPath(DATA_DRIVER_TEST_SG + ".device0.sensor1", 
TSDataType.INT32);
       SeriesScanOperator seriesScanOperator2 =
           new SeriesScanOperator(
+              planNodeId2,
               measurementPath2,
               allSensors,
               TSDataType.INT32,
diff --git 
a/server/src/test/java/org/apache/iotdb/db/mpp/operator/LimitOperatorTest.java 
b/server/src/test/java/org/apache/iotdb/db/mpp/operator/LimitOperatorTest.java
index f0c8e7e5c1..6ad74bd110 100644
--- 
a/server/src/test/java/org/apache/iotdb/db/mpp/operator/LimitOperatorTest.java
+++ 
b/server/src/test/java/org/apache/iotdb/db/mpp/operator/LimitOperatorTest.java
@@ -90,16 +90,19 @@ public class LimitOperatorTest {
       FragmentInstanceContext fragmentInstanceContext =
           new FragmentInstanceContext(
               new FragmentInstanceId(new PlanFragmentId(queryId, 0), 
"stub-instance"), state);
+      PlanNodeId planNodeId1 = new PlanNodeId("1");
       fragmentInstanceContext.addOperatorContext(
-          1, new PlanNodeId("1"), SeriesScanOperator.class.getSimpleName());
+          1, planNodeId1, SeriesScanOperator.class.getSimpleName());
+      PlanNodeId planNodeId2 = new PlanNodeId("2");
       fragmentInstanceContext.addOperatorContext(
-          2, new PlanNodeId("2"), SeriesScanOperator.class.getSimpleName());
+          2, planNodeId2, SeriesScanOperator.class.getSimpleName());
       fragmentInstanceContext.addOperatorContext(
           3, new PlanNodeId("3"), TimeJoinOperator.class.getSimpleName());
       fragmentInstanceContext.addOperatorContext(
           4, new PlanNodeId("4"), LimitOperator.class.getSimpleName());
       SeriesScanOperator seriesScanOperator1 =
           new SeriesScanOperator(
+              planNodeId1,
               measurementPath1,
               allSensors,
               TSDataType.INT32,
@@ -113,6 +116,7 @@ public class LimitOperatorTest {
           new MeasurementPath(TIME_JOIN_OPERATOR_TEST_SG + ".device0.sensor1", 
TSDataType.INT32);
       SeriesScanOperator seriesScanOperator2 =
           new SeriesScanOperator(
+              planNodeId2,
               measurementPath2,
               allSensors,
               TSDataType.INT32,
diff --git 
a/server/src/test/java/org/apache/iotdb/db/mpp/operator/SeriesScanOperatorTest.java
 
b/server/src/test/java/org/apache/iotdb/db/mpp/operator/SeriesScanOperatorTest.java
index 669e374bbd..3793504970 100644
--- 
a/server/src/test/java/org/apache/iotdb/db/mpp/operator/SeriesScanOperatorTest.java
+++ 
b/server/src/test/java/org/apache/iotdb/db/mpp/operator/SeriesScanOperatorTest.java
@@ -84,10 +84,12 @@ public class SeriesScanOperatorTest {
       FragmentInstanceContext fragmentInstanceContext =
           new FragmentInstanceContext(
               new FragmentInstanceId(new PlanFragmentId(queryId, 0), 
"stub-instance"), state);
+      PlanNodeId planNodeId = new PlanNodeId("1");
       fragmentInstanceContext.addOperatorContext(
-          1, new PlanNodeId("1"), SeriesScanOperator.class.getSimpleName());
+          1, planNodeId, SeriesScanOperator.class.getSimpleName());
       SeriesScanOperator seriesScanOperator =
           new SeriesScanOperator(
+              planNodeId,
               measurementPath,
               allSensors,
               TSDataType.INT32,
diff --git 
a/server/src/test/java/org/apache/iotdb/db/mpp/operator/TimeJoinOperatorTest.java
 
b/server/src/test/java/org/apache/iotdb/db/mpp/operator/TimeJoinOperatorTest.java
index 2289ae7f4e..5534418b84 100644
--- 
a/server/src/test/java/org/apache/iotdb/db/mpp/operator/TimeJoinOperatorTest.java
+++ 
b/server/src/test/java/org/apache/iotdb/db/mpp/operator/TimeJoinOperatorTest.java
@@ -86,14 +86,17 @@ public class TimeJoinOperatorTest {
       FragmentInstanceContext fragmentInstanceContext =
           new FragmentInstanceContext(
               new FragmentInstanceId(new PlanFragmentId(queryId, 0), 
"stub-instance"), state);
+      PlanNodeId planNodeId1 = new PlanNodeId("1");
       fragmentInstanceContext.addOperatorContext(
-          1, new PlanNodeId("1"), SeriesScanOperator.class.getSimpleName());
+          1, planNodeId1, SeriesScanOperator.class.getSimpleName());
+      PlanNodeId planNodeId2 = new PlanNodeId("2");
       fragmentInstanceContext.addOperatorContext(
-          2, new PlanNodeId("2"), SeriesScanOperator.class.getSimpleName());
+          2, planNodeId2, SeriesScanOperator.class.getSimpleName());
       fragmentInstanceContext.addOperatorContext(
           3, new PlanNodeId("3"), TimeJoinOperator.class.getSimpleName());
       SeriesScanOperator seriesScanOperator1 =
           new SeriesScanOperator(
+              planNodeId1,
               measurementPath1,
               allSensors,
               TSDataType.INT32,
@@ -107,6 +110,7 @@ public class TimeJoinOperatorTest {
           new MeasurementPath(TIME_JOIN_OPERATOR_TEST_SG + ".device0.sensor1", 
TSDataType.INT32);
       SeriesScanOperator seriesScanOperator2 =
           new SeriesScanOperator(
+              planNodeId2,
               measurementPath2,
               allSensors,
               TSDataType.INT32,

Reply via email to