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 {}
