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 86bb0c27bc6ee29127690a8eb3259cb54879c91b
Merge: b17c657 c83ccfa
Author: Jinrui.Zhang <[email protected]>
AuthorDate: Sun Mar 20 16:44:28 2022 +0800

    Merge branch 'master' into xingtanzjr/logical_to_distributed

 .github/workflows/grafana-plugin.yml               |     2 +-
 .github/workflows/greetings.yml                    |     4 -
 .github/workflows/sonar-coveralls.yml              |     7 -
 .../org/apache/iotdb/db/qp/sql/IoTDBSqlParser.g4   |     8 +-
 .../java/org/apache/iotdb/cluster/ClientMain.java  |     2 +-
 .../org/apache/iotdb/cluster/ClusterIoTDB.java     |    20 +-
 .../iotdb/cluster/config/ClusterDescriptor.java    |     2 +-
 .../iotdb/cluster/coordinator/Coordinator.java     |     2 +-
 .../apache/iotdb/cluster/log/LogDispatcher.java    |     4 +-
 .../iotdb/cluster/log/catchup/LogCatchUpTask.java  |     2 +-
 .../iotdb/cluster/log/manage/RaftLogManager.java   |     2 +-
 .../serializable/SyncLogDequeSerializer.java       |     2 +-
 .../iotdb/cluster/metadata/CSchemaEngine.java      |     2 +-
 .../iotdb/cluster/query/ClusterPlanExecutor.java   |     2 +-
 .../query/last/ClusterLastQueryExecutor.java       |     2 +-
 .../iotdb/cluster/server/ClusterRPCService.java    |    10 +-
 .../cluster/server/ClusterRPCServiceMBean.java     |     2 +-
 .../cluster/server/PullSnapshotHintService.java    |     2 +-
 .../server/clusterinfo/ClusterInfoServer.java      |    10 +-
 .../cluster/server/member/DataGroupMember.java     |     6 +-
 .../cluster/server/member/MetaGroupMember.java     |     6 +-
 .../iotdb/cluster/server/member/RaftMember.java    |     8 +-
 .../cluster/server/raft/AbstractRaftService.java   |     6 +-
 .../server/raft/DataRaftHeartBeatService.java      |     8 +-
 .../iotdb/cluster/server/raft/DataRaftService.java |     8 +-
 .../server/raft/MetaRaftHeartBeatService.java      |     8 +-
 .../iotdb/cluster/server/raft/MetaRaftService.java |     8 +-
 .../cluster/server/service/DataGroupEngine.java    |     6 +-
 .../apache/iotdb/cluster/utils/PlanSerializer.java |     4 +-
 .../cluster/utils/nodetool/ClusterMonitor.java     |    14 +-
 .../org/apache/iotdb/cluster/common/IoTDBTest.java |     4 +-
 .../iotdb/cluster/integration/SingleNodeTest.java  |     2 +-
 .../iotdb/cluster/log/LogDispatcherTest.java       |     2 +-
 .../cluster/log/applier/DataLogApplierTest.java    |     8 +-
 .../cluster/log/catchup/LogCatchUpTaskTest.java    |     2 +-
 .../manage/MetaSingleSnapshotLogManagerTest.java   |     2 +-
 .../serializable/SyncLogDequeSerializerTest.java   |     2 +-
 .../cluster/log/snapshot/DataSnapshotTest.java     |     2 +-
 .../log/snapshot/MetaSimpleSnapshotTest.java       |     5 +-
 .../cluster/log/snapshot/PullSnapshotTaskTest.java |     2 +-
 .../query/ClusterUDTFQueryExecutorTest.java        |     2 +-
 .../iotdb/cluster/server/member/BaseMember.java    |     2 +-
 .../resources/conf/iotdb-confignode.properties     |    41 +-
 .../iotdb/confignode/conf/ConfigNodeConf.java      |   116 +-
 .../iotdb/confignode/service/ConfigNode.java       |    31 +-
 .../confignode/service/register/IService.java      |    51 -
 .../confignode/service/register/JMXService.java    |   105 -
 .../service/register/RegisterManager.java          |    82 -
 .../confignode/service/startup/StartupCheck.java   |    28 -
 .../confignode/service/startup/StartupChecks.java  |    89 -
 .../service/thrift/server/ConfigNodeRPCServer.java |    88 +-
 ...rver.java => ConfigNodeRPCServerProcessor.java} |     6 +-
 .../thrift/server/ConfigNodeRPCServiceHandler.java |    52 +
 .../utils/ConfigNodeEnvironmentUtils.java          |     9 +-
 docs/UserGuide/Query-Data/Overview.md              |    27 +-
 docs/zh/UserGuide/Query-Data/Overview.md           |    27 +-
 grafana-plugin/package.json                        |     4 +-
 grafana-plugin/src/componments/ControlValue.tsx    |     5 +-
 grafana-plugin/src/componments/FromValue.tsx       |     8 +-
 grafana-plugin/src/componments/SelectValue.tsx     |     8 +-
 grafana-plugin/src/componments/WhereValue.tsx      |     5 +-
 grafana-plugin/src/datasource.ts                   |    16 +-
 grafana-plugin/yarn.lock                           | 10529 +++++++++----------
 .../apache/iotdb/db/integration/IoTDBDaemonIT.java |     2 +-
 .../iotdb/db/integration/IoTDBLargeDataIT.java     |     2 +-
 .../iotdb/db/integration/IoTDBMetadataFetchIT.java |     2 +-
 .../iotdb/db/integration/IoTDBMultiSeriesIT.java   |     2 +-
 .../db/integration/IoTDBRecoverUnclosedIT.java     |     2 +-
 .../iotdb/db/integration/IoTDBUDFManagementIT.java |     6 +-
 .../iotdb/session/IoTDBSessionComplexIT.java       |     2 +-
 .../iotdb/session/IoTDBSessionIteratorIT.java      |     2 +-
 .../apache/iotdb/session/IoTDBSessionSimpleIT.java |     2 +-
 .../session/IoTDBSessionSyntaxConventionIT.java    |     2 +-
 integration/src/test/resources/logback.xml         |     4 +-
 {iotdb-commons => node-commons}/pom.xml            |    20 +-
 .../apache/iotdb/commons/ServerCommandLine.java    |     0
 .../apache/iotdb/commons}/concurrent/HashLock.java |     2 +-
 .../concurrent/IoTDBDaemonThreadFactory.java       |     2 +-
 .../IoTDBDefaultThreadExceptionHandler.java        |     2 +-
 .../concurrent/IoTDBThreadPoolFactory.java         |    10 +-
 .../commons}/concurrent/IoTThreadFactory.java      |     2 +-
 .../iotdb/commons}/concurrent/ThreadName.java      |     2 +-
 .../iotdb/commons}/concurrent/WrappedRunnable.java |     2 +-
 .../concurrent/threadpool/IThreadPoolMBean.java    |     2 +-
 .../WrappedScheduledExecutorService.java           |     6 +-
 .../WrappedScheduledExecutorServiceMBean.java      |     2 +-
 .../WrappedSingleThreadExecutorService.java        |     6 +-
 .../WrappedSingleThreadExecutorServiceMBean.java   |     2 +-
 .../WrappedSingleThreadScheduledExecutor.java      |     6 +-
 .../WrappedSingleThreadScheduledExecutorMBean.java |     2 +-
 .../threadpool/WrappedThreadPoolExecutor.java      |     8 +-
 .../threadpool/WrappedThreadPoolExecutorMBean.java |     2 +-
 .../apache/iotdb/commons}/conf/IoTDBConstant.java  |     2 +-
 .../iotdb/commons}/exception/IoTDBException.java   |     2 +-
 .../commons}/exception/ShutdownException.java      |     2 +-
 .../iotdb/commons}/exception/StartupException.java |     2 +-
 .../exception/runtime/RPCServiceException.java     |     2 +-
 .../apache/iotdb/commons/hash/APHashExecutor.java  |     0
 .../iotdb/commons/hash/BKDRHashExecutor.java       |     0
 .../commons/hash/DeviceGroupHashExecutor.java      |     0
 .../apache/iotdb/commons/hash/JSHashExecutor.java  |     0
 .../iotdb/commons/hash/SDBMHashExecutor.java       |     0
 .../service/AbstractThriftServiceThread.java       |    35 +-
 .../apache/iotdb/commons}/service/IService.java    |     6 +-
 .../apache/iotdb/commons}/service/JMXService.java  |     4 +-
 .../iotdb/commons}/service/RegisterManager.java    |     6 +-
 .../apache/iotdb/commons}/service/ServiceType.java |     9 +-
 .../iotdb/commons}/service/StartupCheck.java       |     4 +-
 .../iotdb/commons}/service/StartupChecks.java      |    10 +-
 .../iotdb/commons/service}/ThriftService.java      |    10 +-
 .../iotdb/commons/service/ThriftServiceThread.java |    89 +
 .../apache/iotdb/commons/utils/JVMCommonUtils.java |    81 +
 .../org/apache/iotdb/commons/utils/TestOnly.java   |     0
 .../IoTDBDefaultThreadExceptionHandlerTest.java    |     2 +-
 .../iotdb/commons}/IoTDBThreadPoolFactoryTest.java |     5 +-
 pom.xml                                            |     2 +-
 .../org/apache/iotdb/db/auth/AuthorityChecker.java |     2 +-
 .../iotdb/db/auth/authorizer/BasicAuthorizer.java  |     8 +-
 .../iotdb/db/auth/role/BasicRoleManager.java       |     2 +-
 .../iotdb/db/auth/role/LocalFileRoleAccessor.java  |     2 +-
 .../iotdb/db/auth/user/BasicUserManager.java       |     2 +-
 .../iotdb/db/auth/user/LocalFileUserAccessor.java  |     2 +-
 .../java/org/apache/iotdb/db/conf/IoTDBConfig.java |     1 +
 .../org/apache/iotdb/db/conf/IoTDBConfigCheck.java |     1 +
 .../org/apache/iotdb/db/conf/IoTDBDescriptor.java  |     5 +-
 .../db/conf/directories/DirectoryManager.java      |     2 +-
 .../directories/strategy/DirectoryStrategy.java    |     4 +-
 .../strategy/MaxDiskUsableSpaceFirstStrategy.java  |     6 +-
 .../MinFolderOccupiedSpaceFirstStrategy.java       |     6 +-
 .../strategy/RandomOnDiskUsableSpaceStrategy.java  |     4 +-
 .../directories/strategy/SequenceStrategy.java     |     6 +-
 .../db/conf/rest/IoTDBRestServiceDescriptor.java   |     2 +-
 .../org/apache/iotdb/db/engine/StorageEngine.java  |    10 +-
 .../db/engine/cache/CacheHitRatioMonitor.java      |    10 +-
 .../db/engine/cache/TimeSeriesMetadataCache.java   |     2 +-
 .../engine/compaction/CompactionTaskManager.java   |    12 +-
 .../db/engine/compaction/CompactionUtils.java      |     2 +-
 .../db/engine/compaction/TsFileIdentifier.java     |     2 +-
 .../CrossSpaceCompactionExceptionHandler.java      |     2 +-
 .../RewriteCrossSpaceCompactionSelector.java       |     2 +-
 .../task/RewriteCrossCompactionRecoverTask.java    |     2 +-
 .../task/RewriteCrossSpaceCompactionTask.java      |     2 +-
 .../inner/AbstractInnerSpaceCompactionTask.java    |     2 +-
 .../InnerSpaceCompactionExceptionHandler.java      |     2 +-
 .../SizeTieredCompactionRecoverTask.java           |     2 +-
 .../sizetiered/SizeTieredCompactionSelector.java   |     2 +-
 .../inner/sizetiered/SizeTieredCompactionTask.java |     2 +-
 .../inner/utils/InnerSpaceCompactionUtils.java     |     2 +-
 .../compaction/task/AbstractCompactionTask.java    |     2 +-
 .../compaction/task/CompactionRecoverTask.java     |     2 +-
 .../utils/log/CompactionLogAnalyzer.java           |     2 +-
 .../iotdb/db/engine/cq/ContinuousQueryService.java |     8 +-
 .../iotdb/db/engine/cq/ContinuousQueryTask.java    |     4 +-
 .../engine/cq/ContinuousQueryTaskPoolManager.java  |     4 +-
 .../apache/iotdb/db/engine/flush/FlushManager.java |    10 +-
 .../engine/flush/pool/FlushSubTaskPoolManager.java |     4 +-
 .../db/engine/flush/pool/FlushTaskPoolManager.java |     4 +-
 .../apache/iotdb/db/engine/settle/SettleTask.java  |     2 +-
 .../db/engine/storagegroup/TsFileManager.java      |     2 +-
 .../engine/storagegroup/TsFileNameGenerator.java   |     4 +-
 .../db/engine/storagegroup/TsFileResource.java     |     2 +-
 .../storagegroup/VirtualStorageGroupProcessor.java |     8 +-
 .../virtualSg/StorageGroupManager.java             |     2 +-
 .../service/TriggerRegistrationService.java        |    22 +-
 .../iotdb/db/engine/upgrade/UpgradeTask.java       |     2 +-
 .../iotdb/db/exception/ConfigurationException.java |     1 +
 .../iotdb/db/exception/LoadFileException.java      |     1 +
 .../apache/iotdb/db/exception/MergeException.java  |     1 +
 .../db/exception/QueryIdNotExsitException.java     |     1 +
 .../exception/QueryInBatchStatementException.java  |     1 +
 .../iotdb/db/exception/StorageEngineException.java |     1 +
 .../exception/StorageGroupProcessorException.java  |     1 +
 .../db/exception/SyncConnectionException.java      |     1 +
 .../SyncDeviceOwnerConflictException.java          |     1 +
 .../iotdb/db/exception/SystemCheckException.java   |     1 +
 .../db/exception/TsFileProcessorException.java     |     1 +
 .../iotdb/db/exception/WriteProcessException.java  |     1 +
 .../db/exception/index/IndexManagerException.java  |     2 +-
 .../db/exception/metadata/MetadataException.java   |     2 +-
 .../exception/query/LogicalOperatorException.java  |     2 +-
 .../exception/query/LogicalOptimizeException.java  |     2 +-
 .../db/exception/query/QueryProcessException.java  |     2 +-
 .../{runtime => sql}/SQLParserException.java       |     2 +-
 .../StatementAnalyzeException.java}                |    21 +-
 .../apache/iotdb/db/metadata/MetadataConstant.java |     2 +-
 .../org/apache/iotdb/db/metadata/SchemaEngine.java |     4 +-
 .../db/metadata/StorageGroupSchemaManager.java     |     4 +-
 .../org/apache/iotdb/db/metadata/mnode/MNode.java  |     2 +-
 .../iotdb/db/metadata/mtree/MTreeAboveSG.java      |     4 +-
 .../iotdb/db/metadata/mtree/MTreeBelowSG.java      |     2 +-
 .../db/metadata/mtree/traverser/Traverser.java     |     6 +-
 .../iotdb/db/metadata/path/MeasurementPath.java    |     2 +-
 .../apache/iotdb/db/metadata/path/PartialPath.java |     4 +-
 .../db/metadata/template/TemplateManager.java      |     2 +-
 .../db/metadata/upgrade/MetadataUpgrader.java      |     4 +-
 .../iotdb/db/metadata/utils/MetaFormatUtils.java   |    10 +-
 .../apache/iotdb/db/metadata/utils/MetaUtils.java  |     2 +-
 .../org/apache/iotdb/db/mpp/common/Analysis.java   |    48 +-
 .../org/apache/iotdb/db/mpp/common/DataRegion.java |    39 +-
 .../iotdb/db/mpp/common/DataRegionTimeSlice.java   |     5 +-
 .../mpp/common/{TreeNode.java => FragmentId.java}  |    40 +-
 .../{InstanceId.java => FragmentInstanceId.java}   |     4 +-
 .../{QueryContext.java => MPPQueryContext.java}    |     2 +-
 .../org/apache/iotdb/db/mpp/common/OrderBy.java    |    27 -
 .../org/apache/iotdb/db/mpp/common/QueryId.java    |    90 +-
 .../apache/iotdb/db/mpp/common/SchemaRegion.java   |    10 +-
 .../org/apache/iotdb/db/mpp/common/TsBlock.java    |    25 +
 .../FragmentInfo.java}                             |    33 +-
 ...ceContext.java => FragmentInstanceContext.java} |    35 +-
 .../db/mpp/execution/FragmentInstanceState.java    |    68 +
 .../iotdb/db/mpp/execution/FragmentState.java      |    71 +
 .../iotdb/db/mpp/execution/QueryExecution.java     |    12 +-
 .../mpp/execution/executor/InstanceExecutor.java   |    39 -
 .../db/mpp/execution/executor/InstanceHandle.java  |    31 -
 .../ClusterScheduler.java}                         |    32 +-
 .../scheduler/IScheduler.java}                     |    33 +-
 .../scheduler/StandaloneScheduler.java}            |    39 +-
 .../iotdb/db/mpp/operator/OperatorContext.java     |     4 +
 .../db/mpp/operator/process/LimitOperator.java     |    32 +-
 .../iotdb/db/mpp/operator/sink/SinkOperator.java   |     2 +-
 .../db/mpp/sql/planner/LocalExecutionPlanner.java  |    14 +-
 .../sql/planner/optimization/PlanOptimizer.java    |     5 +-
 .../mpp/sql/planner/plan/DistributedQueryPlan.java |     4 +-
 .../mpp/sql/planner/plan/DistributionPlanner.java  |   142 +-
 .../db/mpp/sql/planner/plan/FragmentInstance.java  |     2 +
 .../mpp/sql/planner/plan/FragmentInstanceId.java   |    30 -
 .../db/mpp/sql/planner/plan/LogicalPlanner.java    |     7 +-
 .../db/mpp/sql/planner/plan/LogicalQueryPlan.java  |     8 +-
 .../sql/planner/plan/node/PlanNodeAllocator.java   |    11 +-
 .../db/mpp/sql/planner/plan/node/PlanNodeUtil.java |    27 +-
 .../planner/plan/node/SimplePlanNodeRewriter.java  |    31 +-
 .../planner/plan/node/process/DeviceMergeNode.java |     2 +-
 .../sql/planner/plan/node/process/OffsetNode.java  |     2 +-
 .../sql/planner/plan/node/process/SortNode.java    |     2 +-
 .../planner/plan/node/process/TimeJoinNode.java    |     8 +-
 .../planner/plan/node/source/SeriesScanNode.java   |     7 +-
 .../influxdb/meta/InfluxDBMetaManager.java         |     2 +-
 .../apache/iotdb/db/protocol/rest/RestService.java |     6 +-
 .../db/protocol/rest/handler/ExceptionHandler.java |     2 +-
 .../protocol/rest/impl/GrafanaApiServiceImpl.java  |     2 +-
 .../db/protocol/rest/impl/PingApiServiceImpl.java  |     2 +-
 .../db/protocol/rest/impl/RestApiServiceImpl.java  |     2 +-
 .../main/java/org/apache/iotdb/db/qp/Planner.java  |     2 +-
 .../apache/iotdb/db/qp/executor/PlanExecutor.java  |    76 +-
 .../db/qp/logical/crud/BasicFunctionOperator.java  |     2 +-
 .../db/qp/logical/crud/BasicOperatorType.java      |     4 +-
 .../db/qp/logical/crud/DeleteDataOperator.java     |     2 +-
 .../db/qp/logical/crud/FillQueryOperator.java      |     2 +-
 .../iotdb/db/qp/logical/crud/InsertOperator.java   |     2 +-
 .../sys/CreateAlignedTimeSeriesOperator.java       |     2 +-
 .../iotdb/db/qp/physical/crud/InsertRowPlan.java   |     2 +-
 .../iotdb/db/qp/physical/sys/SetTemplatePlan.java  |     2 +-
 .../db/qp/physical/sys/UnsetTemplatePlan.java      |     2 +-
 .../apache/iotdb/db/qp/sql/IoTDBSqlVisitor.java    |    15 +-
 .../iotdb/db/qp/strategy/LogicalGenerator.java     |     4 +-
 .../qp/strategy/optimizer/ConcatPathOptimizer.java |     2 +-
 .../iotdb/db/qp/utils/GroupByLevelController.java  |     2 +-
 .../iotdb/db/query/control/QueryTimeManager.java   |     6 +-
 .../iotdb/db/query/control/SessionManager.java     |     2 +-
 .../db/query/control/SessionTimeoutManager.java    |     2 +-
 .../db/query/dataset/NonAlignEngineDataSet.java    |     2 +-
 .../dataset/RawQueryDataSetWithoutValueFilter.java |     2 +-
 .../iotdb/db/query/dataset/ShowDevicesDataSet.java |     6 +-
 .../db/query/dataset/ShowTimeseriesDataSet.java    |    16 +-
 .../db/query/executor/AggregationExecutor.java     |     2 +-
 .../iotdb/db/query/executor/LastQueryExecutor.java |     6 +-
 .../query/expression/unary/FunctionExpression.java |     2 +-
 .../iotdb/db/query/pool/QueryTaskManager.java      |     4 +-
 .../db/query/pool/RawQueryReadTaskPoolManager.java |     4 +-
 .../row/SerializableRowRecordList.java             |     2 +-
 .../datastructure/tv/SerializableBinaryTVList.java |     2 +-
 .../tv/SerializableBooleanTVList.java              |     2 +-
 .../datastructure/tv/SerializableDoubleTVList.java |     2 +-
 .../datastructure/tv/SerializableFloatTVList.java  |     2 +-
 .../datastructure/tv/SerializableIntTVList.java    |     2 +-
 .../datastructure/tv/SerializableLongTVList.java   |     2 +-
 .../udf/service/TemporaryQueryDataFileService.java |     6 +-
 .../query/udf/service/UDFClassLoaderManager.java   |     6 +-
 .../query/udf/service/UDFRegistrationService.java  |     6 +-
 .../org/apache/iotdb/db/rescon/SystemInfo.java     |     2 +-
 .../iotdb/db/service/InfluxDBRPCService.java       |     9 +-
 .../java/org/apache/iotdb/db/service/IoTDB.java    |     9 +-
 .../org/apache/iotdb/db/service/MQTTService.java   |     2 +
 .../org/apache/iotdb/db/service/RPCService.java    |     9 +-
 .../apache/iotdb/db/service/RPCServiceMBean.java   |     2 +-
 .../org/apache/iotdb/db/service/SettleService.java |     6 +-
 .../org/apache/iotdb/db/service/StaticResps.java   |     6 +-
 .../org/apache/iotdb/db/service/UpgradeSevice.java |     4 +-
 .../db/service/basic/QueryFrequencyRecorder.java   |     2 +-
 .../iotdb/db/service/basic/ServiceProvider.java    |     2 +-
 .../iotdb/db/service/metrics/MetricsService.java   |    10 +-
 .../db/service/metrics/MetricsServiceMBean.java    |     2 +-
 .../db/service/thrift/impl/TSServiceImpl.java      |     4 +-
 .../iotdb/db/sql/constant/FilterConstant.java      |   102 +
 .../iotdb/db/sql/constant/StatementType.java       |   134 +
 .../org/apache/iotdb/db/sql/parser/ASTVisitor.java |  1368 +++
 .../iotdb/db/sql/parser/StatementGenerator.java    |   184 +
 .../statement/AggregationQueryStatement.java}      |    35 +-
 .../statement/FillQueryStatement.java}             |    20 +-
 .../statement/GroupByFillQueryStatement.java}      |    28 +-
 .../db/sql/statement/GroupByQueryStatement.java    |    29 +-
 .../statement/LastQueryStatement.java}             |    14 +-
 .../iotdb/db/sql/statement/QueryStatement.java     |   211 +
 .../statement/ShowDevicesStatement.java}           |    16 +-
 .../statement/ShowTimeSeriesStatement.java}        |    16 +-
 .../statement/Statement.java}                      |    39 +-
 .../iotdb/db/sql/statement/UDAFQueryStatement.java |    10 +-
 .../iotdb/db/sql/statement/UDTFQueryStatement.java |    10 +-
 .../db/sql/statement/component/FillComponent.java  |    36 +-
 .../db/sql/statement/component/FromComponent.java  |    29 +-
 .../statement/component/GroupByLevelComponent.java |    31 +-
 .../statement/component/GroupByTimeComponent.java  |    99 +
 .../iotdb/db/sql/statement/component/OrderBy.java  |     9 +-
 .../db/sql/statement/component/ResultColumn.java   |   168 +
 .../sql/statement/component/ResultSetFormat.java   |    10 +-
 .../sql/statement/component/SelectComponent.java   |    99 +
 .../statement/component/WhereCondition.java}       |    27 +-
 .../db/sql/statement/component/WithoutPolicy.java  |    59 +
 .../statement/filter/BasicFilterType.java}         |    18 +-
 .../statement/filter/BasicFunctionFilter.java}     |    35 +-
 .../statement/filter/FunctionFilter.java}          |    36 +-
 .../iotdb/db/sql/statement/filter/InFilter.java    |   201 +
 .../iotdb/db/sql/statement/filter/LikeFilter.java  |   134 +
 .../iotdb/db/sql/statement/filter/QueryFilter.java |   295 +
 .../db/sql/statement/filter/RegexpFilter.java      |   134 +
 .../iotdb/db/sync/conf/SyncSenderDescriptor.java   |     2 +-
 .../iotdb/db/sync/receiver/SyncServerManager.java  |    10 +-
 .../db/sync/receiver/SyncServerManagerMBean.java   |     2 +-
 .../db/sync/receiver/load/FileLoaderManager.java   |     4 +-
 .../db/sync/receiver/transfer/SyncServiceImpl.java |     2 +-
 .../db/sync/sender/manage/SyncFileManager.java     |     2 +-
 .../iotdb/db/sync/sender/transfer/SyncClient.java  |     4 +-
 .../org/apache/iotdb/db/tools/TsFileSplitTool.java |     2 +-
 .../java/org/apache/iotdb/db/utils/AuthUtils.java  |     2 +-
 .../org/apache/iotdb/db/utils/CommonUtils.java     |    57 -
 .../apache/iotdb/db/utils/ErrorHandlingUtils.java  |     4 +-
 .../java/org/apache/iotdb/db/utils/MemUtils.java   |     2 +-
 .../org/apache/iotdb/db/utils/OpenFileNumUtil.java |     2 +-
 .../org/apache/iotdb/db/utils/ThreadUtils.java     |     2 +-
 .../windowing/runtime/WindowEvaluationTask.java    |     2 +-
 .../runtime/WindowEvaluationTaskPoolManager.java   |     6 +-
 .../writelog/manager/MultiFileLogNodeManager.java  |     8 +-
 .../db/writelog/node/ExclusiveWriteLogNode.java    |     4 +-
 .../iotdb/db/writelog/recover/LogReplayer.java     |     2 +-
 .../apache/iotdb/db/conf/IoTDBDescriptorTest.java  |     2 +
 .../strategy/DirectoryStrategyTest.java            |    24 +-
 .../iotdb/db/engine/cache/ChunkCacheTest.java      |     2 +-
 .../db/engine/compaction/CompactionUtilsTest.java  |     2 +-
 .../cross/CrossSpaceCompactionExceptionTest.java   |     2 +-
 .../db/engine/compaction/cross/MergeTest.java      |     2 +-
 .../engine/compaction/cross/MergeUpgradeTest.java  |     2 +-
 .../cross/RewriteCompactionFileSelectorTest.java   |     2 +-
 .../RewriteCrossSpaceCompactionRecoverTest.java    |     2 +-
 .../cross/RewriteCrossSpaceCompactionTest.java     |     4 +-
 .../inner/AbstractInnerSpaceCompactionTest.java    |     4 +-
 .../inner/InnerCompactionMoreDataTest.java         |     4 +-
 .../compaction/inner/InnerCompactionTest.java      |     4 +-
 .../inner/InnerSpaceCompactionUtilsOldTest.java    |     2 +-
 .../SizeTieredCompactionRecoverTest.java           |     2 +-
 .../inner/sizetiered/SizeTieredCompactionTest.java |     4 +-
 ...eCrossSpaceCompactionRecoverCompatibleTest.java |     2 +-
 .../recover/SizeTieredCompactionRecoverTest.java   |     2 +-
 .../compaction/utils/CompactionClearUtils.java     |     2 +-
 .../utils/CompactionFileGeneratorUtils.java        |     2 +-
 .../storagegroup/StorageGroupProcessorTest.java    |     2 +-
 .../iotdb/db/engine/storagegroup/TTLTest.java      |     2 +-
 .../db/engine/storagegroup/TsFileManagerTest.java  |     2 +-
 .../db/metadata/upgrade/MetadataUpgradeTest.java   |     2 +-
 .../db/mpp/sql/plan/DistributionPlannerTest.java   |   101 +-
 .../iotdb/db/protocol/rest/IoTDBRestServiceIT.java |     2 +-
 .../java/org/apache/iotdb/db/qp/PlannerTest.java   |     4 +-
 .../iotdb/db/qp/logical/LogicalPlanSmallTest.java  |     2 +-
 .../iotdb/db/qp/physical/PhysicalPlanTest.java     |     2 +-
 .../apache/iotdb/db/qp/sql/ASTVisitorTest.java}    |    32 +-
 .../iotdb/db/qp/sql/IoTDBsqlVisitorTest.java       |     2 +-
 .../query/reader/series/SeriesReaderTestUtil.java  |     2 +-
 .../iotdb/db/rescon/ResourceManagerTest.java       |     4 +-
 .../iotdb/db/sql/StatementGeneratorTest.java       |    75 +
 .../db/sync/receiver/load/FileLoaderTest.java      |     2 +-
 .../recover/SyncReceiverLogAnalyzerTest.java       |     2 +-
 .../db/sync/sender/manage/SyncFileManagerTest.java |     2 +-
 .../sender/recover/SyncSenderLogAnalyzerTest.java  |     4 +-
 .../sync/sender/recover/SyncSenderLoggerTest.java  |     2 +-
 .../db/sync/sender/transfer/SyncClientTest.java    |     2 +-
 .../apache/iotdb/db/tools/IoTDBWatermarkTest.java  |     2 +-
 .../iotdb/db/writelog/IoTDBLogFileSizeTest.java    |     2 +-
 .../recover/RecoverResourceFromReaderTest.java     |     2 +-
 .../writelog/recover/UnseqTsFileRecoverTest.java   |     2 +-
 server/src/test/resources/logback.xml              |     4 +-
 .../session/IoTDBSessionDisableMemControlIT.java   |     2 +-
 .../iotdb/session/IoTDBSessionVectorInsertIT.java  |     2 +-
 .../java/org/apache/iotdb/session/SessionTest.java |     2 +-
 .../apache/iotdb/session/pool/SessionPoolTest.java |     2 +-
 .../apache/iotdb/session/template/TemplateUT.java  |     2 +-
 .../apache/iotdb/spark/db/EnvironmentUtils.java    |     2 +-
 .../org/apache/iotdb/spark/db/IoTDBTest.scala      |     3 +-
 .../org/apache/iotdb/spark/db/IoTDBWriteTest.scala |     3 +-
 .../iotdb/spark/db/unit/DataFrameToolsTest.scala   |     4 +-
 .../zeppelin/iotdb/IoTDBInterpreterTest.java       |     2 +-
 399 files changed, 10163 insertions(+), 7463 deletions(-)

