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

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

commit 365852c39fbee5bf6ce7f3e3b3449d1293b1fa16
Merge: b877352 b3fff9f
Author: JackieTien97 <[email protected]>
AuthorDate: Fri Mar 18 18:52:14 2022 +0800

    format code and merge master

 .github/workflows/client-go.yml                    |    4 +
 .github/workflows/client.yml                       |    4 +
 .github/workflows/cluster.yml                      |    4 +
 .github/workflows/e2e.yml                          |    4 +
 .github/workflows/grafana-plugin.yml               |    5 +
 .github/workflows/greetings.yml                    |    4 +
 .github/workflows/influxdb-protocol.yml            |    4 +
 .github/workflows/main-unix.yml                    |    4 +
 .github/workflows/main-win.yml                     |    4 +
 .github/workflows/sonar-coveralls.yml              |    4 +
 .../org/apache/iotdb/db/qp/sql/IoTDBSqlParser.g4   |    8 +-
 .../antlr4/org/apache/iotdb/db/qp/sql/SqlLexer.g4  |   11 +-
 client-cpp/src/main/Session.h                      |    3 +-
 client-py/iotdb/utils/IoTDBConstants.py            |    1 +
 .../org/apache/iotdb/cluster/ClusterIoTDB.java     |   59 +-
 .../cluster/ClusterIoTDBServerCommandLine.java     |   94 ++
 .../cluster/client/async/AsyncDataClient.java      |    2 +-
 .../cluster/client/async/AsyncMetaClient.java      |    2 +-
 .../iotdb/cluster/client/sync/SyncDataClient.java  |    2 +-
 .../iotdb/cluster/client/sync/SyncMetaClient.java  |    2 +-
 .../iotdb/cluster/config/ClusterConstant.java      |    2 +-
 .../iotdb/cluster/coordinator/Coordinator.java     |   12 +-
 .../apache/iotdb/cluster/log/LogDispatcher.java    |    2 +-
 .../cluster/log/applier/AsyncDataLogApplier.java   |    8 +-
 .../iotdb/cluster/log/applier/BaseApplier.java     |    2 +-
 .../iotdb/cluster/log/applier/DataLogApplier.java  |    6 +-
 .../iotdb/cluster/log/catchup/CatchUpTask.java     |    2 +-
 .../iotdb/cluster/log/catchup/LogCatchUpTask.java  |    2 +-
 .../cluster/log/manage/CommittedEntryManager.java  |    2 +-
 .../log/manage/MetaSingleSnapshotLogManager.java   |    2 +-
 .../log/manage/PartitionedSnapshotLogManager.java  |    4 +-
 .../iotdb/cluster/log/manage/RaftLogManager.java   |    2 +-
 .../log/manage/UnCommittedEntryManager.java        |    2 +-
 .../serializable/SyncLogDequeSerializer.java       |    2 +-
 .../cluster/log/snapshot/MetaSimpleSnapshot.java   |    4 +-
 .../{CMManager.java => CSchemaEngine.java}         |   20 +-
 .../apache/iotdb/cluster/metadata/MetaPuller.java  |   10 +-
 .../iotdb/cluster/partition/PartitionTable.java    |    4 +-
 .../cluster/query/ClusterPhysicalGenerator.java    |    8 +-
 .../iotdb/cluster/query/ClusterPlanExecutor.java   |   24 +-
 .../iotdb/cluster/query/ClusterPlanRouter.java     |   18 +-
 .../iotdb/cluster/query/LocalQueryExecutor.java    |   31 +-
 .../iotdb/cluster/query/filter/SlotSgFilter.java   |    2 +-
 .../cluster/query/reader/ClusterTimeGenerator.java |    6 +-
 .../cluster/server/member/DataGroupMember.java     |    8 +-
 .../cluster/server/member/MetaGroupMember.java     |    4 +-
 .../iotdb/cluster/server/member/RaftMember.java    |    2 +-
 .../cluster/server/monitor/NodeStatusManager.java  |    2 +-
 .../cluster/server/service/DataAsyncService.java   |   14 +-
 .../cluster/server/service/DataGroupEngine.java    |    2 +-
 .../cluster/server/service/DataSyncService.java    |   12 +-
 .../iotdb/cluster/utils/ClusterQueryUtils.java     |    2 +-
 .../apache/iotdb/cluster/utils/ClusterUtils.java   |    4 +-
 .../apache/iotdb/cluster/utils/PartitionUtils.java |    2 -
 .../log/applier/AsyncDataLogApplierTest.java       |    6 +-
 .../cluster/log/applier/DataLogApplierTest.java    |   12 +-
 .../cluster/log/applier/MetaLogApplierTest.java    |   16 +-
 .../iotdb/cluster/log/catchup/CatchUpTaskTest.java |    4 +-
 .../cluster/log/snapshot/DataSnapshotTest.java     |    2 +-
 .../cluster/log/snapshot/FileSnapshotTest.java     |    8 +-
 .../log/snapshot/MetaSimpleSnapshotTest.java       |    4 +-
 .../log/snapshot/PartitionedSnapshotTest.java      |    5 +-
 .../cluster/log/snapshot/PullSnapshotTaskTest.java |    2 +-
 ...agerWhiteBox.java => SchemaEngineWhiteBox.java} |   20 +-
 .../cluster/partition/SlotPartitionTableTest.java  |   28 +-
 .../cluster/query/ClusterPlanExecutorTest.java     |    2 +-
 .../clusterinfo/ClusterInfoServiceImplTest.java    |    4 +-
 .../iotdb/cluster/server/member/BaseMember.java    |   10 +-
 .../cluster/server/member/DataGroupMemberTest.java |    4 +-
 .../cluster/server/member/MetaGroupMemberTest.java |   22 +-
 confignode/pom.xml                                 |    6 -
 .../assembly/resources/sbin/start-confignode.bat   |    3 +-
 .../assembly/resources/sbin/start-confignode.sh    |    2 +-
 .../iotdb/confignode/conf/ConfigNodeConfCheck.java |    8 +-
 .../iotdb/confignode/manager/ConfigManager.java    |    2 +-
 .../iotdb/confignode/service/ConfigNode.java       |   33 +-
 .../confignode/service/ConfigNodeCommandLine.java  |   80 +
 .../utils/ConfigNodeEnvironmentUtils.java          |    2 +-
 .../org/apache/iotdb/consensus/IConsensus.java     |    1 +
 .../common/request/ByteBufferConsensusRequest.java |   30 +-
 .../common/request/IConsensusRequest.java          |    2 -
 .../consensus/standalone/StandAloneServerImpl.java |    8 +-
 .../consensus/statemachine/EmptyStateMachine.java  |    2 +-
 .../standalone/StandAloneConsensusTest.java        |   30 +-
 docs/UserGuide/API/Programming-Java-Native-API.md  |    1 +
 docs/UserGuide/Data-Concept/Encoding.md            |   14 +-
 .../Maintenance-Tools/Maintenance-Command.md       |    8 -
 docs/UserGuide/Operate-Metadata/Timeseries.md      |    2 +-
 docs/UserGuide/QuickStart/WayToGetIoTDB.md         |   17 +-
 docs/UserGuide/Reference/Config-Manual.md          |   18 +
 docs/UserGuide/Reference/SQL-Reference.md          |    5 -
 .../UserGuide/API/Programming-Java-Native-API.md   |    3 +-
 docs/zh/UserGuide/Data-Concept/Encoding.md         |   14 +-
 .../Maintenance-Tools/Maintenance-Command.md       |    7 -
 docs/zh/UserGuide/Operate-Metadata/Timeseries.md   |    2 +-
 docs/zh/UserGuide/QuickStart/WayToGetIoTDB.md      |   17 +-
 docs/zh/UserGuide/Reference/Config-Manual.md       |   22 +-
 docs/zh/UserGuide/Reference/SQL-Reference.md       |    6 -
 .../java/org/apache/iotdb/SessionPoolExample.java  |   42 +-
 ...=> RowTSRecordOutputFormatIntegrationTest.java} |    2 +-
 ...va => RowTsFileInputFormatIntegrationTest.java} |   58 +-
 .../util/TSFileConfigUtilCompletenessTest.java     |    4 +-
 .../iotdb/db/integration/IoTDBArithmeticIT.java    |   18 +-
 .../iotdb/db/integration/IoTDBCheckConfigIT.java   |    4 +-
 .../integration/IoTDBCompactionWithIDTableIT.java  |  352 ++++
 .../db/integration/IoTDBCreateSnapshotIT.java      |  180 --
 .../iotdb/db/integration/IoTDBEncodingIT.java      |   76 +
 .../apache/iotdb/db/integration/IoTDBLastIT.java   |   14 +-
 .../iotdb/db/integration/IoTDBMetadataFetchIT.java |   45 +-
 .../iotdb/db/integration/IoTDBNestedQueryIT.java   |   12 +-
 .../iotdb/db/integration/IoTDBSelectIntoIT.java    |   18 +-
 .../iotdb/db/integration/IoTDBSimpleQueryIT.java   |    8 +-
 .../db/integration/IoTDBTriggerExecutionIT.java    |   26 +-
 .../db/integration/IoTDBTriggerManagementIT.java   |    8 +-
 .../iotdb/db/integration/IoTDBUDFManagementIT.java |    6 +-
 .../apache/iotdb/session/IoTDBSessionSimpleIT.java |    4 +-
 iotdb-commons/pom.xml                              |  127 ++
 .../apache/iotdb/commons/ServerCommandLine.java    |   67 +
 .../org/apache/iotdb/commons}/utils/TestOnly.java  |    2 +-
 pom.xml                                            |    2 +
 .../resources/conf/iotdb-engine.properties         |   31 +-
 .../java/org/apache/iotdb/db/conf/IoTDBConfig.java |   54 +-
 .../org/apache/iotdb/db/conf/IoTDBConfigCheck.java |   24 +-
 .../org/apache/iotdb/db/conf/IoTDBDescriptor.java  |   44 +-
 .../{ConsensusMain.java => ConsensusExample.java}  |   32 +-
 .../consensus/statemachine/BaseStateMachine.java   |   74 +
 .../DataRegionStateMachine.java}                   |   23 +-
 .../SchemaRegionStateMachine.java}                 |   23 +-
 .../org/apache/iotdb/db/engine/StorageEngine.java  |   22 +-
 .../iotdb/db/engine/cache/BloomFilterCache.java    |    2 +-
 .../apache/iotdb/db/engine/cache/ChunkCache.java   |    2 +-
 .../db/engine/cache/TimeSeriesMetadataCache.java   |    2 +-
 .../engine/compaction/CompactionTaskManager.java   |    2 +-
 .../db/engine/compaction/CompactionUtils.java      |   19 +-
 .../manage/CrossSpaceCompactionResource.java       |    6 -
 .../inner/utils/InnerSpaceCompactionUtils.java     |    8 +-
 .../engine/cq/ContinuousQuerySchemaCheckTask.java  |    2 +-
 .../iotdb/db/engine/cq/ContinuousQueryService.java |    2 +-
 .../db/engine/storagegroup/StorageGroupInfo.java   |    2 +-
 .../db/engine/storagegroup/TsFileProcessor.java    |    2 +-
 .../db/engine/storagegroup/TsFileResource.java     |    2 +-
 .../db/engine/storagegroup/TsFileResourceList.java |    2 +-
 .../storagegroup/VirtualStorageGroupProcessor.java |   16 +-
 .../engine/trigger/executor/TriggerExecutor.java   |    2 +-
 .../service/TriggerRegistrationService.java        |    4 +-
 .../trigger/sink/local/LocalIoTDBHandler.java      |    6 +-
 .../SchemaDirCreationFailureException.java         |   13 +-
 .../db/metadata/IStorageGroupSchemaManager.java    |  232 +++
 .../apache/iotdb/db/metadata/MetadataConstant.java |    3 +
 .../org/apache/iotdb/db/metadata/SchemaEngine.java | 1736 ++++++++++++++++++++
 .../metadata/{MManager.java => SchemaRegion.java}  | 1238 +++-----------
 .../db/metadata/StorageGroupSchemaManager.java     |  251 +++
 .../idtable/AppendOnlyDiskSchemaManager.java       |   41 +-
 .../apache/iotdb/db/metadata/idtable/IDTable.java  |   12 +-
 .../db/metadata/idtable/IDTableHashmapImpl.java    |   41 +-
 .../iotdb/db/metadata/idtable/IDTableManager.java  |   23 +-
 .../db/metadata/idtable/IDiskSchemaManager.java    |    2 +-
 .../db/metadata/idtable/entry/DeviceEntry.java     |    2 +-
 .../db/metadata/idtable/entry/DeviceIDFactory.java |    2 +-
 .../idtable/entry/InsertMeasurementMNode.java      |    9 +-
 .../db/metadata/idtable/entry/SchemaEntry.java     |    2 +-
 .../db/metadata/lastCache/LastCacheManager.java    |    6 +-
 .../iotdb/db/metadata/logfile/MLogReader.java      |    2 +-
 .../iotdb/db/metadata/logfile/MLogTxtReader.java   |    2 +-
 .../iotdb/db/metadata/logfile/MLogUpgrader.java    |  290 ----
 .../iotdb/db/metadata/mnode/EntityMNode.java       |   13 +
 .../org/apache/iotdb/db/metadata/mnode/IMNode.java |    7 +-
 .../db/metadata/mnode/IStorageGroupMNode.java      |    6 +
 .../iotdb/db/metadata/mnode/InternalMNode.java     |   38 +-
 .../org/apache/iotdb/db/metadata/mnode/MNode.java  |   14 +-
 .../iotdb/db/metadata/mnode/MeasurementMNode.java  |    3 +-
 .../db/metadata/mnode/StorageGroupEntityMNode.java |   23 +
 .../iotdb/db/metadata/mnode/StorageGroupMNode.java |   23 +
 .../iotdb/db/metadata/mtree/MTreeAboveSG.java      |  559 +++++++
 .../mtree/{MTree.java => MTreeBelowSG.java}        | 1102 +++----------
 .../db/metadata/mtree/traverser/Traverser.java     |   86 +-
 .../MNodeAboveSGCollector.java}                    |   48 +-
 .../mtree/traverser/collector/MNodeCollector.java  |    2 +-
 .../traverser/collector/MeasurementCollector.java  |   20 +-
 ...lCounter.java => MNodeAboveSGLevelCounter.java} |   47 +-
 .../mtree/traverser/counter/MNodeLevelCounter.java |   10 +-
 .../counter/MeasurementGroupByLevelCounter.java    |   26 +
 .../apache/iotdb/db/metadata/path/AlignedPath.java |    2 +-
 .../iotdb/db/metadata/path/MeasurementPath.java    |    2 +-
 .../apache/iotdb/db/metadata/path/PartialPath.java |    2 +-
 .../db/metadata/rescon/TimeseriesStatistics.java   |  104 ++
 .../apache/iotdb/db/metadata/tag/TagManager.java   |   27 +-
 .../template/TemplateLogReader.java}               |   30 +-
 .../db/metadata/template/TemplateLogWriter.java    |   64 +
 .../db/metadata/template/TemplateManager.java      |  123 +-
 .../db/metadata/upgrade/MetadataUpgrader.java      |  429 +++++
 .../apache/iotdb/db/metadata/utils/MetaUtils.java  |    2 +-
 .../org/apache/iotdb/db/mpp/common/InstanceId.java |    8 +-
 .../iotdb/db/mpp/execution/InstanceContext.java    |   28 +-
 .../iotdb/db/mpp/execution/InstanceState.java      |   77 +-
 .../mpp/execution/executor/InstanceExecutor.java   |   25 +-
 .../db/mpp/execution/executor/InstanceHandle.java  |    9 +-
 .../db/mpp/operator/process/AggregateOperator.java |   51 +-
 .../mpp/operator/process/DeviceMergeOperator.java  |   43 +-
 .../db/mpp/operator/process/FillOperator.java      |   43 +-
 .../mpp/operator/process/FilterNullOperator.java   |   51 +-
 .../mpp/operator/process/GroupByLevelOperator.java |   51 +-
 .../db/mpp/operator/process/LimitOperator.java     |   43 +-
 .../db/mpp/operator/process/OffsetOperator.java    |   51 +-
 .../db/mpp/operator/process/ProcessOperator.java   |    4 +-
 .../db/mpp/operator/process/SortOperator.java      |   51 +-
 .../db/mpp/operator/process/TimeJoinOperator.java  |   51 +-
 .../iotdb/db/mpp/operator/sink/SinkOperator.java   |   30 +-
 .../source/SeriesAggregateScanOperator.java        |   61 +-
 .../db/mpp/operator/source/SeriesScanOperator.java |   61 +-
 .../db/mpp/operator/source/SourceOperator.java     |    2 +-
 .../db/mpp/sql/planner/LocalExecutionPlanner.java  |  168 +-
 .../mpp/sql/planner/plan/DistributedQueryPlan.java |    1 -
 .../db/mpp/sql/planner/plan/LogicalQueryPlan.java  |    1 -
 .../db/mpp/sql/planner/plan/PlanFragment.java      |    1 -
 .../db/mpp/sql/planner/plan/node/PlanNode.java     |    5 +-
 .../db/mpp/sql/planner/plan/node/PlanNodeId.java   |   34 +-
 .../db/mpp/sql/planner/plan/node/PlanVisitor.java  |   74 +-
 .../planner/plan/node/process/AggregateNode.java   |   14 +-
 .../planner/plan/node/process/DeviceMergeNode.java |    2 +-
 .../sql/planner/plan/node/process/FillNode.java    |    3 +-
 .../sql/planner/plan/node/process/FilterNode.java  |    6 +-
 .../planner/plan/node/process/FilterNullNode.java  |    5 +-
 .../plan/node/process/GroupByLevelNode.java        |    3 +-
 .../sql/planner/plan/node/process/LimitNode.java   |    3 +-
 .../sql/planner/plan/node/process/ProcessNode.java |    1 -
 .../sql/planner/plan/node/process/SortNode.java    |    3 +-
 .../planner/plan/node/process/TimeJoinNode.java    |   11 +-
 .../mpp/sql/planner/plan/node/sink/SinkNode.java   |    1 -
 .../plan/node/source/SeriesAggregateScanNode.java  |    3 +-
 .../planner/plan/node/source/SeriesScanNode.java   |    3 +-
 .../sql/planner/plan/node/source/SourceNode.java   |    1 -
 .../apache/iotdb/db/mpp/sql/tree/Expression.java   |    4 +-
 .../apache/iotdb/db/qp/executor/PlanExecutor.java  |   97 +-
 .../iotdb/db/qp/logical/crud/QueryOperator.java    |    4 +-
 .../sys/CreateAlignedTimeSeriesOperator.java       |    2 +-
 .../apache/iotdb/db/qp/physical/PhysicalPlan.java  |   15 +-
 .../iotdb/db/qp/physical/crud/InsertPlan.java      |    2 +-
 .../iotdb/db/qp/physical/crud/InsertRowPlan.java   |    2 +-
 .../iotdb/db/qp/physical/crud/QueryPlan.java       |    2 +-
 .../physical/sys/CreateAlignedTimeSeriesPlan.java  |   82 +-
 .../db/qp/physical/sys/CreateSnapshotPlan.java     |   56 -
 .../db/qp/physical/sys/CreateTemplatePlan.java     |    2 +-
 .../apache/iotdb/db/qp/sql/IoTDBSqlVisitor.java    |   16 +-
 .../iotdb/db/qp/strategy/LogicalGenerator.java     |    2 +-
 .../qp/strategy/optimizer/ConcatPathOptimizer.java |    2 +-
 .../apache/iotdb/db/qp/utils/DatetimeUtils.java    |    2 +-
 .../apache/iotdb/db/qp/utils/WildcardsRemover.java |    4 +-
 .../iotdb/db/query/dataset/ShowDevicesDataSet.java |    2 +-
 .../db/query/dataset/ShowTimeseriesDataSet.java    |    2 +-
 .../query/dataset/groupby/GroupByTimeDataSet.java  |    2 +-
 .../groupby/GroupByWithValueFilterDataSet.java     |    2 +-
 .../iotdb/db/query/executor/LastQueryExecutor.java |   16 +-
 .../query/reader/series/AlignedSeriesReader.java   |    2 +-
 .../query/reader/series/SeriesAggregateReader.java |    2 +-
 .../reader/series/SeriesRawDataBatchReader.java    |    2 +-
 .../iotdb/db/query/reader/series/SeriesReader.java |    2 +-
 .../reader/series/SeriesReaderByTimestamp.java     |    2 +-
 .../query/udf/service/UDFRegistrationService.java  |    2 +-
 .../iotdb/db/rescon/TsFileResourceManager.java     |    2 +-
 .../java/org/apache/iotdb/db/service/IoTDB.java    |   14 +-
 .../apache/iotdb/db/service/RegisterManager.java   |    2 +-
 .../db/service/thrift/impl/TSServiceImpl.java      |   20 +-
 .../db/sync/receiver/transfer/SyncServiceImpl.java |    2 +-
 .../db/sync/sender/manage/SyncFileManager.java     |    2 +-
 .../db/tools/virtualsg/DeviceMappingViewer.java    |   12 +-
 .../org/apache/iotdb/db/utils/CommonUtils.java     |    1 +
 .../apache/iotdb/db/utils/EnvironmentUtils.java    |    4 +-
 .../org/apache/iotdb/db/utils/SchemaTestUtils.java |    2 +-
 .../org/apache/iotdb/db/utils/SchemaUtils.java     |    5 +-
 .../db/utils/datastructure/AlignedTVList.java      |    2 +-
 .../iotdb/db/utils/datastructure/TVList.java       |    2 +-
 .../utils/windowing/window/EvictableBatchList.java |    2 +-
 .../org/apache/iotdb/db/writelog/io/LogWriter.java |    2 +-
 .../iotdb/db/writelog/recover/LogReplayer.java     |    4 +-
 .../iotdb/db/engine/MetadataManagerHelper.java     |   48 +-
 .../iotdb/db/engine/cache/ChunkCacheTest.java      |    8 +-
 .../engine/compaction/AbstractCompactionTest.java  |   10 +-
 .../engine/compaction/CompactionSchedulerTest.java |   64 +-
 .../compaction/TestUtilsForAlignedSeries.java      |    6 +-
 .../compaction/cross/CrossSpaceCompactionTest.java |    8 +-
 .../db/engine/compaction/cross/MergeTest.java      |    8 +-
 .../inner/AbstractInnerSpaceCompactionTest.java    |    8 +-
 .../inner/InnerCompactionMoreDataTest.java         |    4 +-
 .../compaction/inner/InnerCompactionTest.java      |    8 +-
 .../compaction/inner/InnerSeqCompactionTest.java   |    8 +-
 .../InnerSpaceCompactionUtilsAlignedTest.java      |    4 +-
 .../InnerSpaceCompactionUtilsNoAlignedTest.java    |    6 +-
 .../compaction/inner/InnerUnseqCompactionTest.java |    8 +-
 .../inner/sizetiered/SizeTieredCompactionTest.java |    8 +-
 .../recover/SizeTieredCompactionRecoverTest.java   |   12 +-
 .../engine/modification/DeletionFileNodeTest.java  |    4 +-
 .../db/engine/modification/DeletionQueryTest.java  |    4 +-
 .../storagegroup/FileNodeManagerBenchmark.java     |    8 +-
 .../iotdb/db/engine/storagegroup/TTLTest.java      |   16 +-
 .../org/apache/iotdb/db/metadata/MTreeTest.java    | 1060 ------------
 ...ncedTest.java => SchemaEngineAdvancedTest.java} |   72 +-
 ...erBasicTest.java => SchemaEngineBasicTest.java} | 1004 ++++++-----
 ...proveTest.java => SchemaEngineImproveTest.java} |   36 +-
 .../org/apache/iotdb/db/metadata/TemplateTest.java |  112 +-
 .../iotdb/db/metadata/idtable/IDTableTest.java     |   66 +-
 .../db/metadata/idtable/InsertWithIDTableTest.java |   18 +-
 .../iotdb/db/metadata/mlog/MLogUpgraderTest.java   |  176 --
 .../iotdb/db/metadata/mtree/MTreeAboveSGTest.java  |  292 ++++
 .../iotdb/db/metadata/mtree/MTreeBelowSGTest.java  |  796 +++++++++
 .../db/metadata/upgrade/MetadataUpgradeTest.java   |  304 ++++
 .../java/org/apache/iotdb/db/qp/PlannerTest.java   |   34 +-
 .../iotdb/db/qp/logical/LogicalPlanSmallTest.java  |    4 +-
 .../iotdb/db/qp/physical/ConcatOptimizerTest.java  |   18 +-
 .../iotdb/db/qp/physical/InsertRowPlanTest.java    |   12 +-
 .../iotdb/db/qp/physical/InsertTabletPlanTest.java |   10 +-
 .../iotdb/db/qp/physical/PhysicalPlanTest.java     |   12 +-
 .../iotdb/db/qp/physical/SerializationTest.java    |   14 +-
 .../dataset/EngineDataSetWithValueFilterTest.java  |    2 +-
 .../query/dataset/UDTFAlignByTimeDataSetTest.java  |   14 +-
 .../query/dataset/groupby/GroupByDataSetTest.java  |    2 +-
 .../dataset/groupby/GroupByFillDataSetTest.java    |    2 +-
 .../dataset/groupby/GroupByLevelDataSetTest.java   |    2 +-
 .../query/reader/series/SeriesReaderTestUtil.java  |    8 +-
 .../iotdb/db/rescon/ResourceManagerTest.java       |    8 +-
 .../db/sync/receiver/load/FileLoaderTest.java      |   12 +-
 .../recover/SyncReceiverLogAnalyzerTest.java       |   12 +-
 .../db/sync/sender/manage/SyncFileManagerTest.java |    2 +-
 .../sender/recover/SyncSenderLogAnalyzerTest.java  |    2 +-
 .../org/apache/iotdb/db/tools/MLogParserTest.java  |  112 +-
 .../org/apache/iotdb/db/utils/SchemaUtilsTest.java |    8 +-
 .../apache/iotdb/db/writelog/PerformanceTest.java  |   10 +-
 .../db/writelog/recover/DeviceStringTest.java      |   12 +-
 .../iotdb/db/writelog/recover/LogReplayerTest.java |    4 +-
 .../recover/RecoverResourceFromReaderTest.java     |    8 +-
 .../db/writelog/recover/SeqTsFileRecoverTest.java  |    8 +-
 .../writelog/recover/UnseqTsFileRecoverTest.java   |    8 +-
 .../java/org/apache/iotdb/session/Session.java     |   88 +-
 .../org/apache/iotdb/session/pool/SessionPool.java |  178 +-
 .../apache/iotdb/spark/db/EnvironmentUtils.java    |    4 +-
 thrift-confignode/pom.xml                          |    2 +-
 tsfile/pom.xml                                     |   14 +-
 .../iotdb/tsfile/common/conf/TSFileConfig.java     |   20 +
 .../iotdb/tsfile/common/conf/TSFileDescriptor.java |    6 +
 .../iotdb/tsfile/encoding/decoder/Decoder.java     |    3 +-
 .../iotdb/tsfile/encoding/decoder/FreqDecoder.java |  140 ++
 .../iotdb/tsfile/encoding/encoder/FreqEncoder.java |  313 ++++
 .../tsfile/encoding/encoder/TSEncodingBuilder.java |   65 +-
 .../tsfile/file/metadata/enums/TSEncoding.java     |    5 +-
 .../apache/iotdb/tsfile/utils/BitConstructor.java  |   93 ++
 .../org/apache/iotdb/tsfile/utils/BitReader.java   |   70 +
 .../iotdb/tsfile/utils/MeasurementGroup.java       |    3 +-
 .../tsfile/encoding/decoder/FreqDecoderTest.java   |  161 ++
 348 files changed, 10201 insertions(+), 6131 deletions(-)

