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,