diff --cc server/src/main/java/org/apache/iotdb/db/mpp/common/Analysis.java
index 85ff111,88d8884..2d2a397
--- 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
@@@ -25,36 -22,12 +25,36 @@@ import java.util.*
  
  /** Analysis used for planning a query. TODO: This class may need to store 
more info for a query. */
  public class Analysis {
-     // Description for each series. Such as dataType, existence
++  // Description for each series. Such as dataType, existence
 +
-     // Data distribution info for each series. Series -> [VSG, VSG]
++  // Data distribution info for each series. Series -> [VSG, VSG]
 +
-     // Map<PartialPath, List<FullPath>> Used to remove asterisk
++  // Map<PartialPath, List<FullPath>> Used to remove asterisk
  
-     // Statement
-     private String statement;
 -  private Statement statement;
++  // Statement
++  private String statement;
  
-     // DataPartitionInfo
-     private Map<String, Map<DataRegionTimeSlice, List<DataRegion>>> 
dataPartitionInfo;
 -  public Analysis() {}
++  // DataPartitionInfo
++  private Map<String, Map<DataRegionTimeSlice, List<DataRegion>>> 
dataPartitionInfo;
 +
-     // SchemaPartitionInfo
-     private Map<String, List<SchemaRegion>> schemaPartitionInfo;
++  // 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;
-         }
++  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;
 +    }