diff --cc server/src/main/java/org/apache/iotdb/db/mpp/common/InstanceId.java
index a3dd91f,7a4107e..53ed39a
--- a/server/src/main/java/org/apache/iotdb/db/mpp/common/InstanceId.java
+++ b/server/src/main/java/org/apache/iotdb/db/mpp/common/InstanceId.java
@@@ -18,11 -18,7 +18,11 @@@
   */
  package org.apache.iotdb.db.mpp.common;
  
 -public enum WithoutPolicy {
 -  CONTAINS_NULL,
 -  ALL_NULL
 +public class InstanceId {
 +
-     private final String fullId;
++  private final String fullId;
 +
-     public InstanceId(String fullId) {
-         this.fullId = fullId;
-     }
++  public InstanceId(String fullId) {
++    this.fullId = fullId;
++  }
  }
diff --cc 
server/src/main/java/org/apache/iotdb/db/mpp/execution/InstanceContext.java
index b6d27f5,e5a69e9..3a5f99d
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/execution/InstanceContext.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/execution/InstanceContext.java
@@@ -16,29 -16,21 +16,25 @@@
   * specific language governing permissions and limitations
   * under the License.
   */
 -package org.apache.iotdb.db.mpp.plan.node;
 +package org.apache.iotdb.db.mpp.execution;
  
 -public class PlanNodeId {
 -  private String id;
 +import org.apache.iotdb.db.mpp.common.InstanceId;
  
- import java.util.concurrent.atomic.AtomicLong;
- import java.util.concurrent.atomic.AtomicReference;
- 
 -  public PlanNodeId(String id) {
 -    this.id = id;
 -  }
 +public class InstanceContext {
  
-     private InstanceId id;
- 
-     private final long createNanos = System.nanoTime();
 -  public String getId() {
 -    return this.id;
 -  }
++  private InstanceId id;
  
- //    private final GcMonitor gcMonitor;
- //    private final AtomicLong startNanos = new AtomicLong();
- //    private final AtomicLong startFullGcCount = new AtomicLong(-1);
- //    private final AtomicLong startFullGcTimeNanos = new AtomicLong(-1);
- //    private final AtomicLong endNanos = new AtomicLong();
- //    private final AtomicLong endFullGcCount = new AtomicLong(-1);
- //    private final AtomicLong endFullGcTimeNanos = new AtomicLong(-1);
 -  @Override
 -  public String toString() {
 -    return this.id;
++  private final long createNanos = System.nanoTime();
 +
++  //    private final GcMonitor gcMonitor;
++  //    private final AtomicLong startNanos = new AtomicLong();
++  //    private final AtomicLong startFullGcCount = new AtomicLong(-1);
++  //    private final AtomicLong startFullGcTimeNanos = new AtomicLong(-1);
++  //    private final AtomicLong endNanos = new AtomicLong();
++  //    private final AtomicLong endFullGcCount = new AtomicLong(-1);
++  //    private final AtomicLong endFullGcTimeNanos = new AtomicLong(-1);
 +
-     public InstanceContext(InstanceId id) {
-         this.id = id;
-     }
++  public InstanceContext(InstanceId id) {
++    this.id = id;
+   }
  }
diff --cc 
server/src/main/java/org/apache/iotdb/db/mpp/execution/InstanceState.java
index 847a6e7,0000000..202fc59
mode 100644,000000..100644
--- a/server/src/main/java/org/apache/iotdb/db/mpp/execution/InstanceState.java
+++ b/server/src/main/java/org/apache/iotdb/db/mpp/execution/InstanceState.java
@@@ -1,77 -1,0 +1,62 @@@
 +/*
 + * 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.execution;
 +
 +import java.util.Set;
 +import java.util.stream.Stream;
 +
 +import static com.google.common.collect.ImmutableSet.toImmutableSet;
 +
 +public enum InstanceState {
-     /**
-      * Instance is planned but has not been scheduled yet. An instance will
-      * be in the planned state until, the dependencies of the instance
-      * have begun producing output.
-      */
-     PLANNED(false),
-     /**
-      * Instance is running.
-      */
-     RUNNING(false),
-     /**
-      * Instance has finished executing and output is left to be consumed.
-      * In this state, there will be no new drivers, the existing drivers have 
finished
-      * and the output buffer of the instance is at-least in a 
'no-more-tsBlocks' state.
-      */
-     FLUSHING(false),
-     /**
-      * Instance has finished executing and all output has been consumed.
-      */
-     FINISHED(true),
-     /**
-      * Instance was canceled by a user.
-      */
-     CANCELED(true),
-     /**
-      * Instance was aborted due to a failure in the query. The failure
-      * was not in this instance.
-      */
-     ABORTED(true),
-     /**
-      * Instance execution failed.
-      */
-     FAILED(true);
++  /**
++   * Instance is planned but has not been scheduled yet. An instance will be 
in the planned state
++   * until, the dependencies of the instance have begun producing output.
++   */
++  PLANNED(false),
++  /** Instance is running. */
++  RUNNING(false),
++  /**
++   * Instance has finished executing and output is left to be consumed. In 
this state, there will be
++   * no new drivers, the existing drivers have finished and the output buffer 
of the instance is
++   * at-least in a 'no-more-tsBlocks' state.
++   */
++  FLUSHING(false),
++  /** Instance has finished executing and all output has been consumed. */
++  FINISHED(true),
++  /** Instance was canceled by a user. */
++  CANCELED(true),
++  /** Instance was aborted due to a failure in the query. The failure was not 
in this instance. */
++  ABORTED(true),
++  /** Instance execution failed. */
++  FAILED(true);
 +
-     public static final Set<InstanceState> TERMINAL_TASK_STATES = 
Stream.of(InstanceState.values()).filter(InstanceState::isDone).collect(toImmutableSet());
++  public static final Set<InstanceState> TERMINAL_TASK_STATES =
++      
Stream.of(InstanceState.values()).filter(InstanceState::isDone).collect(toImmutableSet());
 +
-     private final boolean doneState;
++  private final boolean doneState;
 +
-     InstanceState(boolean doneState)
-     {
-         this.doneState = doneState;
-     }
++  InstanceState(boolean doneState) {
++    this.doneState = doneState;
++  }
 +
-     /**
-      * Is this a terminal state.
-      */
-     public boolean isDone()
-     {
-         return doneState;
-     }
++  /** Is this a terminal state. */
++  public boolean isDone() {
++    return doneState;
++  }
 +}
diff --cc 
server/src/main/java/org/apache/iotdb/db/mpp/execution/executor/InstanceExecutor.java
index 8e819ff,9954c74..76909e2
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/execution/executor/InstanceExecutor.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/execution/executor/InstanceExecutor.java
@@@ -16,23 -16,19 +16,24 @@@
   * specific language governing permissions and limitations
   * under the License.
   */
 -package org.apache.iotdb.db.mpp.plan;
 +package org.apache.iotdb.db.mpp.execution.executor;
  
- import com.google.common.util.concurrent.ListenableFuture;
- import com.google.common.util.concurrent.SettableFuture;
 -import org.apache.iotdb.db.mpp.common.QueryContext;
 -import org.apache.iotdb.db.mpp.common.TsBlock;
 -import org.apache.iotdb.db.mpp.plan.node.PlanNode;
 +import org.apache.iotdb.db.mpp.execution.ExecFragmentInstance;
  
 -import java.util.List;
++import com.google.common.util.concurrent.ListenableFuture;
++import com.google.common.util.concurrent.SettableFuture;
  
 -public class DistributedQueryPlan {
 -  private QueryContext context;
 -  private PlanNode<TsBlock> rootNode;
 -  private PlanFragment rootFragment;
 +public class InstanceExecutor {
  
-     /**
-      * TODO native implementation, should be replaced later
-      * @param instance executable fragment instance
-      * @param handle instance handle
-      * @return ListenableFuture indicate the instance's end state
-      */
-     public ListenableFuture<Void> enqueueInstance(ExecFragmentInstance 
instance, InstanceHandle handle) {
-         return SettableFuture.create();
-     }
- 
 -  // TODO: consider whether this field is necessary when do the implementation
 -  private List<PlanFragment> fragments;
++  /**
++   * TODO native implementation, should be replaced later
++   *
++   * @param instance executable fragment instance
++   * @param handle instance handle
++   * @return ListenableFuture indicate the instance's end state
++   */
++  public ListenableFuture<Void> enqueueInstance(
++      ExecFragmentInstance instance, InstanceHandle handle) {
++    return SettableFuture.create();
++  }
  }
diff --cc 
server/src/main/java/org/apache/iotdb/db/mpp/execution/executor/InstanceHandle.java
index 08444e7,e5a69e9..cdacf5a
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/execution/executor/InstanceHandle.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/execution/executor/InstanceHandle.java
@@@ -16,17 -16,21 +16,16 @@@
   * specific language governing permissions and limitations
   * under the License.
   */
 -package org.apache.iotdb.db.mpp.plan.node;
 +package org.apache.iotdb.db.mpp.execution.executor;
  
 -public class PlanNodeId {
 -  private String id;
 +import org.apache.iotdb.db.mpp.common.InstanceId;
  
- 
 -  public PlanNodeId(String id) {
 -    this.id = id;
 -  }
 +// TODO should contain more fields, add as you want
 +public class InstanceHandle {
  
-     private final InstanceId taskId;
 -  public String getId() {
 -    return this.id;
 -  }
++  private final InstanceId taskId;
  
-     public InstanceHandle(InstanceId taskId) {
-         this.taskId = taskId;
-     }
 -  @Override
 -  public String toString() {
 -    return this.id;
++  public InstanceHandle(InstanceId taskId) {
++    this.taskId = taskId;
+   }
  }
diff --cc 
server/src/main/java/org/apache/iotdb/db/mpp/operator/process/AggregateOperator.java
index 9604721,63e07b4..6366b29
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/operator/process/AggregateOperator.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/operator/process/AggregateOperator.java
@@@ -16,36 -16,14 +16,37 @@@
   * specific language governing permissions and limitations
   * under the License.
   */
 -package org.apache.iotdb.db.mpp.plan.node.process;
 +package org.apache.iotdb.db.mpp.operator.process;
  
- import com.google.common.util.concurrent.ListenableFuture;
  import org.apache.iotdb.db.mpp.common.TsBlock;
 -import org.apache.iotdb.db.mpp.plan.node.PlanNode;
 -import org.apache.iotdb.db.mpp.plan.node.PlanNodeId;
 +import org.apache.iotdb.db.mpp.operator.OperatorContext;
  
 -public class ProcessNode extends PlanNode<TsBlock> {
 -  public ProcessNode(PlanNodeId id) {
 -    super(id);
++import com.google.common.util.concurrent.ListenableFuture;
++
 +public class AggregateOperator implements ProcessOperator {
 +
-     @Override
-     public OperatorContext getOperatorContext() {
-         return null;
-     }
- 
-     @Override
-     public ListenableFuture<Void> isBlocked() {
-         return ProcessOperator.super.isBlocked();
-     }
- 
-     @Override
-     public TsBlock next() {
-         return null;
-     }
- 
-     @Override
-     public boolean hasNext() {
-         return false;
-     }
- 
-     @Override
-     public void close() throws Exception {
-         ProcessOperator.super.close();
-     }
++  @Override
++  public OperatorContext getOperatorContext() {
++    return null;
++  }
++
++  @Override
++  public ListenableFuture<Void> isBlocked() {
++    return ProcessOperator.super.isBlocked();
++  }
++
++  @Override
++  public TsBlock next() {
++    return null;
++  }
++
++  @Override
++  public boolean hasNext() {
++    return false;
++  }
++
++  @Override
++  public void close() throws Exception {
++    ProcessOperator.super.close();
+   }
  }
diff --cc 
server/src/main/java/org/apache/iotdb/db/mpp/operator/process/DeviceMergeOperator.java
index dab4df5,63e07b4..408296a
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/operator/process/DeviceMergeOperator.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/operator/process/DeviceMergeOperator.java
@@@ -16,35 -16,14 +16,36 @@@
   * specific language governing permissions and limitations
   * under the License.
   */
 -package org.apache.iotdb.db.mpp.plan.node.process;
 +package org.apache.iotdb.db.mpp.operator.process;
  
- import com.google.common.util.concurrent.ListenableFuture;
  import org.apache.iotdb.db.mpp.common.TsBlock;
 -import org.apache.iotdb.db.mpp.plan.node.PlanNode;
 -import org.apache.iotdb.db.mpp.plan.node.PlanNodeId;
 +import org.apache.iotdb.db.mpp.operator.OperatorContext;
  
 -public class ProcessNode extends PlanNode<TsBlock> {
 -  public ProcessNode(PlanNodeId id) {
 -    super(id);
++import com.google.common.util.concurrent.ListenableFuture;
++
 +public class DeviceMergeOperator implements ProcessOperator {
-     @Override
-     public OperatorContext getOperatorContext() {
-         return null;
-     }
++  @Override
++  public OperatorContext getOperatorContext() {
++    return null;
++  }
 +
-     @Override
-     public ListenableFuture<Void> isBlocked() {
-         return ProcessOperator.super.isBlocked();
-     }
++  @Override
++  public ListenableFuture<Void> isBlocked() {
++    return ProcessOperator.super.isBlocked();
++  }
 +
-     @Override
-     public TsBlock next() {
-         return null;
-     }
++  @Override
++  public TsBlock next() {
++    return null;
++  }
 +
-     @Override
-     public boolean hasNext() {
-         return false;
-     }
++  @Override
++  public boolean hasNext() {
++    return false;
++  }
 +
-     @Override
-     public void close() throws Exception {
-         ProcessOperator.super.close();
-     }
++  @Override
++  public void close() throws Exception {
++    ProcessOperator.super.close();
+   }
  }
diff --cc 
server/src/main/java/org/apache/iotdb/db/mpp/operator/process/FillOperator.java
index 3b35646,63e07b4..66dc3ea
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/operator/process/FillOperator.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/operator/process/FillOperator.java
@@@ -16,35 -16,14 +16,36 @@@
   * specific language governing permissions and limitations
   * under the License.
   */
 -package org.apache.iotdb.db.mpp.plan.node.process;
 +package org.apache.iotdb.db.mpp.operator.process;
  
- import com.google.common.util.concurrent.ListenableFuture;
  import org.apache.iotdb.db.mpp.common.TsBlock;
 -import org.apache.iotdb.db.mpp.plan.node.PlanNode;
 -import org.apache.iotdb.db.mpp.plan.node.PlanNodeId;
 +import org.apache.iotdb.db.mpp.operator.OperatorContext;
  
 -public class ProcessNode extends PlanNode<TsBlock> {
 -  public ProcessNode(PlanNodeId id) {
 -    super(id);
++import com.google.common.util.concurrent.ListenableFuture;
++
 +public class FillOperator implements ProcessOperator {
-     @Override
-     public OperatorContext getOperatorContext() {
-         return null;
-     }
++  @Override
++  public OperatorContext getOperatorContext() {
++    return null;
++  }
 +
-     @Override
-     public ListenableFuture<Void> isBlocked() {
-         return ProcessOperator.super.isBlocked();
-     }
++  @Override
++  public ListenableFuture<Void> isBlocked() {
++    return ProcessOperator.super.isBlocked();
++  }
 +
-     @Override
-     public TsBlock next() {
-         return null;
-     }
++  @Override
++  public TsBlock next() {
++    return null;
++  }
 +
-     @Override
-     public boolean hasNext() {
-         return false;
-     }
++  @Override
++  public boolean hasNext() {
++    return false;
++  }
 +
-     @Override
-     public void close() throws Exception {
-         ProcessOperator.super.close();
-     }
++  @Override
++  public void close() throws Exception {
++    ProcessOperator.super.close();
+   }
  }
diff --cc 
server/src/main/java/org/apache/iotdb/db/mpp/operator/process/FilterNullOperator.java
index f0a707c,63e07b4..8d15250
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/operator/process/FilterNullOperator.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/operator/process/FilterNullOperator.java
@@@ -16,36 -16,14 +16,37 @@@
   * specific language governing permissions and limitations
   * under the License.
   */
 -package org.apache.iotdb.db.mpp.plan.node.process;
 +package org.apache.iotdb.db.mpp.operator.process;
  
- import com.google.common.util.concurrent.ListenableFuture;
  import org.apache.iotdb.db.mpp.common.TsBlock;
 -import org.apache.iotdb.db.mpp.plan.node.PlanNode;
 -import org.apache.iotdb.db.mpp.plan.node.PlanNodeId;
 +import org.apache.iotdb.db.mpp.operator.OperatorContext;
  
 -public class ProcessNode extends PlanNode<TsBlock> {
 -  public ProcessNode(PlanNodeId id) {
 -    super(id);
++import com.google.common.util.concurrent.ListenableFuture;
++
 +public class FilterNullOperator implements ProcessOperator {
 +
-     @Override
-     public OperatorContext getOperatorContext() {
-         return null;
-     }
- 
-     @Override
-     public ListenableFuture<Void> isBlocked() {
-         return ProcessOperator.super.isBlocked();
-     }
- 
-     @Override
-     public TsBlock next() {
-         return null;
-     }
- 
-     @Override
-     public boolean hasNext() {
-         return false;
-     }
- 
-     @Override
-     public void close() throws Exception {
-         ProcessOperator.super.close();
-     }
++  @Override
++  public OperatorContext getOperatorContext() {
++    return null;
++  }
++
++  @Override
++  public ListenableFuture<Void> isBlocked() {
++    return ProcessOperator.super.isBlocked();
++  }
++
++  @Override
++  public TsBlock next() {
++    return null;
++  }
++
++  @Override
++  public boolean hasNext() {
++    return false;
++  }
++
++  @Override
++  public void close() throws Exception {
++    ProcessOperator.super.close();
+   }
  }
diff --cc 
server/src/main/java/org/apache/iotdb/db/mpp/operator/process/GroupByLevelOperator.java
index 435e016,63e07b4..10e9daa
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/operator/process/GroupByLevelOperator.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/operator/process/GroupByLevelOperator.java
@@@ -16,36 -16,14 +16,37 @@@
   * specific language governing permissions and limitations
   * under the License.
   */
 -package org.apache.iotdb.db.mpp.plan.node.process;
 +package org.apache.iotdb.db.mpp.operator.process;
  
- import com.google.common.util.concurrent.ListenableFuture;
  import org.apache.iotdb.db.mpp.common.TsBlock;
 -import org.apache.iotdb.db.mpp.plan.node.PlanNode;
 -import org.apache.iotdb.db.mpp.plan.node.PlanNodeId;
 +import org.apache.iotdb.db.mpp.operator.OperatorContext;
  
 -public class ProcessNode extends PlanNode<TsBlock> {
 -  public ProcessNode(PlanNodeId id) {
 -    super(id);
++import com.google.common.util.concurrent.ListenableFuture;
++
 +public class GroupByLevelOperator implements ProcessOperator {
 +
-     @Override
-     public OperatorContext getOperatorContext() {
-         return null;
-     }
- 
-     @Override
-     public ListenableFuture<Void> isBlocked() {
-         return ProcessOperator.super.isBlocked();
-     }
- 
-     @Override
-     public TsBlock next() {
-         return null;
-     }
- 
-     @Override
-     public boolean hasNext() {
-         return false;
-     }
- 
-     @Override
-     public void close() throws Exception {
-         ProcessOperator.super.close();
-     }
++  @Override
++  public OperatorContext getOperatorContext() {
++    return null;
++  }
++
++  @Override
++  public ListenableFuture<Void> isBlocked() {
++    return ProcessOperator.super.isBlocked();
++  }
++
++  @Override
++  public TsBlock next() {
++    return null;
++  }
++
++  @Override
++  public boolean hasNext() {
++    return false;
++  }
++
++  @Override
++  public void close() throws Exception {
++    ProcessOperator.super.close();
+   }
  }
diff --cc 
server/src/main/java/org/apache/iotdb/db/mpp/operator/process/LimitOperator.java
index 60d68ff,63e07b4..0efd325
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/operator/process/LimitOperator.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/operator/process/LimitOperator.java
@@@ -16,35 -16,14 +16,36 @@@
   * specific language governing permissions and limitations
   * under the License.
   */
 -package org.apache.iotdb.db.mpp.plan.node.process;
 +package org.apache.iotdb.db.mpp.operator.process;
  
- import com.google.common.util.concurrent.ListenableFuture;
  import org.apache.iotdb.db.mpp.common.TsBlock;
 -import org.apache.iotdb.db.mpp.plan.node.PlanNode;
 -import org.apache.iotdb.db.mpp.plan.node.PlanNodeId;
 +import org.apache.iotdb.db.mpp.operator.OperatorContext;
  
 -public class ProcessNode extends PlanNode<TsBlock> {
 -  public ProcessNode(PlanNodeId id) {
 -    super(id);
++import com.google.common.util.concurrent.ListenableFuture;
++
 +public class LimitOperator implements ProcessOperator {
-     @Override
-     public OperatorContext getOperatorContext() {
-         return null;
-     }
++  @Override
++  public OperatorContext getOperatorContext() {
++    return null;
++  }
 +
-     @Override
-     public ListenableFuture<Void> isBlocked() {
-         return ProcessOperator.super.isBlocked();
-     }
++  @Override
++  public ListenableFuture<Void> isBlocked() {
++    return ProcessOperator.super.isBlocked();
++  }
 +
-     @Override
-     public TsBlock next() {
-         return null;
-     }
++  @Override
++  public TsBlock next() {
++    return null;
++  }
 +
-     @Override
-     public boolean hasNext() {
-         return false;
-     }
++  @Override
++  public boolean hasNext() {
++    return false;
++  }
 +
-     @Override
-     public void close() throws Exception {
-         ProcessOperator.super.close();
-     }
++  @Override
++  public void close() throws Exception {
++    ProcessOperator.super.close();
+   }
  }
diff --cc 
server/src/main/java/org/apache/iotdb/db/mpp/operator/process/OffsetOperator.java
index 08f971b,63e07b4..25b4bc9
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/operator/process/OffsetOperator.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/operator/process/OffsetOperator.java
@@@ -16,36 -16,14 +16,37 @@@
   * specific language governing permissions and limitations
   * under the License.
   */
 -package org.apache.iotdb.db.mpp.plan.node.process;
 +package org.apache.iotdb.db.mpp.operator.process;
  
- import com.google.common.util.concurrent.ListenableFuture;
  import org.apache.iotdb.db.mpp.common.TsBlock;
 -import org.apache.iotdb.db.mpp.plan.node.PlanNode;
 -import org.apache.iotdb.db.mpp.plan.node.PlanNodeId;
 +import org.apache.iotdb.db.mpp.operator.OperatorContext;
  
 -public class ProcessNode extends PlanNode<TsBlock> {
 -  public ProcessNode(PlanNodeId id) {
 -    super(id);
++import com.google.common.util.concurrent.ListenableFuture;
++
 +public class OffsetOperator implements ProcessOperator {
 +
-     @Override
-     public OperatorContext getOperatorContext() {
-         return null;
-     }
- 
-     @Override
-     public ListenableFuture<Void> isBlocked() {
-         return ProcessOperator.super.isBlocked();
-     }
- 
-     @Override
-     public TsBlock next() {
-         return null;
-     }
- 
-     @Override
-     public boolean hasNext() {
-         return false;
-     }
- 
-     @Override
-     public void close() throws Exception {
-         ProcessOperator.super.close();
-     }
++  @Override
++  public OperatorContext getOperatorContext() {
++    return null;
++  }
++
++  @Override
++  public ListenableFuture<Void> isBlocked() {
++    return ProcessOperator.super.isBlocked();
++  }
++
++  @Override
++  public TsBlock next() {
++    return null;
++  }
++
++  @Override
++  public boolean hasNext() {
++    return false;
++  }
++
++  @Override
++  public void close() throws Exception {
++    ProcessOperator.super.close();
+   }
  }
diff --cc 
server/src/main/java/org/apache/iotdb/db/mpp/operator/process/ProcessOperator.java
index 320e505,39f8d17..aeb9535
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/operator/process/ProcessOperator.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/operator/process/ProcessOperator.java
@@@ -16,11 -16,12 +16,9 @@@
   * specific language governing permissions and limitations
   * under the License.
   */
 -package org.apache.iotdb.db.mpp.plan;
 +package org.apache.iotdb.db.mpp.operator.process;
  
 -public class PlanFragmentId {
 -  private String id;
 +import org.apache.iotdb.db.mpp.operator.Operator;
  
 -  public PlanFragmentId(String id) {
 -    this.id = id;
 -  }
 -}
 +// TODO should think about what interfaces should this ProcessOperator have
- public interface ProcessOperator extends Operator {
- 
- }
++public interface ProcessOperator extends Operator {}
diff --cc 
server/src/main/java/org/apache/iotdb/db/mpp/operator/process/SortOperator.java
index 45b5da4,63e07b4..0199d52
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/operator/process/SortOperator.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/operator/process/SortOperator.java
@@@ -16,36 -16,14 +16,37 @@@
   * specific language governing permissions and limitations
   * under the License.
   */
 -package org.apache.iotdb.db.mpp.plan.node.process;
 +package org.apache.iotdb.db.mpp.operator.process;
  
- import com.google.common.util.concurrent.ListenableFuture;
  import org.apache.iotdb.db.mpp.common.TsBlock;
 -import org.apache.iotdb.db.mpp.plan.node.PlanNode;
 -import org.apache.iotdb.db.mpp.plan.node.PlanNodeId;
 +import org.apache.iotdb.db.mpp.operator.OperatorContext;
  
 -public class ProcessNode extends PlanNode<TsBlock> {
 -  public ProcessNode(PlanNodeId id) {
 -    super(id);
++import com.google.common.util.concurrent.ListenableFuture;
++
 +public class SortOperator implements ProcessOperator {
 +
-     @Override
-     public OperatorContext getOperatorContext() {
-         return null;
-     }
- 
-     @Override
-     public ListenableFuture<Void> isBlocked() {
-         return ProcessOperator.super.isBlocked();
-     }
- 
-     @Override
-     public TsBlock next() {
-         return null;
-     }
- 
-     @Override
-     public boolean hasNext() {
-         return false;
-     }
- 
-     @Override
-     public void close() throws Exception {
-         ProcessOperator.super.close();
-     }
++  @Override
++  public OperatorContext getOperatorContext() {
++    return null;
++  }
++
++  @Override
++  public ListenableFuture<Void> isBlocked() {
++    return ProcessOperator.super.isBlocked();
++  }
++
++  @Override
++  public TsBlock next() {
++    return null;
++  }
++
++  @Override
++  public boolean hasNext() {
++    return false;
++  }
++
++  @Override
++  public void close() throws Exception {
++    ProcessOperator.super.close();
+   }
  }
diff --cc 
server/src/main/java/org/apache/iotdb/db/mpp/operator/process/TimeJoinOperator.java
index 0c95752,63e07b4..11cee59
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/operator/process/TimeJoinOperator.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/operator/process/TimeJoinOperator.java
@@@ -16,36 -16,14 +16,37 @@@
   * specific language governing permissions and limitations
   * under the License.
   */
 -package org.apache.iotdb.db.mpp.plan.node.process;
 +package org.apache.iotdb.db.mpp.operator.process;
  
- import com.google.common.util.concurrent.ListenableFuture;
  import org.apache.iotdb.db.mpp.common.TsBlock;
 -import org.apache.iotdb.db.mpp.plan.node.PlanNode;
 -import org.apache.iotdb.db.mpp.plan.node.PlanNodeId;
 +import org.apache.iotdb.db.mpp.operator.OperatorContext;
  
 -public class ProcessNode extends PlanNode<TsBlock> {
 -  public ProcessNode(PlanNodeId id) {
 -    super(id);
++import com.google.common.util.concurrent.ListenableFuture;
++
 +public class TimeJoinOperator implements ProcessOperator {
 +
-     @Override
-     public OperatorContext getOperatorContext() {
-         return null;
-     }
- 
-     @Override
-     public ListenableFuture<Void> isBlocked() {
-         return ProcessOperator.super.isBlocked();
-     }
- 
-     @Override
-     public TsBlock next() {
-         return null;
-     }
- 
-     @Override
-     public boolean hasNext() {
-         return false;
-     }
- 
-     @Override
-     public void close() throws Exception {
-         ProcessOperator.super.close();
-     }
++  @Override
++  public OperatorContext getOperatorContext() {
++    return null;
++  }
++
++  @Override
++  public ListenableFuture<Void> isBlocked() {
++    return ProcessOperator.super.isBlocked();
++  }
++
++  @Override
++  public TsBlock next() {
++    return null;
++  }
++
++  @Override
++  public boolean hasNext() {
++    return false;
++  }
++
++  @Override
++  public void close() throws Exception {
++    ProcessOperator.super.close();
+   }
  }
diff --cc 
server/src/main/java/org/apache/iotdb/db/mpp/operator/sink/SinkOperator.java
index 7f75c4a,a4cb88c..b03e7ec
--- 
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
@@@ -16,29 -16,23 +16,29 @@@
   * specific language governing permissions and limitations
   * under the License.
   */
 -package org.apache.iotdb.db.mpp.plan.node.process;
 +package org.apache.iotdb.db.mpp.operator.sink;
  
 -import org.apache.iotdb.db.mpp.plan.node.PlanNodeId;
 -import org.apache.iotdb.db.qp.logical.crud.FilterOperator;
 +import org.apache.iotdb.db.mpp.operator.Operator;
  
 -/** The FilterNode is responsible to filter the RowRecord from TsBlock. */
 -public class FilterNode extends ProcessNode {
 +import java.nio.ByteBuffer;
  
 -  // The filter
 -  private FilterOperator rowFilter;
 +public interface SinkOperator extends Operator {
  
-     /**
-      * Sends a tsBlock to an unpartitioned buffer. If no-more-tsBlocks has 
been set, the send tsBlock
-      * call is ignored. This can happen with limit queries.
-      */
-     void send(ByteBuffer tsBlock);
 -  public FilterNode(PlanNodeId id) {
 -    super(id);
 -  }
++  /**
++   * Sends a tsBlock to an unpartitioned buffer. If no-more-tsBlocks has been 
set, the send tsBlock
++   * call is ignored. This can happen with limit queries.
++   */
++  void send(ByteBuffer tsBlock);
  
-     /**
-      * Notify SinkHandle that no more tsBlocks will be sent. Any future calls 
to send a tsBlock are
-      * ignored.
-      */
-     void setNoMoreTsBlocks();
 -  public FilterNode(PlanNodeId id, FilterOperator rowFilter) {
 -    this(id);
 -    this.rowFilter = rowFilter;
 -  }
++  /**
++   * Notify SinkHandle that no more tsBlocks will be sent. Any future calls 
to send a tsBlock are
++   * ignored.
++   */
++  void setNoMoreTsBlocks();
 +
-     /**
-      * Abort the sink handle, discarding all tsBlocks which may still in 
memory buffer, but blocking
-      * readers. It is expected that readers will be unblocked when the failed 
query is cleaned up.
-      */
-     void abort();
++  /**
++   * Abort the sink handle, discarding all tsBlocks which may still in memory 
buffer, but blocking
++   * readers. It is expected that readers will be unblocked when the failed 
query is cleaned up.
++   */
++  void abort();
  }
diff --cc 
server/src/main/java/org/apache/iotdb/db/mpp/operator/source/SeriesAggregateScanOperator.java
index 83c6309,63e07b4..fcd4605
--- 
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
@@@ -16,41 -16,14 +16,42 @@@
   * specific language governing permissions and limitations
   * under the License.
   */
 -package org.apache.iotdb.db.mpp.plan.node.process;
 +package org.apache.iotdb.db.mpp.operator.source;
  
- import com.google.common.util.concurrent.ListenableFuture;
  import org.apache.iotdb.db.mpp.common.TsBlock;
 -import org.apache.iotdb.db.mpp.plan.node.PlanNode;
 -import org.apache.iotdb.db.mpp.plan.node.PlanNodeId;
 +import org.apache.iotdb.db.mpp.operator.OperatorContext;
 +import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanNodeId;
  
 -public class ProcessNode extends PlanNode<TsBlock> {
 -  public ProcessNode(PlanNodeId id) {
 -    super(id);
++import com.google.common.util.concurrent.ListenableFuture;
++
 +public class SeriesAggregateScanOperator implements SourceOperator {
-     @Override
-     public OperatorContext getOperatorContext() {
-         return null;
-     }
- 
-     @Override
-     public ListenableFuture<Void> isBlocked() {
-         return SourceOperator.super.isBlocked();
-     }
- 
-     @Override
-     public TsBlock next() {
-         return null;
-     }
- 
-     @Override
-     public boolean hasNext() {
-         return false;
-     }
- 
-     @Override
-     public void close() throws Exception {
-         SourceOperator.super.close();
-     }
- 
-     @Override
-     public PlanNodeId getSourceId() {
-         return null;
-     }
++  @Override
++  public OperatorContext getOperatorContext() {
++    return null;
++  }
++
++  @Override
++  public ListenableFuture<Void> isBlocked() {
++    return SourceOperator.super.isBlocked();
++  }
++
++  @Override
++  public TsBlock next() {
++    return null;
++  }
++
++  @Override
++  public boolean hasNext() {
++    return false;
++  }
++
++  @Override
++  public void close() throws Exception {
++    SourceOperator.super.close();
++  }
++
++  @Override
++  public PlanNodeId getSourceId() {
++    return null;
+   }
  }
diff --cc 
server/src/main/java/org/apache/iotdb/db/mpp/operator/source/SeriesScanOperator.java
index 42085ae,63e07b4..c0bc9ca
--- 
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
@@@ -16,42 -16,14 +16,43 @@@
   * specific language governing permissions and limitations
   * under the License.
   */
 -package org.apache.iotdb.db.mpp.plan.node.process;
 +package org.apache.iotdb.db.mpp.operator.source;
  
- import com.google.common.util.concurrent.ListenableFuture;
  import org.apache.iotdb.db.mpp.common.TsBlock;
 -import org.apache.iotdb.db.mpp.plan.node.PlanNode;
 -import org.apache.iotdb.db.mpp.plan.node.PlanNodeId;
 +import org.apache.iotdb.db.mpp.operator.OperatorContext;
 +import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanNodeId;
  
 -public class ProcessNode extends PlanNode<TsBlock> {
 -  public ProcessNode(PlanNodeId id) {
 -    super(id);
++import com.google.common.util.concurrent.ListenableFuture;
++
 +public class SeriesScanOperator implements SourceOperator {
 +
-     @Override
-     public OperatorContext getOperatorContext() {
-         return null;
-     }
- 
-     @Override
-     public ListenableFuture<Void> isBlocked() {
-         return SourceOperator.super.isBlocked();
-     }
- 
-     @Override
-     public TsBlock next() {
-         return null;
-     }
- 
-     @Override
-     public boolean hasNext() {
-         return false;
-     }
- 
-     @Override
-     public void close() throws Exception {
-         SourceOperator.super.close();
-     }
- 
-     @Override
-     public PlanNodeId getSourceId() {
-         return null;
-     }
++  @Override
++  public OperatorContext getOperatorContext() {
++    return null;
++  }
++
++  @Override
++  public ListenableFuture<Void> isBlocked() {
++    return SourceOperator.super.isBlocked();
++  }
++
++  @Override
++  public TsBlock next() {
++    return null;
++  }
++
++  @Override
++  public boolean hasNext() {
++    return false;
++  }
++
++  @Override
++  public void close() throws Exception {
++    SourceOperator.super.close();
++  }
++
++  @Override
++  public PlanNodeId getSourceId() {
++    return null;
+   }
  }
diff --cc 
server/src/main/java/org/apache/iotdb/db/mpp/operator/source/SourceOperator.java
index 63b9acd,39f8d17..8454fd6
--- 
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
@@@ -16,12 -16,12 +16,12 @@@
   * specific language governing permissions and limitations
   * under the License.
   */
 -package org.apache.iotdb.db.mpp.plan;
 +package org.apache.iotdb.db.mpp.operator.source;
  
 -public class PlanFragmentId {
 -  private String id;
 +import org.apache.iotdb.db.mpp.operator.Operator;
 +import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanNodeId;
  
 -  public PlanFragmentId(String id) {
 -    this.id = id;
 -  }
 +public interface SourceOperator extends Operator {
 +
-     PlanNodeId getSourceId();
++  PlanNodeId getSourceId();
  }
diff --cc 
server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/LocalExecutionPlanner.java
index 920c52f,0000000..a4dd7539
mode 100644,000000..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
@@@ -1,126 -1,0 +1,126 @@@
 +/*
 + * 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.planner;
 +
 +import org.apache.iotdb.db.mpp.execution.InstanceContext;
 +import org.apache.iotdb.db.mpp.operator.Operator;
 +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.*;
 +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.qp.logical.crud.FilterOperator;
 +
 +import java.util.List;
 +
 +/**
-  * used to plan a fragment instance. Currently, we simply change it from 
PlanNode to executable Operator tree,
-  * but in the future, we may split one fragment instance into multiple 
pipeline to run a fragment instance parallel and take full advantage of 
multi-cores
++ * used to plan a fragment instance. Currently, we simply change it from 
PlanNode to executable
++ * Operator tree, but in the future, we may split one fragment instance into 
multiple pipeline to
++ * run a fragment instance parallel and take full advantage of multi-cores
 + */
 +public class LocalExecutionPlanner {
 +
++  /** This Visitor is responsible for transferring PlanNode Tree to Operator 
Tree */
++  private class Visitor extends PlanVisitor<Operator, 
LocalExecutionPlanContext> {
 +
-     /**
-      * This Visitor is responsible for transferring PlanNode Tree to Operator 
Tree
-      */
-     private class Visitor extends PlanVisitor<Operator, 
LocalExecutionPlanContext> {
- 
-         @Override
-         public Operator visitPlan(PlanNode node, LocalExecutionPlanContext 
context) {
-             throw new UnsupportedOperationException("should call the concrete 
visitXX() method");
-         }
- 
-         @Override
-         public Operator visitSeriesScan(SeriesScanNode node, 
LocalExecutionPlanContext context) {
-             return super.visitSeriesScan(node, context);
-         }
- 
-         @Override
-         public Operator visitSeriesAggregate(SeriesAggregateScanNode node, 
LocalExecutionPlanContext context) {
-             return super.visitSeriesAggregate(node, context);
-         }
- 
-         @Override
-         public Operator visitDeviceMerge(DeviceMergeNode node, 
LocalExecutionPlanContext context) {
-             return super.visitDeviceMerge(node, context);
-         }
- 
-         @Override
-         public Operator visitFill(FillNode node, LocalExecutionPlanContext 
context) {
-             return super.visitFill(node, context);
-         }
- 
-         @Override
-         public Operator visitFilter(FilterNode node, 
LocalExecutionPlanContext context) {
-             PlanNode child = node.getChild();
- 
-             FilterOperator filterExpression = node.getPredicate();
-             List<String> outputSymbols = node.getOutputColumnNames();
-             return super.visitFilter(node, context);
-         }
- 
-         @Override
-         public Operator visitFilterNull(FilterNullNode node, 
LocalExecutionPlanContext context) {
-             return super.visitFilterNull(node, context);
-         }
- 
-         @Override
-         public Operator visitGroupByLevel(GroupByLevelNode node, 
LocalExecutionPlanContext context) {
-             return super.visitGroupByLevel(node, context);
-         }
- 
-         @Override
-         public Operator visitLimit(LimitNode node, LocalExecutionPlanContext 
context) {
-             return super.visitLimit(node, context);
-         }
- 
-         @Override
-         public Operator visitOffset(OffsetNode node, 
LocalExecutionPlanContext context) {
-             return super.visitOffset(node, context);
-         }
- 
-         @Override
-         public Operator visitRowBasedSeriesAggregate(AggregateNode node, 
LocalExecutionPlanContext context) {
-             return super.visitRowBasedSeriesAggregate(node, context);
-         }
- 
-         @Override
-         public Operator visitSort(SortNode node, LocalExecutionPlanContext 
context) {
-             return super.visitSort(node, context);
-         }
- 
-         @Override
-         public Operator visitTimeJoin(TimeJoinNode node, 
LocalExecutionPlanContext context) {
-             return super.visitTimeJoin(node, context);
-         }
++    @Override
++    public Operator visitPlan(PlanNode node, LocalExecutionPlanContext 
context) {
++      throw new UnsupportedOperationException("should call the concrete 
visitXX() method");
 +    }
 +
-     private static class LocalExecutionPlanContext {
-         private final InstanceContext taskContext;
-         private int nextOperatorId = 0;
++    @Override
++    public Operator visitSeriesScan(SeriesScanNode node, 
LocalExecutionPlanContext context) {
++      return super.visitSeriesScan(node, context);
++    }
++
++    @Override
++    public Operator visitSeriesAggregate(
++        SeriesAggregateScanNode node, LocalExecutionPlanContext context) {
++      return super.visitSeriesAggregate(node, context);
++    }
++
++    @Override
++    public Operator visitDeviceMerge(DeviceMergeNode node, 
LocalExecutionPlanContext context) {
++      return super.visitDeviceMerge(node, context);
++    }
++
++    @Override
++    public Operator visitFill(FillNode node, LocalExecutionPlanContext 
context) {
++      return super.visitFill(node, context);
++    }
++
++    @Override
++    public Operator visitFilter(FilterNode node, LocalExecutionPlanContext 
context) {
++      PlanNode child = node.getChild();
++
++      FilterOperator filterExpression = node.getPredicate();
++      List<String> outputSymbols = node.getOutputColumnNames();
++      return super.visitFilter(node, context);
++    }
++
++    @Override
++    public Operator visitFilterNull(FilterNullNode node, 
LocalExecutionPlanContext context) {
++      return super.visitFilterNull(node, context);
++    }
++
++    @Override
++    public Operator visitGroupByLevel(GroupByLevelNode node, 
LocalExecutionPlanContext context) {
++      return super.visitGroupByLevel(node, context);
++    }
++
++    @Override
++    public Operator visitLimit(LimitNode node, LocalExecutionPlanContext 
context) {
++      return super.visitLimit(node, context);
++    }
++
++    @Override
++    public Operator visitOffset(OffsetNode node, LocalExecutionPlanContext 
context) {
++      return super.visitOffset(node, context);
++    }
++
++    @Override
++    public Operator visitRowBasedSeriesAggregate(
++        AggregateNode node, LocalExecutionPlanContext context) {
++      return super.visitRowBasedSeriesAggregate(node, context);
++    }
 +
-         public LocalExecutionPlanContext(InstanceContext taskContext) {
-             this.taskContext = taskContext;
-         }
++    @Override
++    public Operator visitSort(SortNode node, LocalExecutionPlanContext 
context) {
++      return super.visitSort(node, context);
++    }
++
++    @Override
++    public Operator visitTimeJoin(TimeJoinNode node, 
LocalExecutionPlanContext context) {
++      return super.visitTimeJoin(node, context);
++    }
++  }
++
++  private static class LocalExecutionPlanContext {
++    private final InstanceContext taskContext;
++    private int nextOperatorId = 0;
++
++    public LocalExecutionPlanContext(InstanceContext taskContext) {
++      this.taskContext = taskContext;
++    }
 +
-         private int getNextOperatorId() {
-             return nextOperatorId++;
-         }
++    private int getNextOperatorId() {
++      return nextOperatorId++;
 +    }
++  }
 +}
diff --cc 
server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/DistributedQueryPlan.java
index 96185bd,9954c74..50835dd
--- 
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
@@@ -16,11 -16,11 +16,10 @@@
   * specific language governing permissions and limitations
   * under the License.
   */
 -package org.apache.iotdb.db.mpp.plan;
 +package org.apache.iotdb.db.mpp.sql.planner.plan;
  
  import org.apache.iotdb.db.mpp.common.QueryContext;
--import org.apache.iotdb.db.mpp.common.TsBlock;
 -import org.apache.iotdb.db.mpp.plan.node.PlanNode;
 +import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanNode;
  
  import java.util.List;
  
diff --cc 
server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/LogicalQueryPlan.java
index f3dce2f,5094df6..666fbf3
--- 
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
@@@ -16,11 -16,11 +16,10 @@@
   * specific language governing permissions and limitations
   * under the License.
   */
 -package org.apache.iotdb.db.mpp.plan;
 +package org.apache.iotdb.db.mpp.sql.planner.plan;
  
  import org.apache.iotdb.db.mpp.common.QueryContext;
--import org.apache.iotdb.db.mpp.common.TsBlock;
 -import org.apache.iotdb.db.mpp.plan.node.PlanNode;
 +import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanNode;
  
  /**
   * LogicalQueryPlan represents a logical query plan. It stores the root node 
of corresponding query
diff --cc 
server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/PlanFragment.java
index bf13247,fc49264..0aa3ac5
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/PlanFragment.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/PlanFragment.java
@@@ -16,10 -16,10 +16,9 @@@
   * specific language governing permissions and limitations
   * under the License.
   */
 -package org.apache.iotdb.db.mpp.plan;
 +package org.apache.iotdb.db.mpp.sql.planner.plan;
  
--import org.apache.iotdb.db.mpp.common.TsBlock;
 -import org.apache.iotdb.db.mpp.plan.node.PlanNode;
 +import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanNode;
  
  // TODO: consider whether it is necessary to make PlanFragment as a TreeNode
  /** PlanFragment contains a sub-query of distributed query. */
diff --cc 
server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/PlanNode.java
index 583d4a3,1a3f103..3a3975f
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/PlanNode.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/PlanNode.java
@@@ -16,35 -16,19 +16,32 @@@
   * specific language governing permissions and limitations
   * under the License.
   */
 -package org.apache.iotdb.db.mpp.plan.node;
 +package org.apache.iotdb.db.mpp.sql.planner.plan.node;
  
- 
 -import org.apache.iotdb.db.mpp.common.TreeNode;
 +import java.util.List;
 +
 +import static java.util.Objects.requireNonNull;
 +
- /**
-  * The base class of query executable operators, which is used to compose 
logical query plan.
-  */
++/** The base class of query executable operators, which is used to compose 
logical query plan. */
 +// TODO: consider how to restrict the children type for each type of 
ExecOperator
 +public abstract class PlanNode {
  
 -/**
 - * @author xingtanzjr The base class of query executable operators, which is 
used to compose logical
 - *     query plan. TODO: consider how to restrict the children type for each 
type of ExecOperator
 - *     TODO: consider to fix the Template type as TsBlock
 - */
 -public abstract class PlanNode<T> extends TreeNode<PlanNode<T>> {
    private PlanNodeId id;
  
 -  public PlanNode(PlanNodeId id) {
 +  protected PlanNode(PlanNodeId id) {
 +    requireNonNull(id, "id is null");
      this.id = id;
    }
 +
 +  public PlanNodeId getId() {
 +    return id;
 +  }
 +
 +  public abstract List<PlanNode> getChildren();
 +
 +  public abstract List<String> getOutputColumnNames();
 +
 +  public <R, C> R accept(PlanVisitor<R, C> visitor, C context) {
 +    return visitor.visitPlan(this, context);
 +  }
  }
diff --cc 
server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/PlanNodeId.java
index 426c778,e5a69e9..f829dfb
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/PlanNodeId.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/PlanNodeId.java
@@@ -1,20 -1,22 +1,22 @@@
- // 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.
+ /*
+  * 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.plan.node;
 +package org.apache.iotdb.db.mpp.sql.planner.plan.node;
  
  public class PlanNodeId {
    private String id;
diff --cc 
server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/PlanVisitor.java
index ecbc043,0000000..0be251d
mode 100644,000000..100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/PlanVisitor.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/PlanVisitor.java
@@@ -1,76 -1,0 +1,76 @@@
 +/*
 + * 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.planner.plan.node;
 +
 +import org.apache.iotdb.db.mpp.sql.planner.plan.node.process.*;
 +import 
org.apache.iotdb.db.mpp.sql.planner.plan.node.source.SeriesAggregateScanNode;
 +import org.apache.iotdb.db.mpp.sql.planner.plan.node.source.SeriesScanNode;
 +
 +public abstract class PlanVisitor<R, C> {
 +
-     public abstract R visitPlan(PlanNode node, C context);
++  public abstract R visitPlan(PlanNode node, C context);
 +
-     public R visitSeriesScan(SeriesScanNode node, C context) {
-         return visitPlan(node, context);
-     }
++  public R visitSeriesScan(SeriesScanNode node, C context) {
++    return visitPlan(node, context);
++  }
 +
-     public R visitSeriesAggregate(SeriesAggregateScanNode node, C context) {
-         return visitPlan(node, context);
-     }
++  public R visitSeriesAggregate(SeriesAggregateScanNode node, C context) {
++    return visitPlan(node, context);
++  }
 +
-     public R visitDeviceMerge(DeviceMergeNode node, C context) {
-         return visitPlan(node, context);
-     }
++  public R visitDeviceMerge(DeviceMergeNode node, C context) {
++    return visitPlan(node, context);
++  }
 +
-     public R visitFill(FillNode node, C context) {
-         return visitPlan(node, context);
-     }
++  public R visitFill(FillNode node, C context) {
++    return visitPlan(node, context);
++  }
 +
-     public R visitFilter(FilterNode node, C context) {
-         return visitPlan(node, context);
-     }
++  public R visitFilter(FilterNode node, C context) {
++    return visitPlan(node, context);
++  }
 +
-     public R visitFilterNull(FilterNullNode node, C context) {
-         return visitPlan(node, context);
-     }
++  public R visitFilterNull(FilterNullNode node, C context) {
++    return visitPlan(node, context);
++  }
 +
-     public R visitGroupByLevel(GroupByLevelNode node, C context) {
-         return visitPlan(node, context);
-     }
++  public R visitGroupByLevel(GroupByLevelNode node, C context) {
++    return visitPlan(node, context);
++  }
 +
-     public R visitLimit(LimitNode node, C context) {
-         return visitPlan(node, context);
-     }
++  public R visitLimit(LimitNode node, C context) {
++    return visitPlan(node, context);
++  }
 +
-     public R visitOffset(OffsetNode node, C context) {
-         return visitPlan(node, context);
-     }
++  public R visitOffset(OffsetNode node, C context) {
++    return visitPlan(node, context);
++  }
 +
-     public R visitRowBasedSeriesAggregate(AggregateNode node, C context) {
-         return visitPlan(node, context);
-     }
++  public R visitRowBasedSeriesAggregate(AggregateNode node, C context) {
++    return visitPlan(node, context);
++  }
 +
-     public R visitSort(SortNode node, C context) {
-         return visitPlan(node, context);
-     }
++  public R visitSort(SortNode node, C context) {
++    return visitPlan(node, context);
++  }
 +
-     public R visitTimeJoin(TimeJoinNode node, C context) {
-         return visitPlan(node, context);
-     }
++  public R visitTimeJoin(TimeJoinNode node, C context) {
++    return visitPlan(node, context);
++  }
 +}
diff --cc 
server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/AggregateNode.java
index 4588ea4,9d7b943..9382fbf
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/AggregateNode.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/AggregateNode.java
@@@ -25,14 -23,13 +25,14 @@@ import org.apache.iotdb.db.mpp.sql.plan
  import org.apache.iotdb.db.query.expression.unary.FunctionExpression;
  
  import java.util.List;
 +import java.util.Map;
  
  /**
-  * This node is used to aggregate required series from multiple sources. The 
source data will be input as a TsBlock,
-  * it may be raw data or partial aggregation result.
-  * This node will output the final series aggregated result represented by 
TsBlock.
 - * This node is used to aggregate required series by raw data. The raw data 
will be input as a
 - * TsBlock. This node will output the series aggregated result represented by 
TsBlock Thus, the
 - * columns in output TsBlock will be different from input TsBlock.
++ * This node is used to aggregate required series from multiple sources. The 
source data will be
++ * input as a TsBlock, it may be raw data or partial aggregation result. This 
node will output the
++ * final series aggregated result represented by TsBlock.
   */
 -public class RowBasedSeriesAggregateNode extends ProcessNode {
 +public class AggregateNode extends ProcessNode {
    // The parameter of `group by time`
    // Its value will be null if there is no `group by time` clause,
    private GroupByTimeParameter groupByTimeParameter;
@@@ -41,32 -38,22 +41,34 @@@
    // result TsBlock
    // (Currently we only support one series in the aggregation function)
    // TODO: need consider whether it is suitable the aggregation function 
using FunctionExpression
 -  private List<FunctionExpression> aggregateFuncList;
 +  private Map<String, FunctionExpression> aggregateFuncMap;
  
 -  public RowBasedSeriesAggregateNode(PlanNodeId id) {
 +  private final List<PlanNode> children;
 +  private final List<String> columnNames;
 +
-   public AggregateNode(PlanNodeId id, Map<String, FunctionExpression> 
aggregateFuncMap, List<PlanNode> children, List<String> columnNames) {
++  public AggregateNode(
++      PlanNodeId id,
++      Map<String, FunctionExpression> aggregateFuncMap,
++      List<PlanNode> children,
++      List<String> columnNames) {
      super(id);
 +    this.aggregateFuncMap = aggregateFuncMap;
 +    this.children = children;
 +    this.columnNames = columnNames;
    }
  
 -  public RowBasedSeriesAggregateNode(PlanNodeId id, List<FunctionExpression> 
aggregateFuncList) {
 -    this(id);
 -    this.aggregateFuncList = aggregateFuncList;
 +  @Override
 +  public List<PlanNode> getChildren() {
 +    return children;
    }
  
 -  public RowBasedSeriesAggregateNode(
 -      PlanNodeId id,
 -      List<FunctionExpression> aggregateFuncList,
 -      GroupByTimeParameter groupByTimeParameter) {
 -    this(id, aggregateFuncList);
 -    this.groupByTimeParameter = groupByTimeParameter;
 +  @Override
 +  public List<String> getOutputColumnNames() {
 +    return columnNames;
 +  }
 +
 +  @Override
 +  public <R, C> R accept(PlanVisitor<R, C> visitor, C context) {
 +    return visitor.visitRowBasedSeriesAggregate(this, context);
    }
- 
- 
  }
diff --cc 
server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/DeviceMergeNode.java
index 1f61e88,9269545..54a0cf8
--- 
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
@@@ -16,15 -16,14 +16,15 @@@
   * specific language governing permissions and limitations
   * under the License.
   */
 -package org.apache.iotdb.db.mpp.plan.node.process;
 +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.common.FilterNullPolicy;
+ import org.apache.iotdb.db.mpp.common.OrderBy;
 -import org.apache.iotdb.db.mpp.common.TsBlock;
 -import org.apache.iotdb.db.mpp.common.WithoutPolicy;
 -import org.apache.iotdb.db.mpp.plan.node.PlanNode;
 -import org.apache.iotdb.db.mpp.plan.node.PlanNodeId;
 +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 java.util.List;
  import java.util.Map;
  
  /**
diff --cc 
server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/FillNode.java
index 29396c1,31e57cd..af88895
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/FillNode.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/FillNode.java
@@@ -16,15 -16,10 +16,16 @@@
   * specific language governing permissions and limitations
   * under the License.
   */
 -package org.apache.iotdb.db.mpp.plan.node.process;
 +package org.apache.iotdb.db.mpp.sql.planner.plan.node.process;
  
- import com.google.common.collect.ImmutableList;
  import org.apache.iotdb.db.mpp.common.FillPolicy;
 -import org.apache.iotdb.db.mpp.plan.node.PlanNodeId;
 +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 com.google.common.collect.ImmutableList;
++
 +import java.util.List;
  
  /** FillNode is used to fill the empty field in one row. */
  public class FillNode extends ProcessNode {
diff --cc 
server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/FilterNode.java
index a6108fe,0000000..4504890
mode 100644,000000..100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/FilterNode.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/FilterNode.java
@@@ -1,64 -1,0 +1,66 @@@
 +/*
 + * 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.planner.plan.node.process;
 +
- import com.google.common.collect.ImmutableList;
 +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.qp.logical.crud.FilterOperator;
 +
++import com.google.common.collect.ImmutableList;
++
 +import java.util.List;
 +
 +/** The FilterNode is responsible to filter the RowRecord from TsBlock. */
 +public class FilterNode extends ProcessNode {
 +
 +  private final PlanNode child;
-   // TODO we need to rename it to something like expression in order to 
distinguish from Operator class
++  // TODO we need to rename it to something like expression in order to 
distinguish from Operator
++  // class
 +  private final FilterOperator predicate;
 +
 +  public FilterNode(PlanNodeId id, PlanNode child, FilterOperator predicate) {
 +    super(id);
 +    this.child = child;
 +    this.predicate = predicate;
 +  }
 +
 +  @Override
 +  public List<PlanNode> getChildren() {
 +    return ImmutableList.of(child);
 +  }
 +
 +  @Override
 +  public List<String> getOutputColumnNames() {
 +    return child.getOutputColumnNames();
 +  }
 +
 +  @Override
 +  public <R, C> R accept(PlanVisitor<R, C> visitor, C context) {
 +    return visitor.visitFilter(this, context);
 +  }
 +
 +  public FilterOperator getPredicate() {
 +    return predicate;
 +  }
 +
 +  public PlanNode getChild() {
 +    return child;
 +  }
 +}
diff --cc 
server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/FilterNullNode.java
index 976af21,0000000..fbf8cc9
mode 100644,000000..100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/FilterNullNode.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/FilterNullNode.java
@@@ -1,70 -1,0 +1,69 @@@
 +/*
 + * 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.planner.plan.node.process;
 +
- import com.google.common.collect.ImmutableList;
 +import org.apache.iotdb.db.mpp.common.FilterNullPolicy;
 +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.qp.logical.crud.FilterOperator;
++
++import com.google.common.collect.ImmutableList;
 +
 +import java.util.List;
 +
 +/** WithoutNode is used to discard specific rows from upstream node. */
 +public class FilterNullNode extends ProcessNode {
 +
 +  // The policy to discard the result from upstream operator
 +  private FilterNullPolicy discardPolicy;
 +
 +  private PlanNode child;
 +
 +  private List<String> filterNullColumnNames;
 +
- 
 +  public FilterNullNode(PlanNodeId id, PlanNode child) {
 +    super(id);
 +    this.child = child;
 +  }
 +
 +  public FilterNullNode(PlanNodeId id, PlanNode child, List<String> 
filterNullColumnNames) {
 +    super(id);
 +    this.child = child;
 +    this.filterNullColumnNames = filterNullColumnNames;
 +  }
 +
 +  @Override
 +  public List<PlanNode> getChildren() {
 +    return ImmutableList.of(child);
 +  }
 +
 +  @Override
 +  public List<String> getOutputColumnNames() {
 +    return child.getOutputColumnNames();
 +  }
 +
 +  @Override
 +  public <R, C> R accept(PlanVisitor<R, C> visitor, C context) {
 +    return visitor.visitFilterNull(this, context);
 +  }
 +
 +  public void setFilterNullColumnNames(List<String> filterNullColumnNames) {
 +    this.filterNullColumnNames = filterNullColumnNames;
 +  }
 +}
diff --cc 
server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/GroupByLevelNode.java
index b877933,538d6d8..7dcca76
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/GroupByLevelNode.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/GroupByLevelNode.java
@@@ -36,47 -32,10 +36,48 @@@ import java.util.List
   */
  public class GroupByLevelNode extends ProcessNode {
  
 +  private PlanNode child;
 +
    private int[] groupByLevels;
  
 -  public GroupByLevelNode(PlanNodeId id, int[] groupByLevels) {
 +  private List<String> columnNames;
 +
-   public GroupByLevelNode(PlanNodeId id, PlanNode child, int[] groupByLevels, 
List<String> columnNames) {
++  public GroupByLevelNode(
++      PlanNodeId id, PlanNode child, int[] groupByLevels, List<String> 
columnNames) {
      super(id);
 +    this.child = child;
 +    this.groupByLevels = groupByLevels;
 +    this.columnNames = columnNames;
 +  }
 +
 +  @Override
 +  public List<PlanNode> getChildren() {
 +    return child.getChildren();
 +  }
 +
 +  @Override
 +  public List<String> getOutputColumnNames() {
 +    return columnNames;
 +  }
 +
 +  @Override
 +  public <R, C> R accept(PlanVisitor<R, C> visitor, C context) {
 +    return visitor.visitGroupByLevel(this, context);
 +  }
 +
 +  public int[] getGroupByLevels() {
 +    return groupByLevels;
 +  }
 +
 +  public void setGroupByLevels(int[] groupByLevels) {
      this.groupByLevels = groupByLevels;
    }
 +
 +  public List<String> getColumnNames() {
 +    return columnNames;
 +  }
 +
 +  public void setColumnNames(List<String> columnNames) {
 +    this.columnNames = columnNames;
 +  }
  }
diff --cc 
server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/LimitNode.java
index 67b7203,9596c1a..8fdc9b4
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/LimitNode.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/LimitNode.java
@@@ -16,14 -16,9 +16,15 @@@
   * specific language governing permissions and limitations
   * under the License.
   */
 -package org.apache.iotdb.db.mpp.plan.node.process;
 +package org.apache.iotdb.db.mpp.sql.planner.plan.node.process;
  
- import com.google.common.collect.ImmutableList;
 -import org.apache.iotdb.db.mpp.plan.node.PlanNodeId;
 +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 com.google.common.collect.ImmutableList;
++
 +import java.util.List;
  
  /** LimitNode is used to select top n result. It uses the default order of 
upstream nodes */
  public class LimitNode extends ProcessNode {
diff --cc 
server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/ProcessNode.java
index 8348a21,63e07b4..9c1fec5
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/ProcessNode.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/ProcessNode.java
@@@ -16,14 -16,13 +16,13 @@@
   * specific language governing permissions and limitations
   * under the License.
   */
 -package org.apache.iotdb.db.mpp.plan.node.process;
 +package org.apache.iotdb.db.mpp.sql.planner.plan.node.process;
  
--import org.apache.iotdb.db.mpp.common.TsBlock;
 -import org.apache.iotdb.db.mpp.plan.node.PlanNode;
 -import org.apache.iotdb.db.mpp.plan.node.PlanNodeId;
 +import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanNode;
 +import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanNodeId;
 +
 +public abstract class ProcessNode extends PlanNode {
  
 -public class ProcessNode extends PlanNode<TsBlock> {
    public ProcessNode(PlanNodeId id) {
      super(id);
    }
diff --cc 
server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/SortNode.java
index 49b59a8,1e83783..19464a2
--- 
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
@@@ -16,15 -16,10 +16,16 @@@
   * specific language governing permissions and limitations
   * under the License.
   */
 -package org.apache.iotdb.db.mpp.plan.node.process;
 +package org.apache.iotdb.db.mpp.sql.planner.plan.node.process;
  
- import com.google.common.collect.ImmutableList;
  import org.apache.iotdb.db.mpp.common.OrderBy;
 -import org.apache.iotdb.db.mpp.plan.node.PlanNodeId;
 +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 com.google.common.collect.ImmutableList;
++
 +import java.util.List;
  
  /**
   * In general, the parameter in sortNode should be pushed down to the 
upstream operators. In our
diff --cc 
server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/TimeJoinNode.java
index 998594c,ab48cc4..74fac14
--- 
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
@@@ -42,34 -41,19 +42,39 @@@ public class TimeJoinNode extends Proce
    // The without policy is able to be push down to the TimeJoinOperator 
because we can know whether
    // a row contains
    // null or not.
 -  private WithoutPolicy withoutPolicy;
 +  private FilterNullPolicy filterNullPolicy;
  
 -  public TimeJoinNode(PlanNodeId id) {
 +  private List<PlanNode> children;
 +
- 
-   public TimeJoinNode(PlanNodeId id, OrderBy mergeOrder, FilterNullPolicy 
filterNullPolicy, List<PlanNode> children) {
++  public TimeJoinNode(
++      PlanNodeId id,
++      OrderBy mergeOrder,
++      FilterNullPolicy filterNullPolicy,
++      List<PlanNode> children) {
      super(id);
 -    this.mergeOrder = OrderBy.TIMESTAMP_ASC;
 +    this.mergeOrder = mergeOrder;
 +    this.filterNullPolicy = filterNullPolicy;
 +    this.children = children;
    }
  
 -  public TimeJoinNode(PlanNodeId id, PlanNode<TsBlock>... children) {
 -    super(id);
 -    this.children.addAll(Arrays.asList(children));
 +  @Override
 +  public List<PlanNode> getChildren() {
 +    return children;
 +  }
 +
 +  @Override
 +  public List<String> getOutputColumnNames() {
-     return children.stream().flatMap(child -> 
child.getOutputColumnNames().stream()).collect(Collectors.toList());
++    return children.stream()
++        .flatMap(child -> child.getOutputColumnNames().stream())
++        .collect(Collectors.toList());
    }
  
 -  public void addChild(PlanNode<TsBlock> child) {
 +  @Override
 +  public <R, C> R accept(PlanVisitor<R, C> visitor, C context) {
 +    return visitor.visitTimeJoin(this, context);
 +  }
 +
 +  public void addChild(PlanNode child) {
      this.children.add(child);
    }
  
diff --cc 
server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/sink/SinkNode.java
index 7221fda,f59effb..e01a2c4
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/sink/SinkNode.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/sink/SinkNode.java
@@@ -16,13 -16,13 +16,12 @@@
   * specific language governing permissions and limitations
   * under the License.
   */
 -package org.apache.iotdb.db.mpp.plan.node.sink;
 +package org.apache.iotdb.db.mpp.sql.planner.plan.node.sink;
  
--import org.apache.iotdb.db.mpp.common.TsBlock;
 -import org.apache.iotdb.db.mpp.plan.node.PlanNode;
 -import org.apache.iotdb.db.mpp.plan.node.PlanNodeId;
 +import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanNode;
 +import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanNodeId;
  
 -public abstract class SinkNode extends PlanNode<TsBlock> implements 
AutoCloseable {
 +public abstract class SinkNode extends PlanNode implements AutoCloseable {
  
    public SinkNode(PlanNodeId id) {
      super(id);
diff --cc 
server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/source/SeriesAggregateScanNode.java
index 6bdf1c4,80ea58f..ee01ac1
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/source/SeriesAggregateScanNode.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/source/SeriesAggregateScanNode.java
@@@ -16,18 -16,13 +16,19 @@@
   * specific language governing permissions and limitations
   * under the License.
   */
 -package org.apache.iotdb.db.mpp.plan.node.source;
 +package org.apache.iotdb.db.mpp.sql.planner.plan.node.source;
  
- import com.google.common.collect.ImmutableList;
  import org.apache.iotdb.db.mpp.common.GroupByTimeParameter;
 -import org.apache.iotdb.db.mpp.plan.node.PlanNodeId;
 +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.query.expression.unary.FunctionExpression;
  import org.apache.iotdb.tsfile.read.filter.basic.Filter;
  
++import com.google.common.collect.ImmutableList;
++
 +import java.util.List;
 +
  /**
   * SeriesAggregateOperator is responsible to do the aggregation calculation 
for one series. It will
   * read the target series and calculate the aggregation result by the 
aggregation digest or raw data
diff --cc 
server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/source/SeriesScanNode.java
index aa02fbe,ccfae5c..8858074
--- 
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
@@@ -16,18 -16,13 +16,19 @@@
   * specific language governing permissions and limitations
   * under the License.
   */
 -package org.apache.iotdb.db.mpp.plan.node.source;
 +package org.apache.iotdb.db.mpp.sql.planner.plan.node.source;
  
- import com.google.common.collect.ImmutableList;
  import org.apache.iotdb.db.metadata.path.PartialPath;
  import org.apache.iotdb.db.mpp.common.OrderBy;
 -import org.apache.iotdb.db.mpp.plan.node.PlanNodeId;
 +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.tsfile.read.filter.basic.Filter;
  
++import com.google.common.collect.ImmutableList;
++
 +import java.util.List;
 +
  /**
   * SeriesScanOperator is responsible for read data a specific series. When 
reading data, the
   * SeriesScanOperator can read the raw data batch by batch. And also, it can 
leverage the filter and
diff --cc 
server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/source/SourceNode.java
index ee25e5a,c83da97..551e9d2
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/source/SourceNode.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/source/SourceNode.java
@@@ -16,13 -16,13 +16,12 @@@
   * specific language governing permissions and limitations
   * under the License.
   */
 -package org.apache.iotdb.db.mpp.plan.node.source;
 +package org.apache.iotdb.db.mpp.sql.planner.plan.node.source;
  
--import org.apache.iotdb.db.mpp.common.TsBlock;
 -import org.apache.iotdb.db.mpp.plan.node.PlanNode;
 -import org.apache.iotdb.db.mpp.plan.node.PlanNodeId;
 +import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanNode;
 +import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanNodeId;
  
 -public abstract class SourceNode extends PlanNode<TsBlock> implements 
AutoCloseable {
 +public abstract class SourceNode extends PlanNode implements AutoCloseable {
  
    public SourceNode(PlanNodeId id) {
      super(id);
diff --cc server/src/main/java/org/apache/iotdb/db/mpp/sql/tree/Expression.java
index 45decc1,7a4107e..a75d326
--- a/server/src/main/java/org/apache/iotdb/db/mpp/sql/tree/Expression.java
+++ b/server/src/main/java/org/apache/iotdb/db/mpp/sql/tree/Expression.java
@@@ -16,8 -16,9 +16,6 @@@
   * specific language governing permissions and limitations
   * under the License.
   */
 -package org.apache.iotdb.db.mpp.common;
 +package org.apache.iotdb.db.mpp.sql.tree;
  
- public class Expression {
- 
 -public enum WithoutPolicy {
 -  CONTAINS_NULL,
 -  ALL_NULL
--}
++public class Expression {}

Reply via email to