++  }
  
-     public void setDataPartitionInfo(Map<String, Map<DataRegionTimeSlice, 
List<DataRegion>>> dataPartitionInfo) {
-         this.dataPartitionInfo = dataPartitionInfo;
-     }
 -  public void setStatement(Statement rewrittenStatement) {
 -    this.statement = rewrittenStatement;
++  public void setDataPartitionInfo(
++      Map<String, Map<DataRegionTimeSlice, List<DataRegion>>> 
dataPartitionInfo) {
++    this.dataPartitionInfo = dataPartitionInfo;
+   }
  }
diff --cc server/src/main/java/org/apache/iotdb/db/mpp/common/DataRegion.java
index 5618e1f,612e930..04100e1
--- 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
@@@ -16,34 -16,33 +16,35 @@@
   * specific language governing permissions and limitations
   * under the License.
   */
 -package org.apache.iotdb.db.mpp.sql.planner.plan.node.source;
  
 -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 version. */
 -public class CsvSourceNode extends SourceNode {
 -
 -  public CsvSourceNode(PlanNodeId id) {
 -    super(id);
 +/**
-  * This class is used to represent the data partition info including the 
DataRegionId and physical node IP address
++ * 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
++// TODO: (xingtanzjr) This class should be substituted with the class defined 
in Consensus level
 +public class DataRegion {
-     private Integer dataRegionId;
-     private String endpoint;
++  private Integer dataRegionId;
++  private String endpoint;
 +
-     public DataRegion(Integer dataRegionId, String endpoint) {
-         this.dataRegionId = dataRegionId;
-         this.endpoint = endpoint;
-     }
++  public DataRegion(Integer dataRegionId, String endpoint) {
++    this.dataRegionId = dataRegionId;
++    this.endpoint = endpoint;
+   }
  
-     public int hashCode() {
-         return dataRegionId.hashCode();
-     }
 -  @Override
 -  public List<PlanNode> getChildren() {
 -    return null;
++  public int hashCode() {
++    return dataRegionId.hashCode();
+   }
  
-     public boolean equals(Object obj) {
-         if (obj instanceof DataRegion) {
-             return this.dataRegionId.equals(((DataRegion)obj).dataRegionId);
-         }
-         return false;
 -  @Override
 -  public List<String> getOutputColumnNames() {
 -    return null;
++  public boolean equals(Object obj) {
++    if (obj instanceof DataRegion) {
++      return this.dataRegionId.equals(((DataRegion) obj).dataRegionId);
 +    }
++    return false;
+   }
  
-     public String toString() {
-         return String.format("%s/%d",this.endpoint, this.dataRegionId);
-     }
 -  @Override
 -  public void close() throws Exception {}
 -
 -  @Override
 -  public void open() throws Exception {}
++  public String toString() {
++    return String.format("%s/%d", this.endpoint, this.dataRegionId);
++  }
  }
diff --cc 
server/src/main/java/org/apache/iotdb/db/mpp/common/DataRegionTimeSlice.java
index 52f48cc,a8d9b94..65dcd39
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/common/DataRegionTimeSlice.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/common/DataRegionTimeSlice.java
@@@ -16,10 -16,12 +16,9 @@@
   * specific language governing permissions and limitations
   * under the License.
   */
- 
  package org.apache.iotdb.db.mpp.common;
  
- //TODO: (xingtanzjr) This class should be substituted with the class defined 
in Consensus level
 -/** The traversal order for operators by timestamp */
 -public enum OrderBy {
 -  TIMESTAMP_ASC,
 -  TIMESTAMP_DESC,
 -  DEVICE_NAME_ASC,
 -  DEVICE_NAME_DESC,
++// TODO: (xingtanzjr) This class should be substituted with the class defined 
in Consensus level
 +public class DataRegionTimeSlice {
-     long startTimestamp;
++  long startTimestamp;
  }
diff --cc 
server/src/main/java/org/apache/iotdb/db/mpp/common/MPPQueryContext.java
index 2f6715d,2f6715d..4a29263
--- a/server/src/main/java/org/apache/iotdb/db/mpp/common/MPPQueryContext.java
+++ b/server/src/main/java/org/apache/iotdb/db/mpp/common/MPPQueryContext.java
@@@ -22,7 -22,7 +22,7 @@@ package org.apache.iotdb.db.mpp.common
   * This class is used to record the context of a query including QueryId, 
query statement, session
   * info and so on
   */
--public class QueryContext {
++public class MPPQueryContext {
    private String statement;
    private QueryId queryId;
    private QuerySession session;
diff --cc server/src/main/java/org/apache/iotdb/db/mpp/common/SchemaRegion.java
index 8352f10,2f6715d..a7ffe29
--- a/server/src/main/java/org/apache/iotdb/db/mpp/common/SchemaRegion.java
+++ b/server/src/main/java/org/apache/iotdb/db/mpp/common/SchemaRegion.java
@@@ -19,11 -19,11 +19,11 @@@
  package org.apache.iotdb.db.mpp.common;
  
  /**
-  * This class is used to represent the schema partition info including the 
DataRegionId and physical node IP address
 - * This class is used to record the context of a query including QueryId, 
query statement, session
 - * info and so on
++ * 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 QueryContext {
 -  private String statement;
 -  private QueryId queryId;
 -  private QuerySession session;
++// TODO: (xingtanzjr) This class should be substituted with the class defined 
in Consensus level
 +public class SchemaRegion {
-     private Integer DataRegionId;
-     private String endpoint;
++  private Integer DataRegionId;
++  private String endpoint;
  }
diff --cc 
server/src/main/java/org/apache/iotdb/db/mpp/execution/QueryExecution.java
index 9852350,62a2c65..ad2b176
--- a/server/src/main/java/org/apache/iotdb/db/mpp/execution/QueryExecution.java
+++ b/server/src/main/java/org/apache/iotdb/db/mpp/execution/QueryExecution.java
@@@ -19,7 -19,9 +19,9 @@@
  package org.apache.iotdb.db.mpp.execution;
  
  import org.apache.iotdb.db.mpp.common.Analysis;
--import org.apache.iotdb.db.mpp.common.QueryContext;
++import org.apache.iotdb.db.mpp.common.MPPQueryContext;
+ import org.apache.iotdb.db.mpp.execution.scheduler.ClusterScheduler;
+ import org.apache.iotdb.db.mpp.execution.scheduler.IScheduler;
  import org.apache.iotdb.db.mpp.sql.planner.optimization.PlanOptimizer;
  import org.apache.iotdb.db.mpp.sql.planner.plan.*;
  
@@@ -33,8 -35,8 +35,8 @@@ import java.util.List
   * corresponding physical nodes. 3. Collect and monitor the progress/states 
of this query.
   */
  public class QueryExecution {
--  private QueryContext context;
-   private QueryScheduler scheduler;
++  private MPPQueryContext context;
+   private IScheduler scheduler;
    private QueryStateMachine stateMachine;
  
    private List<PlanOptimizer> planOptimizers;
@@@ -45,7 -47,7 +47,7 @@@
    private List<PlanFragment> fragments;
    private List<FragmentInstance> fragmentInstances;
  
--  public QueryExecution(QueryContext context) {
++  public QueryExecution(MPPQueryContext context) {
      this.context = context;
    }
  
diff --cc 
server/src/main/java/org/apache/iotdb/db/mpp/operator/sink/SinkOperator.java
index c3f16ca,c3f16ca..d362c1d
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/operator/sink/SinkOperator.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/operator/sink/SinkOperator.java
@@@ -7,7 -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
diff --cc 
server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/optimization/PlanOptimizer.java
index ff2baae,b091b91..94ada95
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/optimization/PlanOptimizer.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/optimization/PlanOptimizer.java
@@@ -18,10 -18,9 +18,9 @@@
   */
  package org.apache.iotdb.db.mpp.sql.planner.optimization;
  
--import org.apache.iotdb.db.mpp.common.QueryContext;
- import org.apache.iotdb.db.mpp.common.TsBlock;
++import org.apache.iotdb.db.mpp.common.MPPQueryContext;
  import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanNode;
  
  public interface PlanOptimizer {
--  PlanNode optimize(PlanNode plan, QueryContext context);
++  PlanNode optimize(PlanNode plan, MPPQueryContext context);
  }
diff --cc 
server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/DistributedQueryPlan.java
index 50835dd,50835dd..8cf0341
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/DistributedQueryPlan.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/DistributedQueryPlan.java
@@@ -18,13 -18,13 +18,13 @@@
   */
  package org.apache.iotdb.db.mpp.sql.planner.plan;
  
--import org.apache.iotdb.db.mpp.common.QueryContext;
++import org.apache.iotdb.db.mpp.common.MPPQueryContext;
  import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanNode;
  
  import java.util.List;
  
  public class DistributedQueryPlan {
--  private QueryContext context;
++  private MPPQueryContext context;
    private PlanNode rootNode;
    private PlanFragment rootFragment;
  
diff --cc 
server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/DistributionPlanner.java
index 9c55a46,fc428cf..9b146d5
--- 
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,91 -19,17 +19,101 @@@
  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 PlanNode rewriteSource() {
-         SourceRewriter rewriter = new SourceRewriter();
-         return rewriter.visit(logicalPlan.getRootNode(), new 
DistributionPlanContext());
-     }
+   public DistributionPlanner(Analysis analysis, LogicalQueryPlan logicalPlan) 
{
+     this.analysis = analysis;
+     this.logicalPlan = logicalPlan;
+   }
  
-     public DistributedQueryPlan planFragments() {
-         return null;
-     }
++  public PlanNode rewriteSource() {
++    SourceRewriter rewriter = new SourceRewriter();
++    return rewriter.visit(logicalPlan.getRootNode(), new 
DistributionPlanContext());
++  }
 +
-     private class SourceRewriter extends 
SimplePlanNodeRewriter<DistributionPlanContext> {
-         public PlanNode visitTimeJoin(TimeJoinNode node, 
DistributionPlanContext context) {
-             TimeJoinNode root = (TimeJoinNode) node.clone();
+   public DistributedQueryPlan planFragments() {
+     return null;
+   }
 +
-             // 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(visit(child, context));
-                 }
-             }
++  private class SourceRewriter extends 
SimplePlanNodeRewriter<DistributionPlanContext> {
++    public PlanNode visitTimeJoin(TimeJoinNode node, DistributionPlanContext 
context) {
++      TimeJoinNode root = (TimeJoinNode) node.clone();
 +
-             // 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 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(visit(child, context));
++        }
++      }
 +
-             // 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);
-                 }
-             });
++      // 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);
++          });
 +
-             return root;
-         }
++      // 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);
++            }
++          });
 +
-         public PlanNode visit(PlanNode node, DistributionPlanContext context) 
{
-             return node.accept(this, context);
-         }
++      return root;
 +    }
 +
-     private class DistributionPlanContext {
- 
++    public PlanNode visit(PlanNode node, DistributionPlanContext context) {
++      return node.accept(this, context);
 +    }
++  }
++
++  private class DistributionPlanContext {}
  }
diff --cc 
server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/LogicalPlanner.java
index 0a45b6f,0a45b6f..ad0937d
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/LogicalPlanner.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/LogicalPlanner.java
@@@ -19,17 -19,17 +19,18 @@@
  package org.apache.iotdb.db.mpp.sql.planner.plan;
  
  import org.apache.iotdb.db.mpp.common.Analysis;
--import org.apache.iotdb.db.mpp.common.QueryContext;
++import org.apache.iotdb.db.mpp.common.MPPQueryContext;
  import org.apache.iotdb.db.mpp.sql.planner.optimization.PlanOptimizer;
  
  import java.util.List;
  
  public class LogicalPlanner {
    private Analysis analysis;
--  private QueryContext context;
++  private MPPQueryContext context;
    private List<PlanOptimizer> optimizers;
  
--  public LogicalPlanner(Analysis analysis, QueryContext context, 
List<PlanOptimizer> optimizers) {
++  public LogicalPlanner(
++      Analysis analysis, MPPQueryContext context, List<PlanOptimizer> 
optimizers) {
      this.analysis = analysis;
      this.context = context;
      this.optimizers = optimizers;
diff --cc 
server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/LogicalQueryPlan.java
index 788746a,666fbf3..664497d
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/LogicalQueryPlan.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/LogicalQueryPlan.java
@@@ -18,7 -18,7 +18,7 @@@
   */
  package org.apache.iotdb.db.mpp.sql.planner.plan;
  
--import org.apache.iotdb.db.mpp.common.QueryContext;
++import org.apache.iotdb.db.mpp.common.MPPQueryContext;
  import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanNode;
  
  /**
@@@ -26,19 -26,6 +26,19 @@@
   * plan node tree.
   */
  public class LogicalQueryPlan {
--  private QueryContext context;
++  private MPPQueryContext context;
    private PlanNode rootNode;
 +
-   public LogicalQueryPlan(QueryContext context, PlanNode rootNode) {
++  public LogicalQueryPlan(MPPQueryContext context, PlanNode rootNode) {
 +    this.context = context;
 +    this.rootNode = rootNode;
 +  }
 +
 +  public PlanNode getRootNode() {
 +    return rootNode;
 +  }
 +
-   public QueryContext getContext() {
++  public MPPQueryContext getContext() {
 +    return context;
 +  }
  }
diff --cc 
server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/PlanNodeAllocator.java
index edac517,7979b27..63c9a58
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/PlanNodeAllocator.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/PlanNodeAllocator.java
@@@ -16,13 -16,9 +16,14 @@@
   * specific language governing permissions and limitations
   * under the License.
   */
 -package org.apache.iotdb.db.mpp.common;
  
 -public enum FilterNullPolicy {
 -  CONTAINS_NULL,
 -  ALL_NULL
 +package org.apache.iotdb.db.mpp.sql.planner.plan.node;
 +
 +public class PlanNodeAllocator {
-     public static int initialId = 0;
-     public static synchronized PlanNodeId generateId() {
-         initialId++;
-         return new PlanNodeId(String.valueOf(initialId));
-     }
++  public static int initialId = 0;
++
++  public static synchronized PlanNodeId generateId() {
++    initialId++;
++    return new PlanNodeId(String.valueOf(initialId));
++  }
  }
diff --cc 
server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/PlanNodeUtil.java
index f174b74,2fddbdd..1a597b2
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/PlanNodeUtil.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/PlanNodeUtil.java
@@@ -16,25 -16,31 +16,24 @@@
   * specific language governing permissions and limitations
   * under the License.
   */
- 
 -package org.apache.iotdb.db.mpp.sql.planner.plan.node.sink;
 +package org.apache.iotdb.db.mpp.sql.planner.plan.node;
  
 -import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanNode;
 -import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanNodeId;
 -
 -import java.util.List;
 -
 -public class CsvSinkNode extends SinkNode {
 -  public CsvSinkNode(PlanNodeId id) {
 -    super(id);
 +public class PlanNodeUtil {
-     public static void printPlanNode(PlanNode root) {
-         printPlanNodeWithLevel(root, 0);
-     }
++  public static void printPlanNode(PlanNode root) {
++    printPlanNodeWithLevel(root, 0);
+   }
  
-     private static void printPlanNodeWithLevel(PlanNode root, int level) {
-         printTab(level);
-         System.out.println(root.toString());
-         for (PlanNode child : root.getChildren()) {
-             printPlanNodeWithLevel(child, level + 1);
-         }
 -  @Override
 -  public List<PlanNode> getChildren() {
 -    return null;
++  private static void printPlanNodeWithLevel(PlanNode root, int level) {
++    printTab(level);
++    System.out.println(root.toString());
++    for (PlanNode child : root.getChildren()) {
++      printPlanNodeWithLevel(child, level + 1);
 +    }
+   }
  
-     private static void printTab(int count) {
-         for(int i = 0 ; i < count; i ++) {
-             System.out.print("\t");
-         }
 -  @Override
 -  public List<String> getOutputColumnNames() {
 -    return null;
++  private static void printTab(int count) {
++    for (int i = 0; i < count; i++) {
++      System.out.print("\t");
 +    }
+   }
 -
 -  @Override
 -  public void close() throws Exception {}
 -
 -  @Override
 -  public void send() {}
  }
diff --cc 
server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/SimplePlanNodeRewriter.java
index 639e805,bb343f4..e0ca8f6
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/SimplePlanNodeRewriter.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/SimplePlanNodeRewriter.java
@@@ -21,25 -23,26 +21,24 @@@ package org.apache.iotdb.db.mpp.sql.pla
  
  import java.util.List;
  
- import static com.google.common.base.Verify.verifyNotNull;
 -/** not implemented in current IoTDB yet */
 -public class ThriftSinkNode extends SinkNode {
 -
 -  public ThriftSinkNode(PlanNodeId id) {
 -    super(id);
 -  }
 +import static com.google.common.collect.ImmutableList.toImmutableList;
  
- public class SimplePlanNodeRewriter<C> extends PlanVisitor<PlanNode, C>{
-     @Override
-     public PlanNode visitPlan(PlanNode node, C context) {
-         return defaultRewrite(node, context);
-     }
++public class SimplePlanNodeRewriter<C> extends PlanVisitor<PlanNode, C> {
+   @Override
 -  public List<PlanNode> getChildren() {
 -    return null;
++  public PlanNode visitPlan(PlanNode node, C context) {
++    return defaultRewrite(node, context);
+   }
  
-     public PlanNode defaultRewrite(PlanNode node, C context) {
-         List<PlanNode> children = node.getChildren().stream()
-                 .map(child -> rewrite(child, context))
-                 .collect(toImmutableList());
 -  @Override
 -  public List<String> getOutputColumnNames() {
 -    return null;
 -  }
++  public PlanNode defaultRewrite(PlanNode node, C context) {
++    List<PlanNode> children =
++        node.getChildren().stream()
++            .map(child -> rewrite(child, context))
++            .collect(toImmutableList());
  
-         return node.cloneWithChildren(children);
-     }
 -  @Override
 -  public void close() throws Exception {}
++    return node.cloneWithChildren(children);
++  }
  
-     public PlanNode rewrite(PlanNode node, C userContext)
-     {
-         return node.accept(this, userContext);
-     }
 -  @Override
 -  public void send() {}
++  public PlanNode rewrite(PlanNode node, C userContext) {
++    return node.accept(this, userContext);
++  }
  }
diff --cc 
server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/DeviceMergeNode.java
index 5073d15,54a0cf8..de1c7ea
--- 
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
@@@ -19,10 -19,10 +19,10 @@@
  package org.apache.iotdb.db.mpp.sql.planner.plan.node.process;
  
  import org.apache.iotdb.db.mpp.common.FilterNullPolicy;
--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;
  import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanVisitor;
++import org.apache.iotdb.db.sql.statement.component.OrderBy;
  
  import java.util.List;
  import java.util.Map;
diff --cc 
server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/OffsetNode.java
index bde203c,2e3fc78..f009fd4
--- 
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
@@@ -7,7 -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
diff --cc 
server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/SortNode.java
index b0752f0,19464a2..7e2b1c4
--- 
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
@@@ -18,10 -18,10 +18,10 @@@
   */
  package org.apache.iotdb.db.mpp.sql.planner.plan.node.process;
  
--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;
  import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanVisitor;
++import org.apache.iotdb.db.sql.statement.component.OrderBy;
  
  import com.google.common.collect.ImmutableList;
  
diff --cc 
server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/TimeJoinNode.java
index 907c91a,74fac14..7be2140
--- 
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
@@@ -19,13 -19,11 +19,13 @@@
  package org.apache.iotdb.db.mpp.sql.planner.plan.node.process;
  
  import org.apache.iotdb.db.mpp.common.FilterNullPolicy;
--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.PlanNodeAllocator;
  import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanNodeId;
  import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanVisitor;
++import org.apache.iotdb.db.sql.statement.component.OrderBy;
  
 +import java.util.ArrayList;
  import java.util.List;
  import java.util.stream.Collectors;
  
@@@ -48,16 -46,6 +48,13 @@@ public class TimeJoinNode extends Proce
  
    private List<PlanNode> children;
  
-   public TimeJoinNode(
-           PlanNodeId id,
-           OrderBy mergeOrder,
-           FilterNullPolicy filterNullPolicy) {
++  public TimeJoinNode(PlanNodeId id, OrderBy mergeOrder, FilterNullPolicy 
filterNullPolicy) {
 +    super(id);
 +    this.mergeOrder = mergeOrder;
 +    this.filterNullPolicy = filterNullPolicy;
 +    this.children = new ArrayList<>();
 +  }
 +
    public TimeJoinNode(
        PlanNodeId id,
        OrderBy mergeOrder,
@@@ -111,9 -89,4 +108,8 @@@
    public void setWithoutPolicy(FilterNullPolicy filterNullPolicy) {
      this.filterNullPolicy = filterNullPolicy;
    }
 +
 +  public String toString() {
 +    return "TimeJoinNode-" + this.getId();
 +  }
- 
  }
diff --cc 
server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/source/SeriesScanNode.java
index f6ab1c6,8858074..628d4da
--- 
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,12 -19,10 +19,12 @@@
  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.OrderBy;
 +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.PlanNodeAllocator;
  import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanNodeId;
  import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanVisitor;
++import org.apache.iotdb.db.sql.statement.component.OrderBy;
  import org.apache.iotdb.tsfile.read.filter.basic.Filter;
  
  import com.google.common.collect.ImmutableList;
@@@ -120,25 -105,4 +120,26 @@@ public class SeriesScanNode extends Sou
    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;
 +  }
 +
 +  public String toString() {
-     return String.format("SeriesScanNode-%s:[SeriesPath: %s, DataRegion: %s]",
-             this.getId(), this.getSeriesPath(), this.getDataRegion());
++    return String.format(
++        "SeriesScanNode-%s:[SeriesPath: %s, DataRegion: %s]",
++        this.getId(), this.getSeriesPath(), this.getDataRegion());
 +  }
  }
diff --cc 
server/src/test/java/org/apache/iotdb/db/mpp/sql/plan/DistributionPlannerTest.java
index 3a222fa,0000000..0d9c798
mode 100644,000000..100644
--- 
a/server/src/test/java/org/apache/iotdb/db/mpp/sql/plan/DistributionPlannerTest.java
+++ 
b/server/src/test/java/org/apache/iotdb/db/mpp/sql/plan/DistributionPlannerTest.java
@@@ -1,91 -1,0 +1,98 @@@
 +/*
 + * Licensed to the Apache Software Foundation (ASF) under one
 + * or more contributor license agreements.  See the NOTICE file
 + * distributed with this work for additional information
 + * regarding copyright ownership.  The ASF licenses this file
 + * to you under the Apache License, Version 2.0 (the
 + * "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
 + *
 + * Unless required by applicable law or agreed to in writing,
 + * software distributed under the License is distributed on an
 + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
 + * KIND, either express or implied.  See the License for the
 + * specific language governing permissions and limitations
 + * under the License.
 + */
 +
 +package org.apache.iotdb.db.mpp.sql.plan;
 +
 +import org.apache.iotdb.db.exception.metadata.IllegalPathException;
 +import org.apache.iotdb.db.metadata.path.PartialPath;
 +import org.apache.iotdb.db.mpp.common.*;
 +import org.apache.iotdb.db.mpp.sql.planner.plan.DistributionPlanner;
 +import org.apache.iotdb.db.mpp.sql.planner.plan.LogicalQueryPlan;
 +import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanNode;
 +import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanNodeAllocator;
- import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanNodeId;
 +import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanNodeUtil;
 +import org.apache.iotdb.db.mpp.sql.planner.plan.node.process.LimitNode;
 +import org.apache.iotdb.db.mpp.sql.planner.plan.node.process.TimeJoinNode;
 +import org.apache.iotdb.db.mpp.sql.planner.plan.node.source.SeriesScanNode;
++import org.apache.iotdb.db.sql.statement.component.OrderBy;
++
 +import org.junit.Test;
 +
 +import java.util.ArrayList;
 +import java.util.HashMap;
 +import java.util.List;
 +import java.util.Map;
 +
 +import static org.junit.Assert.assertEquals;
 +
 +public class DistributionPlannerTest {
 +
-     @Test
-     public void TestRewriteSourceNode() throws IllegalPathException {
-         TimeJoinNode timeJoinNode = new 
TimeJoinNode(PlanNodeAllocator.generateId(), OrderBy.TIMESTAMP_ASC, 
FilterNullPolicy.NO_FILTER);
- 
-         timeJoinNode.addChild(new 
SeriesScanNode(PlanNodeAllocator.generateId(), new 
PartialPath("root.sg.d1.s1")));
-         timeJoinNode.addChild(new 
SeriesScanNode(PlanNodeAllocator.generateId(), new 
PartialPath("root.sg.d1.s2")));
-         timeJoinNode.addChild(new 
SeriesScanNode(PlanNodeAllocator.generateId(), new 
PartialPath("root.sg.d2.s1")));
- 
-         LimitNode root = new LimitNode(PlanNodeAllocator.generateId(), 10, 
timeJoinNode);
- 
-         Analysis analysis = constructAnalysis();
- 
-         DistributionPlanner planner = new DistributionPlanner(analysis, new 
LogicalQueryPlan(new QueryContext(), root));
-         PlanNode newRoot = planner.rewriteSource();
- 
-         System.out.println("\nLogical-Plan:");
-         System.out.println("------------------");
-         PlanNodeUtil.printPlanNode(root);
-         System.out.println("\nDistributed-Plan:");
-         System.out.println("------------------");
-         PlanNodeUtil.printPlanNode(newRoot);
-         assertEquals(newRoot.getChildren().get(0).getChildren().size(), 3);
-         
assertEquals(newRoot.getChildren().get(0).getChildren().get(0).getChildren().size(),
 2);
-         
assertEquals(newRoot.getChildren().get(0).getChildren().get(1).getChildren().size(),
 2);
-     }
- 
-     private Analysis constructAnalysis() {
-         Analysis analysis = new Analysis();
-         Map<String, Map<DataRegionTimeSlice, List<DataRegion>>> 
dataPartitionInfo = new HashMap<>();
-         List<DataRegion> d1DataRegions = new ArrayList<>();
-         d1DataRegions.add(new DataRegion(1, "192.0.0.1"));
-         d1DataRegions.add(new DataRegion(2, "192.0.0.1"));
-         Map<DataRegionTimeSlice, List<DataRegion>> d1DataRegionMap = new 
HashMap<>();
-         d1DataRegionMap.put(new DataRegionTimeSlice(), d1DataRegions);
- 
-         List<DataRegion> d2DataRegions = new ArrayList<>();
-         d2DataRegions.add(new DataRegion(3, "192.0.0.1"));
-         Map<DataRegionTimeSlice, List<DataRegion>> d2DataRegionMap = new 
HashMap<>();
-         d2DataRegionMap.put(new DataRegionTimeSlice(), d2DataRegions);
- 
-         dataPartitionInfo.put("root.sg.d1", d1DataRegionMap);
-         dataPartitionInfo.put("root.sg.d2", d2DataRegionMap);
- 
-         analysis.setDataPartitionInfo(dataPartitionInfo);
-         return analysis;
-     }
++  @Test
++  public void TestRewriteSourceNode() throws IllegalPathException {
++    TimeJoinNode timeJoinNode =
++        new TimeJoinNode(
++            PlanNodeAllocator.generateId(), OrderBy.TIMESTAMP_ASC, 
FilterNullPolicy.NO_FILTER);
++
++    timeJoinNode.addChild(
++        new SeriesScanNode(PlanNodeAllocator.generateId(), new 
PartialPath("root.sg.d1.s1")));
++    timeJoinNode.addChild(
++        new SeriesScanNode(PlanNodeAllocator.generateId(), new 
PartialPath("root.sg.d1.s2")));
++    timeJoinNode.addChild(
++        new SeriesScanNode(PlanNodeAllocator.generateId(), new 
PartialPath("root.sg.d2.s1")));
++
++    LimitNode root = new LimitNode(PlanNodeAllocator.generateId(), 10, 
timeJoinNode);
++
++    Analysis analysis = constructAnalysis();
++
++    DistributionPlanner planner =
++        new DistributionPlanner(analysis, new LogicalQueryPlan(new 
MPPQueryContext(), root));
++    PlanNode newRoot = planner.rewriteSource();
++
++    System.out.println("\nLogical-Plan:");
++    System.out.println("------------------");
++    PlanNodeUtil.printPlanNode(root);
++    System.out.println("\nDistributed-Plan:");
++    System.out.println("------------------");
++    PlanNodeUtil.printPlanNode(newRoot);
++    assertEquals(newRoot.getChildren().get(0).getChildren().size(), 3);
++    
assertEquals(newRoot.getChildren().get(0).getChildren().get(0).getChildren().size(),
 2);
++    
assertEquals(newRoot.getChildren().get(0).getChildren().get(1).getChildren().size(),
 2);
++  }
++
++  private Analysis constructAnalysis() {
++    Analysis analysis = new Analysis();
++    Map<String, Map<DataRegionTimeSlice, List<DataRegion>>> dataPartitionInfo 
= new HashMap<>();
++    List<DataRegion> d1DataRegions = new ArrayList<>();
++    d1DataRegions.add(new DataRegion(1, "192.0.0.1"));
++    d1DataRegions.add(new DataRegion(2, "192.0.0.1"));
++    Map<DataRegionTimeSlice, List<DataRegion>> d1DataRegionMap = new 
HashMap<>();
++    d1DataRegionMap.put(new DataRegionTimeSlice(), d1DataRegions);
++
++    List<DataRegion> d2DataRegions = new ArrayList<>();
++    d2DataRegions.add(new DataRegion(3, "192.0.0.1"));
++    Map<DataRegionTimeSlice, List<DataRegion>> d2DataRegionMap = new 
HashMap<>();
++    d2DataRegionMap.put(new DataRegionTimeSlice(), d2DataRegions);
++
++    dataPartitionInfo.put("root.sg.d1", d1DataRegionMap);
++    dataPartitionInfo.put("root.sg.d2", d2DataRegionMap);
++
++    analysis.setDataPartitionInfo(dataPartitionInfo);
++    return analysis;
++  }
 +}

Reply via email to