This is an automated email from the ASF dual-hosted git repository. xingtanzjr pushed a commit to branch xingtanzjr/logical_to_distributed in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit 86bb0c27bc6ee29127690a8eb3259cb54879c91b Merge: b17c657 c83ccfa Author: Jinrui.Zhang <[email protected]> AuthorDate: Sun Mar 20 16:44:28 2022 +0800 Merge branch 'master' into xingtanzjr/logical_to_distributed .github/workflows/grafana-plugin.yml | 2 +- .github/workflows/greetings.yml | 4 - .github/workflows/sonar-coveralls.yml | 7 - .../org/apache/iotdb/db/qp/sql/IoTDBSqlParser.g4 | 8 +- .../java/org/apache/iotdb/cluster/ClientMain.java | 2 +- .../org/apache/iotdb/cluster/ClusterIoTDB.java | 20 +- .../iotdb/cluster/config/ClusterDescriptor.java | 2 +- .../iotdb/cluster/coordinator/Coordinator.java | 2 +- .../apache/iotdb/cluster/log/LogDispatcher.java | 4 +- .../iotdb/cluster/log/catchup/LogCatchUpTask.java | 2 +- .../iotdb/cluster/log/manage/RaftLogManager.java | 2 +- .../serializable/SyncLogDequeSerializer.java | 2 +- .../iotdb/cluster/metadata/CSchemaEngine.java | 2 +- .../iotdb/cluster/query/ClusterPlanExecutor.java | 2 +- .../query/last/ClusterLastQueryExecutor.java | 2 +- .../iotdb/cluster/server/ClusterRPCService.java | 10 +- .../cluster/server/ClusterRPCServiceMBean.java | 2 +- .../cluster/server/PullSnapshotHintService.java | 2 +- .../server/clusterinfo/ClusterInfoServer.java | 10 +- .../cluster/server/member/DataGroupMember.java | 6 +- .../cluster/server/member/MetaGroupMember.java | 6 +- .../iotdb/cluster/server/member/RaftMember.java | 8 +- .../cluster/server/raft/AbstractRaftService.java | 6 +- .../server/raft/DataRaftHeartBeatService.java | 8 +- .../iotdb/cluster/server/raft/DataRaftService.java | 8 +- .../server/raft/MetaRaftHeartBeatService.java | 8 +- .../iotdb/cluster/server/raft/MetaRaftService.java | 8 +- .../cluster/server/service/DataGroupEngine.java | 6 +- .../apache/iotdb/cluster/utils/PlanSerializer.java | 4 +- .../cluster/utils/nodetool/ClusterMonitor.java | 14 +- .../org/apache/iotdb/cluster/common/IoTDBTest.java | 4 +- .../iotdb/cluster/integration/SingleNodeTest.java | 2 +- .../iotdb/cluster/log/LogDispatcherTest.java | 2 +- .../cluster/log/applier/DataLogApplierTest.java | 8 +- .../cluster/log/catchup/LogCatchUpTaskTest.java | 2 +- .../manage/MetaSingleSnapshotLogManagerTest.java | 2 +- .../serializable/SyncLogDequeSerializerTest.java | 2 +- .../cluster/log/snapshot/DataSnapshotTest.java | 2 +- .../log/snapshot/MetaSimpleSnapshotTest.java | 5 +- .../cluster/log/snapshot/PullSnapshotTaskTest.java | 2 +- .../query/ClusterUDTFQueryExecutorTest.java | 2 +- .../iotdb/cluster/server/member/BaseMember.java | 2 +- .../resources/conf/iotdb-confignode.properties | 41 +- .../iotdb/confignode/conf/ConfigNodeConf.java | 116 +- .../iotdb/confignode/service/ConfigNode.java | 31 +- .../confignode/service/register/IService.java | 51 - .../confignode/service/register/JMXService.java | 105 - .../service/register/RegisterManager.java | 82 - .../confignode/service/startup/StartupCheck.java | 28 - .../confignode/service/startup/StartupChecks.java | 89 - .../service/thrift/server/ConfigNodeRPCServer.java | 88 +- ...rver.java => ConfigNodeRPCServerProcessor.java} | 6 +- .../thrift/server/ConfigNodeRPCServiceHandler.java | 52 + .../utils/ConfigNodeEnvironmentUtils.java | 9 +- docs/UserGuide/Query-Data/Overview.md | 27 +- docs/zh/UserGuide/Query-Data/Overview.md | 27 +- grafana-plugin/package.json | 4 +- grafana-plugin/src/componments/ControlValue.tsx | 5 +- grafana-plugin/src/componments/FromValue.tsx | 8 +- grafana-plugin/src/componments/SelectValue.tsx | 8 +- grafana-plugin/src/componments/WhereValue.tsx | 5 +- grafana-plugin/src/datasource.ts | 16 +- grafana-plugin/yarn.lock | 10529 +++++++++---------- .../apache/iotdb/db/integration/IoTDBDaemonIT.java | 2 +- .../iotdb/db/integration/IoTDBLargeDataIT.java | 2 +- .../iotdb/db/integration/IoTDBMetadataFetchIT.java | 2 +- .../iotdb/db/integration/IoTDBMultiSeriesIT.java | 2 +- .../db/integration/IoTDBRecoverUnclosedIT.java | 2 +- .../iotdb/db/integration/IoTDBUDFManagementIT.java | 6 +- .../iotdb/session/IoTDBSessionComplexIT.java | 2 +- .../iotdb/session/IoTDBSessionIteratorIT.java | 2 +- .../apache/iotdb/session/IoTDBSessionSimpleIT.java | 2 +- .../session/IoTDBSessionSyntaxConventionIT.java | 2 +- integration/src/test/resources/logback.xml | 4 +- {iotdb-commons => node-commons}/pom.xml | 20 +- .../apache/iotdb/commons/ServerCommandLine.java | 0 .../apache/iotdb/commons}/concurrent/HashLock.java | 2 +- .../concurrent/IoTDBDaemonThreadFactory.java | 2 +- .../IoTDBDefaultThreadExceptionHandler.java | 2 +- .../concurrent/IoTDBThreadPoolFactory.java | 10 +- .../commons}/concurrent/IoTThreadFactory.java | 2 +- .../iotdb/commons}/concurrent/ThreadName.java | 2 +- .../iotdb/commons}/concurrent/WrappedRunnable.java | 2 +- .../concurrent/threadpool/IThreadPoolMBean.java | 2 +- .../WrappedScheduledExecutorService.java | 6 +- .../WrappedScheduledExecutorServiceMBean.java | 2 +- .../WrappedSingleThreadExecutorService.java | 6 +- .../WrappedSingleThreadExecutorServiceMBean.java | 2 +- .../WrappedSingleThreadScheduledExecutor.java | 6 +- .../WrappedSingleThreadScheduledExecutorMBean.java | 2 +- .../threadpool/WrappedThreadPoolExecutor.java | 8 +- .../threadpool/WrappedThreadPoolExecutorMBean.java | 2 +- .../apache/iotdb/commons}/conf/IoTDBConstant.java | 2 +- .../iotdb/commons}/exception/IoTDBException.java | 2 +- .../commons}/exception/ShutdownException.java | 2 +- .../iotdb/commons}/exception/StartupException.java | 2 +- .../exception/runtime/RPCServiceException.java | 2 +- .../apache/iotdb/commons/hash/APHashExecutor.java | 0 .../iotdb/commons/hash/BKDRHashExecutor.java | 0 .../commons/hash/DeviceGroupHashExecutor.java | 0 .../apache/iotdb/commons/hash/JSHashExecutor.java | 0 .../iotdb/commons/hash/SDBMHashExecutor.java | 0 .../service/AbstractThriftServiceThread.java | 35 +- .../apache/iotdb/commons}/service/IService.java | 6 +- .../apache/iotdb/commons}/service/JMXService.java | 4 +- .../iotdb/commons}/service/RegisterManager.java | 6 +- .../apache/iotdb/commons}/service/ServiceType.java | 9 +- .../iotdb/commons}/service/StartupCheck.java | 4 +- .../iotdb/commons}/service/StartupChecks.java | 10 +- .../iotdb/commons/service}/ThriftService.java | 10 +- .../iotdb/commons/service/ThriftServiceThread.java | 89 + .../apache/iotdb/commons/utils/JVMCommonUtils.java | 81 + .../org/apache/iotdb/commons/utils/TestOnly.java | 0 .../IoTDBDefaultThreadExceptionHandlerTest.java | 2 +- .../iotdb/commons}/IoTDBThreadPoolFactoryTest.java | 5 +- pom.xml | 2 +- .../org/apache/iotdb/db/auth/AuthorityChecker.java | 2 +- .../iotdb/db/auth/authorizer/BasicAuthorizer.java | 8 +- .../iotdb/db/auth/role/BasicRoleManager.java | 2 +- .../iotdb/db/auth/role/LocalFileRoleAccessor.java | 2 +- .../iotdb/db/auth/user/BasicUserManager.java | 2 +- .../iotdb/db/auth/user/LocalFileUserAccessor.java | 2 +- .../java/org/apache/iotdb/db/conf/IoTDBConfig.java | 1 + .../org/apache/iotdb/db/conf/IoTDBConfigCheck.java | 1 + .../org/apache/iotdb/db/conf/IoTDBDescriptor.java | 5 +- .../db/conf/directories/DirectoryManager.java | 2 +- .../directories/strategy/DirectoryStrategy.java | 4 +- .../strategy/MaxDiskUsableSpaceFirstStrategy.java | 6 +- .../MinFolderOccupiedSpaceFirstStrategy.java | 6 +- .../strategy/RandomOnDiskUsableSpaceStrategy.java | 4 +- .../directories/strategy/SequenceStrategy.java | 6 +- .../db/conf/rest/IoTDBRestServiceDescriptor.java | 2 +- .../org/apache/iotdb/db/engine/StorageEngine.java | 10 +- .../db/engine/cache/CacheHitRatioMonitor.java | 10 +- .../db/engine/cache/TimeSeriesMetadataCache.java | 2 +- .../engine/compaction/CompactionTaskManager.java | 12 +- .../db/engine/compaction/CompactionUtils.java | 2 +- .../db/engine/compaction/TsFileIdentifier.java | 2 +- .../CrossSpaceCompactionExceptionHandler.java | 2 +- .../RewriteCrossSpaceCompactionSelector.java | 2 +- .../task/RewriteCrossCompactionRecoverTask.java | 2 +- .../task/RewriteCrossSpaceCompactionTask.java | 2 +- .../inner/AbstractInnerSpaceCompactionTask.java | 2 +- .../InnerSpaceCompactionExceptionHandler.java | 2 +- .../SizeTieredCompactionRecoverTask.java | 2 +- .../sizetiered/SizeTieredCompactionSelector.java | 2 +- .../inner/sizetiered/SizeTieredCompactionTask.java | 2 +- .../inner/utils/InnerSpaceCompactionUtils.java | 2 +- .../compaction/task/AbstractCompactionTask.java | 2 +- .../compaction/task/CompactionRecoverTask.java | 2 +- .../utils/log/CompactionLogAnalyzer.java | 2 +- .../iotdb/db/engine/cq/ContinuousQueryService.java | 8 +- .../iotdb/db/engine/cq/ContinuousQueryTask.java | 4 +- .../engine/cq/ContinuousQueryTaskPoolManager.java | 4 +- .../apache/iotdb/db/engine/flush/FlushManager.java | 10 +- .../engine/flush/pool/FlushSubTaskPoolManager.java | 4 +- .../db/engine/flush/pool/FlushTaskPoolManager.java | 4 +- .../apache/iotdb/db/engine/settle/SettleTask.java | 2 +- .../db/engine/storagegroup/TsFileManager.java | 2 +- .../engine/storagegroup/TsFileNameGenerator.java | 4 +- .../db/engine/storagegroup/TsFileResource.java | 2 +- .../storagegroup/VirtualStorageGroupProcessor.java | 8 +- .../virtualSg/StorageGroupManager.java | 2 +- .../service/TriggerRegistrationService.java | 22 +- .../iotdb/db/engine/upgrade/UpgradeTask.java | 2 +- .../iotdb/db/exception/ConfigurationException.java | 1 + .../iotdb/db/exception/LoadFileException.java | 1 + .../apache/iotdb/db/exception/MergeException.java | 1 + .../db/exception/QueryIdNotExsitException.java | 1 + .../exception/QueryInBatchStatementException.java | 1 + .../iotdb/db/exception/StorageEngineException.java | 1 + .../exception/StorageGroupProcessorException.java | 1 + .../db/exception/SyncConnectionException.java | 1 + .../SyncDeviceOwnerConflictException.java | 1 + .../iotdb/db/exception/SystemCheckException.java | 1 + .../db/exception/TsFileProcessorException.java | 1 + .../iotdb/db/exception/WriteProcessException.java | 1 + .../db/exception/index/IndexManagerException.java | 2 +- .../db/exception/metadata/MetadataException.java | 2 +- .../exception/query/LogicalOperatorException.java | 2 +- .../exception/query/LogicalOptimizeException.java | 2 +- .../db/exception/query/QueryProcessException.java | 2 +- .../{runtime => sql}/SQLParserException.java | 2 +- .../StatementAnalyzeException.java} | 21 +- .../apache/iotdb/db/metadata/MetadataConstant.java | 2 +- .../org/apache/iotdb/db/metadata/SchemaEngine.java | 4 +- .../db/metadata/StorageGroupSchemaManager.java | 4 +- .../org/apache/iotdb/db/metadata/mnode/MNode.java | 2 +- .../iotdb/db/metadata/mtree/MTreeAboveSG.java | 4 +- .../iotdb/db/metadata/mtree/MTreeBelowSG.java | 2 +- .../db/metadata/mtree/traverser/Traverser.java | 6 +- .../iotdb/db/metadata/path/MeasurementPath.java | 2 +- .../apache/iotdb/db/metadata/path/PartialPath.java | 4 +- .../db/metadata/template/TemplateManager.java | 2 +- .../db/metadata/upgrade/MetadataUpgrader.java | 4 +- .../iotdb/db/metadata/utils/MetaFormatUtils.java | 10 +- .../apache/iotdb/db/metadata/utils/MetaUtils.java | 2 +- .../org/apache/iotdb/db/mpp/common/Analysis.java | 48 +- .../org/apache/iotdb/db/mpp/common/DataRegion.java | 39 +- .../iotdb/db/mpp/common/DataRegionTimeSlice.java | 5 +- .../mpp/common/{TreeNode.java => FragmentId.java} | 40 +- .../{InstanceId.java => FragmentInstanceId.java} | 4 +- .../{QueryContext.java => MPPQueryContext.java} | 2 +- .../org/apache/iotdb/db/mpp/common/OrderBy.java | 27 - .../org/apache/iotdb/db/mpp/common/QueryId.java | 90 +- .../apache/iotdb/db/mpp/common/SchemaRegion.java | 10 +- .../org/apache/iotdb/db/mpp/common/TsBlock.java | 25 + .../FragmentInfo.java} | 33 +- ...ceContext.java => FragmentInstanceContext.java} | 35 +- .../db/mpp/execution/FragmentInstanceState.java | 68 + .../iotdb/db/mpp/execution/FragmentState.java | 71 + .../iotdb/db/mpp/execution/QueryExecution.java | 12 +- .../mpp/execution/executor/InstanceExecutor.java | 39 - .../db/mpp/execution/executor/InstanceHandle.java | 31 - .../ClusterScheduler.java} | 32 +- .../scheduler/IScheduler.java} | 33 +- .../scheduler/StandaloneScheduler.java} | 39 +- .../iotdb/db/mpp/operator/OperatorContext.java | 4 + .../db/mpp/operator/process/LimitOperator.java | 32 +- .../iotdb/db/mpp/operator/sink/SinkOperator.java | 2 +- .../db/mpp/sql/planner/LocalExecutionPlanner.java | 14 +- .../sql/planner/optimization/PlanOptimizer.java | 5 +- .../mpp/sql/planner/plan/DistributedQueryPlan.java | 4 +- .../mpp/sql/planner/plan/DistributionPlanner.java | 142 +- .../db/mpp/sql/planner/plan/FragmentInstance.java | 2 + .../mpp/sql/planner/plan/FragmentInstanceId.java | 30 - .../db/mpp/sql/planner/plan/LogicalPlanner.java | 7 +- .../db/mpp/sql/planner/plan/LogicalQueryPlan.java | 8 +- .../sql/planner/plan/node/PlanNodeAllocator.java | 11 +- .../db/mpp/sql/planner/plan/node/PlanNodeUtil.java | 27 +- .../planner/plan/node/SimplePlanNodeRewriter.java | 31 +- .../planner/plan/node/process/DeviceMergeNode.java | 2 +- .../sql/planner/plan/node/process/OffsetNode.java | 2 +- .../sql/planner/plan/node/process/SortNode.java | 2 +- .../planner/plan/node/process/TimeJoinNode.java | 8 +- .../planner/plan/node/source/SeriesScanNode.java | 7 +- .../influxdb/meta/InfluxDBMetaManager.java | 2 +- .../apache/iotdb/db/protocol/rest/RestService.java | 6 +- .../db/protocol/rest/handler/ExceptionHandler.java | 2 +- .../protocol/rest/impl/GrafanaApiServiceImpl.java | 2 +- .../db/protocol/rest/impl/PingApiServiceImpl.java | 2 +- .../db/protocol/rest/impl/RestApiServiceImpl.java | 2 +- .../main/java/org/apache/iotdb/db/qp/Planner.java | 2 +- .../apache/iotdb/db/qp/executor/PlanExecutor.java | 76 +- .../db/qp/logical/crud/BasicFunctionOperator.java | 2 +- .../db/qp/logical/crud/BasicOperatorType.java | 4 +- .../db/qp/logical/crud/DeleteDataOperator.java | 2 +- .../db/qp/logical/crud/FillQueryOperator.java | 2 +- .../iotdb/db/qp/logical/crud/InsertOperator.java | 2 +- .../sys/CreateAlignedTimeSeriesOperator.java | 2 +- .../iotdb/db/qp/physical/crud/InsertRowPlan.java | 2 +- .../iotdb/db/qp/physical/sys/SetTemplatePlan.java | 2 +- .../db/qp/physical/sys/UnsetTemplatePlan.java | 2 +- .../apache/iotdb/db/qp/sql/IoTDBSqlVisitor.java | 15 +- .../iotdb/db/qp/strategy/LogicalGenerator.java | 4 +- .../qp/strategy/optimizer/ConcatPathOptimizer.java | 2 +- .../iotdb/db/qp/utils/GroupByLevelController.java | 2 +- .../iotdb/db/query/control/QueryTimeManager.java | 6 +- .../iotdb/db/query/control/SessionManager.java | 2 +- .../db/query/control/SessionTimeoutManager.java | 2 +- .../db/query/dataset/NonAlignEngineDataSet.java | 2 +- .../dataset/RawQueryDataSetWithoutValueFilter.java | 2 +- .../iotdb/db/query/dataset/ShowDevicesDataSet.java | 6 +- .../db/query/dataset/ShowTimeseriesDataSet.java | 16 +- .../db/query/executor/AggregationExecutor.java | 2 +- .../iotdb/db/query/executor/LastQueryExecutor.java | 6 +- .../query/expression/unary/FunctionExpression.java | 2 +- .../iotdb/db/query/pool/QueryTaskManager.java | 4 +- .../db/query/pool/RawQueryReadTaskPoolManager.java | 4 +- .../row/SerializableRowRecordList.java | 2 +- .../datastructure/tv/SerializableBinaryTVList.java | 2 +- .../tv/SerializableBooleanTVList.java | 2 +- .../datastructure/tv/SerializableDoubleTVList.java | 2 +- .../datastructure/tv/SerializableFloatTVList.java | 2 +- .../datastructure/tv/SerializableIntTVList.java | 2 +- .../datastructure/tv/SerializableLongTVList.java | 2 +- .../udf/service/TemporaryQueryDataFileService.java | 6 +- .../query/udf/service/UDFClassLoaderManager.java | 6 +- .../query/udf/service/UDFRegistrationService.java | 6 +- .../org/apache/iotdb/db/rescon/SystemInfo.java | 2 +- .../iotdb/db/service/InfluxDBRPCService.java | 9 +- .../java/org/apache/iotdb/db/service/IoTDB.java | 9 +- .../org/apache/iotdb/db/service/MQTTService.java | 2 + .../org/apache/iotdb/db/service/RPCService.java | 9 +- .../apache/iotdb/db/service/RPCServiceMBean.java | 2 +- .../org/apache/iotdb/db/service/SettleService.java | 6 +- .../org/apache/iotdb/db/service/StaticResps.java | 6 +- .../org/apache/iotdb/db/service/UpgradeSevice.java | 4 +- .../db/service/basic/QueryFrequencyRecorder.java | 2 +- .../iotdb/db/service/basic/ServiceProvider.java | 2 +- .../iotdb/db/service/metrics/MetricsService.java | 10 +- .../db/service/metrics/MetricsServiceMBean.java | 2 +- .../db/service/thrift/impl/TSServiceImpl.java | 4 +- .../iotdb/db/sql/constant/FilterConstant.java | 102 + .../iotdb/db/sql/constant/StatementType.java | 134 + .../org/apache/iotdb/db/sql/parser/ASTVisitor.java | 1368 +++ .../iotdb/db/sql/parser/StatementGenerator.java | 184 + .../statement/AggregationQueryStatement.java} | 35 +- .../statement/FillQueryStatement.java} | 20 +- .../statement/GroupByFillQueryStatement.java} | 28 +- .../db/sql/statement/GroupByQueryStatement.java | 29 +- .../statement/LastQueryStatement.java} | 14 +- .../iotdb/db/sql/statement/QueryStatement.java | 211 + .../statement/ShowDevicesStatement.java} | 16 +- .../statement/ShowTimeSeriesStatement.java} | 16 +- .../statement/Statement.java} | 39 +- .../iotdb/db/sql/statement/UDAFQueryStatement.java | 10 +- .../iotdb/db/sql/statement/UDTFQueryStatement.java | 10 +- .../db/sql/statement/component/FillComponent.java | 36 +- .../db/sql/statement/component/FromComponent.java | 29 +- .../statement/component/GroupByLevelComponent.java | 31 +- .../statement/component/GroupByTimeComponent.java | 99 + .../iotdb/db/sql/statement/component/OrderBy.java | 9 +- .../db/sql/statement/component/ResultColumn.java | 168 + .../sql/statement/component/ResultSetFormat.java | 10 +- .../sql/statement/component/SelectComponent.java | 99 + .../statement/component/WhereCondition.java} | 27 +- .../db/sql/statement/component/WithoutPolicy.java | 59 + .../statement/filter/BasicFilterType.java} | 18 +- .../statement/filter/BasicFunctionFilter.java} | 35 +- .../statement/filter/FunctionFilter.java} | 36 +- .../iotdb/db/sql/statement/filter/InFilter.java | 201 + .../iotdb/db/sql/statement/filter/LikeFilter.java | 134 + .../iotdb/db/sql/statement/filter/QueryFilter.java | 295 + .../db/sql/statement/filter/RegexpFilter.java | 134 + .../iotdb/db/sync/conf/SyncSenderDescriptor.java | 2 +- .../iotdb/db/sync/receiver/SyncServerManager.java | 10 +- .../db/sync/receiver/SyncServerManagerMBean.java | 2 +- .../db/sync/receiver/load/FileLoaderManager.java | 4 +- .../db/sync/receiver/transfer/SyncServiceImpl.java | 2 +- .../db/sync/sender/manage/SyncFileManager.java | 2 +- .../iotdb/db/sync/sender/transfer/SyncClient.java | 4 +- .../org/apache/iotdb/db/tools/TsFileSplitTool.java | 2 +- .../java/org/apache/iotdb/db/utils/AuthUtils.java | 2 +- .../org/apache/iotdb/db/utils/CommonUtils.java | 57 - .../apache/iotdb/db/utils/ErrorHandlingUtils.java | 4 +- .../java/org/apache/iotdb/db/utils/MemUtils.java | 2 +- .../org/apache/iotdb/db/utils/OpenFileNumUtil.java | 2 +- .../org/apache/iotdb/db/utils/ThreadUtils.java | 2 +- .../windowing/runtime/WindowEvaluationTask.java | 2 +- .../runtime/WindowEvaluationTaskPoolManager.java | 6 +- .../writelog/manager/MultiFileLogNodeManager.java | 8 +- .../db/writelog/node/ExclusiveWriteLogNode.java | 4 +- .../iotdb/db/writelog/recover/LogReplayer.java | 2 +- .../apache/iotdb/db/conf/IoTDBDescriptorTest.java | 2 + .../strategy/DirectoryStrategyTest.java | 24 +- .../iotdb/db/engine/cache/ChunkCacheTest.java | 2 +- .../db/engine/compaction/CompactionUtilsTest.java | 2 +- .../cross/CrossSpaceCompactionExceptionTest.java | 2 +- .../db/engine/compaction/cross/MergeTest.java | 2 +- .../engine/compaction/cross/MergeUpgradeTest.java | 2 +- .../cross/RewriteCompactionFileSelectorTest.java | 2 +- .../RewriteCrossSpaceCompactionRecoverTest.java | 2 +- .../cross/RewriteCrossSpaceCompactionTest.java | 4 +- .../inner/AbstractInnerSpaceCompactionTest.java | 4 +- .../inner/InnerCompactionMoreDataTest.java | 4 +- .../compaction/inner/InnerCompactionTest.java | 4 +- .../inner/InnerSpaceCompactionUtilsOldTest.java | 2 +- .../SizeTieredCompactionRecoverTest.java | 2 +- .../inner/sizetiered/SizeTieredCompactionTest.java | 4 +- ...eCrossSpaceCompactionRecoverCompatibleTest.java | 2 +- .../recover/SizeTieredCompactionRecoverTest.java | 2 +- .../compaction/utils/CompactionClearUtils.java | 2 +- .../utils/CompactionFileGeneratorUtils.java | 2 +- .../storagegroup/StorageGroupProcessorTest.java | 2 +- .../iotdb/db/engine/storagegroup/TTLTest.java | 2 +- .../db/engine/storagegroup/TsFileManagerTest.java | 2 +- .../db/metadata/upgrade/MetadataUpgradeTest.java | 2 +- .../db/mpp/sql/plan/DistributionPlannerTest.java | 101 +- .../iotdb/db/protocol/rest/IoTDBRestServiceIT.java | 2 +- .../java/org/apache/iotdb/db/qp/PlannerTest.java | 4 +- .../iotdb/db/qp/logical/LogicalPlanSmallTest.java | 2 +- .../iotdb/db/qp/physical/PhysicalPlanTest.java | 2 +- .../apache/iotdb/db/qp/sql/ASTVisitorTest.java} | 32 +- .../iotdb/db/qp/sql/IoTDBsqlVisitorTest.java | 2 +- .../query/reader/series/SeriesReaderTestUtil.java | 2 +- .../iotdb/db/rescon/ResourceManagerTest.java | 4 +- .../iotdb/db/sql/StatementGeneratorTest.java | 75 + .../db/sync/receiver/load/FileLoaderTest.java | 2 +- .../recover/SyncReceiverLogAnalyzerTest.java | 2 +- .../db/sync/sender/manage/SyncFileManagerTest.java | 2 +- .../sender/recover/SyncSenderLogAnalyzerTest.java | 4 +- .../sync/sender/recover/SyncSenderLoggerTest.java | 2 +- .../db/sync/sender/transfer/SyncClientTest.java | 2 +- .../apache/iotdb/db/tools/IoTDBWatermarkTest.java | 2 +- .../iotdb/db/writelog/IoTDBLogFileSizeTest.java | 2 +- .../recover/RecoverResourceFromReaderTest.java | 2 +- .../writelog/recover/UnseqTsFileRecoverTest.java | 2 +- server/src/test/resources/logback.xml | 4 +- .../session/IoTDBSessionDisableMemControlIT.java | 2 +- .../iotdb/session/IoTDBSessionVectorInsertIT.java | 2 +- .../java/org/apache/iotdb/session/SessionTest.java | 2 +- .../apache/iotdb/session/pool/SessionPoolTest.java | 2 +- .../apache/iotdb/session/template/TemplateUT.java | 2 +- .../apache/iotdb/spark/db/EnvironmentUtils.java | 2 +- .../org/apache/iotdb/spark/db/IoTDBTest.scala | 3 +- .../org/apache/iotdb/spark/db/IoTDBWriteTest.scala | 3 +- .../iotdb/spark/db/unit/DataFrameToolsTest.scala | 4 +- .../zeppelin/iotdb/IoTDBInterpreterTest.java | 2 +- 399 files changed, 10163 insertions(+), 7463 deletions(-) diff --cc server/src/main/java/org/apache/iotdb/db/mpp/common/Analysis.java index 85ff111,88d8884..2d2a397 --- a/server/src/main/java/org/apache/iotdb/db/mpp/common/Analysis.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/common/Analysis.java @@@ -25,36 -22,12 +25,36 @@@ import java.util.* /** Analysis used for planning a query. TODO: This class may need to store more info for a query. */ public class Analysis { - // Description for each series. Such as dataType, existence ++ // Description for each series. Such as dataType, existence + - // Data distribution info for each series. Series -> [VSG, VSG] ++ // Data distribution info for each series. Series -> [VSG, VSG] + - // Map<PartialPath, List<FullPath>> Used to remove asterisk ++ // Map<PartialPath, List<FullPath>> Used to remove asterisk - // Statement - private String statement; - private Statement statement; ++ // Statement ++ private String statement; - // DataPartitionInfo - private Map<String, Map<DataRegionTimeSlice, List<DataRegion>>> dataPartitionInfo; - public Analysis() {} ++ // DataPartitionInfo ++ private Map<String, Map<DataRegionTimeSlice, List<DataRegion>>> dataPartitionInfo; + - // SchemaPartitionInfo - private Map<String, List<SchemaRegion>> schemaPartitionInfo; ++ // SchemaPartitionInfo ++ private Map<String, List<SchemaRegion>> schemaPartitionInfo; + - - public Set<DataRegion> getPartitionInfo(PartialPath seriesPath, Filter timefilter) { - if (timefilter == null) { - //TODO: (xingtanzjr) we need to have a method to get the deviceGroup by device - String deviceGroup = seriesPath.getDevice(); - Set<DataRegion> result = new HashSet<>(); - this.dataPartitionInfo.get(deviceGroup).values().forEach(result::addAll); - return result; - } else { - //TODO: (xingtanzjr) complete this branch - return null; - } ++ public Set<DataRegion> getPartitionInfo(PartialPath seriesPath, Filter timefilter) { ++ if (timefilter == null) { ++ // TODO: (xingtanzjr) we need to have a method to get the deviceGroup by device ++ String deviceGroup = seriesPath.getDevice(); ++ Set<DataRegion> result = new HashSet<>(); ++ this.dataPartitionInfo.get(deviceGroup).values().forEach(result::addAll); ++ return result; ++ } else { ++ // TODO: (xingtanzjr) complete this branch ++ return null; + } ++ } - public void setDataPartitionInfo(Map<String, Map<DataRegionTimeSlice, List<DataRegion>>> dataPartitionInfo) { - this.dataPartitionInfo = dataPartitionInfo; - } - public void setStatement(Statement rewrittenStatement) { - this.statement = rewrittenStatement; ++ public void setDataPartitionInfo( ++ Map<String, Map<DataRegionTimeSlice, List<DataRegion>>> dataPartitionInfo) { ++ this.dataPartitionInfo = dataPartitionInfo; + } } diff --cc server/src/main/java/org/apache/iotdb/db/mpp/common/DataRegion.java index 5618e1f,612e930..04100e1 --- a/server/src/main/java/org/apache/iotdb/db/mpp/common/DataRegion.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/common/DataRegion.java @@@ -16,34 -16,33 +16,35 @@@ * specific language governing permissions and limitations * under the License. */ -package org.apache.iotdb.db.mpp.sql.planner.plan.node.source; -import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanNode; -import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanNodeId; +package org.apache.iotdb.db.mpp.common; -import java.util.List; - -/** Not implemented in current version. */ -public class CsvSourceNode extends SourceNode { - - public CsvSourceNode(PlanNodeId id) { - super(id); +/** - * This class is used to represent the data partition info including the DataRegionId and physical node IP address ++ * This class is used to represent the data partition info including the DataRegionId and physical ++ * node IP address + */ - //TODO: (xingtanzjr) This class should be substituted with the class defined in Consensus level ++// TODO: (xingtanzjr) This class should be substituted with the class defined in Consensus level +public class DataRegion { - private Integer dataRegionId; - private String endpoint; ++ private Integer dataRegionId; ++ private String endpoint; + - public DataRegion(Integer dataRegionId, String endpoint) { - this.dataRegionId = dataRegionId; - this.endpoint = endpoint; - } ++ public DataRegion(Integer dataRegionId, String endpoint) { ++ this.dataRegionId = dataRegionId; ++ this.endpoint = endpoint; + } - public int hashCode() { - return dataRegionId.hashCode(); - } - @Override - public List<PlanNode> getChildren() { - return null; ++ public int hashCode() { ++ return dataRegionId.hashCode(); + } - public boolean equals(Object obj) { - if (obj instanceof DataRegion) { - return this.dataRegionId.equals(((DataRegion)obj).dataRegionId); - } - return false; - @Override - public List<String> getOutputColumnNames() { - return null; ++ public boolean equals(Object obj) { ++ if (obj instanceof DataRegion) { ++ return this.dataRegionId.equals(((DataRegion) obj).dataRegionId); + } ++ return false; + } - public String toString() { - return String.format("%s/%d",this.endpoint, this.dataRegionId); - } - @Override - public void close() throws Exception {} - - @Override - public void open() throws Exception {} ++ public String toString() { ++ return String.format("%s/%d", this.endpoint, this.dataRegionId); ++ } } diff --cc server/src/main/java/org/apache/iotdb/db/mpp/common/DataRegionTimeSlice.java index 52f48cc,a8d9b94..65dcd39 --- a/server/src/main/java/org/apache/iotdb/db/mpp/common/DataRegionTimeSlice.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/common/DataRegionTimeSlice.java @@@ -16,10 -16,12 +16,9 @@@ * specific language governing permissions and limitations * under the License. */ - package org.apache.iotdb.db.mpp.common; - //TODO: (xingtanzjr) This class should be substituted with the class defined in Consensus level -/** The traversal order for operators by timestamp */ -public enum OrderBy { - TIMESTAMP_ASC, - TIMESTAMP_DESC, - DEVICE_NAME_ASC, - DEVICE_NAME_DESC, ++// TODO: (xingtanzjr) This class should be substituted with the class defined in Consensus level +public class DataRegionTimeSlice { - long startTimestamp; ++ long startTimestamp; } diff --cc server/src/main/java/org/apache/iotdb/db/mpp/common/MPPQueryContext.java index 2f6715d,2f6715d..4a29263 --- a/server/src/main/java/org/apache/iotdb/db/mpp/common/MPPQueryContext.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/common/MPPQueryContext.java @@@ -22,7 -22,7 +22,7 @@@ package org.apache.iotdb.db.mpp.common * This class is used to record the context of a query including QueryId, query statement, session * info and so on */ --public class QueryContext { ++public class MPPQueryContext { private String statement; private QueryId queryId; private QuerySession session; diff --cc server/src/main/java/org/apache/iotdb/db/mpp/common/SchemaRegion.java index 8352f10,2f6715d..a7ffe29 --- a/server/src/main/java/org/apache/iotdb/db/mpp/common/SchemaRegion.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/common/SchemaRegion.java @@@ -19,11 -19,11 +19,11 @@@ package org.apache.iotdb.db.mpp.common; /** - * This class is used to represent the schema partition info including the DataRegionId and physical node IP address - * This class is used to record the context of a query including QueryId, query statement, session - * info and so on ++ * This class is used to represent the schema partition info including the DataRegionId and physical ++ * node IP address */ - //TODO: (xingtanzjr) This class should be substituted with the class defined in Consensus level -public class QueryContext { - private String statement; - private QueryId queryId; - private QuerySession session; ++// TODO: (xingtanzjr) This class should be substituted with the class defined in Consensus level +public class SchemaRegion { - private Integer DataRegionId; - private String endpoint; ++ private Integer DataRegionId; ++ private String endpoint; } diff --cc server/src/main/java/org/apache/iotdb/db/mpp/execution/QueryExecution.java index 9852350,62a2c65..ad2b176 --- a/server/src/main/java/org/apache/iotdb/db/mpp/execution/QueryExecution.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/execution/QueryExecution.java @@@ -19,7 -19,9 +19,9 @@@ package org.apache.iotdb.db.mpp.execution; import org.apache.iotdb.db.mpp.common.Analysis; --import org.apache.iotdb.db.mpp.common.QueryContext; ++import org.apache.iotdb.db.mpp.common.MPPQueryContext; + import org.apache.iotdb.db.mpp.execution.scheduler.ClusterScheduler; + import org.apache.iotdb.db.mpp.execution.scheduler.IScheduler; import org.apache.iotdb.db.mpp.sql.planner.optimization.PlanOptimizer; import org.apache.iotdb.db.mpp.sql.planner.plan.*; @@@ -33,8 -35,8 +35,8 @@@ import java.util.List * corresponding physical nodes. 3. Collect and monitor the progress/states of this query. */ public class QueryExecution { -- private QueryContext context; - private QueryScheduler scheduler; ++ private MPPQueryContext context; + private IScheduler scheduler; private QueryStateMachine stateMachine; private List<PlanOptimizer> planOptimizers; @@@ -45,7 -47,7 +47,7 @@@ private List<PlanFragment> fragments; private List<FragmentInstance> fragmentInstances; -- public QueryExecution(QueryContext context) { ++ public QueryExecution(MPPQueryContext context) { this.context = context; } diff --cc server/src/main/java/org/apache/iotdb/db/mpp/operator/sink/SinkOperator.java index c3f16ca,c3f16ca..d362c1d --- a/server/src/main/java/org/apache/iotdb/db/mpp/operator/sink/SinkOperator.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/operator/sink/SinkOperator.java @@@ -7,7 -7,7 +7,7 @@@ * "License"); you may not use this file except in compliance * with the License. You may obtain a copy of the License at * -- * http://www.apache.org/licenses/LICENSE-2.0 ++ * http://www.apache.org/licenses/LICENSE-2.0 * * Unless required by applicable law or agreed to in writing, * software distributed under the License is distributed on an diff --cc server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/optimization/PlanOptimizer.java index ff2baae,b091b91..94ada95 --- a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/optimization/PlanOptimizer.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/optimization/PlanOptimizer.java @@@ -18,10 -18,9 +18,9 @@@ */ package org.apache.iotdb.db.mpp.sql.planner.optimization; --import org.apache.iotdb.db.mpp.common.QueryContext; - import org.apache.iotdb.db.mpp.common.TsBlock; ++import org.apache.iotdb.db.mpp.common.MPPQueryContext; import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanNode; public interface PlanOptimizer { -- PlanNode optimize(PlanNode plan, QueryContext context); ++ PlanNode optimize(PlanNode plan, MPPQueryContext context); } diff --cc server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/DistributedQueryPlan.java index 50835dd,50835dd..8cf0341 --- a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/DistributedQueryPlan.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/DistributedQueryPlan.java @@@ -18,13 -18,13 +18,13 @@@ */ package org.apache.iotdb.db.mpp.sql.planner.plan; --import org.apache.iotdb.db.mpp.common.QueryContext; ++import org.apache.iotdb.db.mpp.common.MPPQueryContext; import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanNode; import java.util.List; public class DistributedQueryPlan { -- private QueryContext context; ++ private MPPQueryContext context; private PlanNode rootNode; private PlanFragment rootFragment; diff --cc server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/DistributionPlanner.java index 9c55a46,fc428cf..9b146d5 --- a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/DistributionPlanner.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/DistributionPlanner.java @@@ -19,91 -19,17 +19,101 @@@ package org.apache.iotdb.db.mpp.sql.planner.plan; import org.apache.iotdb.db.mpp.common.Analysis; +import org.apache.iotdb.db.mpp.common.DataRegion; +import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanNode; +import org.apache.iotdb.db.mpp.sql.planner.plan.node.SimplePlanNodeRewriter; +import org.apache.iotdb.db.mpp.sql.planner.plan.node.process.TimeJoinNode; +import org.apache.iotdb.db.mpp.sql.planner.plan.node.source.SeriesAggregateScanNode; +import org.apache.iotdb.db.mpp.sql.planner.plan.node.source.SeriesScanNode; + +import java.util.*; public class DistributionPlanner { - private Analysis analysis; - private LogicalQueryPlan logicalPlan; + private Analysis analysis; + private LogicalQueryPlan logicalPlan; - public DistributionPlanner(Analysis analysis, LogicalQueryPlan logicalPlan) { - this.analysis = analysis; - this.logicalPlan = logicalPlan; - } - - public PlanNode rewriteSource() { - SourceRewriter rewriter = new SourceRewriter(); - return rewriter.visit(logicalPlan.getRootNode(), new DistributionPlanContext()); - } + public DistributionPlanner(Analysis analysis, LogicalQueryPlan logicalPlan) { + this.analysis = analysis; + this.logicalPlan = logicalPlan; + } - public DistributedQueryPlan planFragments() { - return null; - } ++ public PlanNode rewriteSource() { ++ SourceRewriter rewriter = new SourceRewriter(); ++ return rewriter.visit(logicalPlan.getRootNode(), new DistributionPlanContext()); ++ } + - private class SourceRewriter extends SimplePlanNodeRewriter<DistributionPlanContext> { - public PlanNode visitTimeJoin(TimeJoinNode node, DistributionPlanContext context) { - TimeJoinNode root = (TimeJoinNode) node.clone(); + public DistributedQueryPlan planFragments() { + return null; + } + - // Step 1: Get all source nodes. For the node which is not source, add it as the child of current TimeJoinNode - List<SeriesScanNode> sources = new ArrayList<>(); - for (PlanNode child : node.getChildren()) { - if (child instanceof SeriesScanNode) { - // If the child is SeriesScanNode, we need to check whether this node should be seperated into several splits. - SeriesScanNode handle = (SeriesScanNode) child; - Set<DataRegion> dataDistribution = analysis.getPartitionInfo(handle.getSeriesPath(), handle.getTimeFilter()); - // If the size of dataDistribution is m, this SeriesScanNode should be seperated into m SeriesScanNode. - for (DataRegion dataRegion : dataDistribution) { - SeriesScanNode split = (SeriesScanNode) handle.clone(); - split.setDataRegion(dataRegion); - sources.add(split); - } - } else if (child instanceof SeriesAggregateScanNode) { - //TODO: (xingtanzjr) We should do the same thing for SeriesAggregateScanNode. Consider to make SeriesAggregateScanNode - // and SeriesScanNode to derived from the same parent Class because they have similar process logic in many scenarios - } else { - // In a general logical query plan, the children of TimeJoinNode should only be SeriesScanNode or SeriesAggregateScanNode - // So this branch should not be touched. - root.addChild(visit(child, context)); - } - } ++ private class SourceRewriter extends SimplePlanNodeRewriter<DistributionPlanContext> { ++ public PlanNode visitTimeJoin(TimeJoinNode node, DistributionPlanContext context) { ++ TimeJoinNode root = (TimeJoinNode) node.clone(); + - // Step 2: For the source nodes, group them by the DataRegion. - Map<DataRegion, List<SeriesScanNode>> sourceGroup = new HashMap<>(); - sources.forEach(source -> { - List<SeriesScanNode> group = sourceGroup.containsKey(source.getDataRegion()) ? - sourceGroup.get(source.getDataRegion()) : new ArrayList<>(); - group.add(source); - sourceGroup.put(source.getDataRegion(), group); - }); ++ // Step 1: Get all source nodes. For the node which is not source, add it as the child of ++ // current TimeJoinNode ++ List<SeriesScanNode> sources = new ArrayList<>(); ++ for (PlanNode child : node.getChildren()) { ++ if (child instanceof SeriesScanNode) { ++ // If the child is SeriesScanNode, we need to check whether this node should be seperated ++ // into several splits. ++ SeriesScanNode handle = (SeriesScanNode) child; ++ Set<DataRegion> dataDistribution = ++ analysis.getPartitionInfo(handle.getSeriesPath(), handle.getTimeFilter()); ++ // If the size of dataDistribution is m, this SeriesScanNode should be seperated into m ++ // SeriesScanNode. ++ for (DataRegion dataRegion : dataDistribution) { ++ SeriesScanNode split = (SeriesScanNode) handle.clone(); ++ split.setDataRegion(dataRegion); ++ sources.add(split); ++ } ++ } else if (child instanceof SeriesAggregateScanNode) { ++ // TODO: (xingtanzjr) We should do the same thing for SeriesAggregateScanNode. Consider to ++ // make SeriesAggregateScanNode ++ // and SeriesScanNode to derived from the same parent Class because they have similar ++ // process logic in many scenarios ++ } else { ++ // In a general logical query plan, the children of TimeJoinNode should only be ++ // SeriesScanNode or SeriesAggregateScanNode ++ // So this branch should not be touched. ++ root.addChild(visit(child, context)); ++ } ++ } + - // Step 3: For the source nodes which belong to same data region, add a TimeJoinNode for them and make the - // new TimeJoinNode as the child of current TimeJoinNode - sourceGroup.forEach((dataRegion, seriesScanNodes) -> { - if (seriesScanNodes.size() == 1) { - root.addChild(seriesScanNodes.get(0)); - } else { - // We clone a TimeJoinNode from root to make the params to be consistent - TimeJoinNode parentOfGroup = (TimeJoinNode) root.clone(); - seriesScanNodes.forEach(parentOfGroup::addChild); - root.addChild(parentOfGroup); - } - }); ++ // Step 2: For the source nodes, group them by the DataRegion. ++ Map<DataRegion, List<SeriesScanNode>> sourceGroup = new HashMap<>(); ++ sources.forEach( ++ source -> { ++ List<SeriesScanNode> group = ++ sourceGroup.containsKey(source.getDataRegion()) ++ ? sourceGroup.get(source.getDataRegion()) ++ : new ArrayList<>(); ++ group.add(source); ++ sourceGroup.put(source.getDataRegion(), group); ++ }); + - return root; - } ++ // Step 3: For the source nodes which belong to same data region, add a TimeJoinNode for them ++ // and make the ++ // new TimeJoinNode as the child of current TimeJoinNode ++ sourceGroup.forEach( ++ (dataRegion, seriesScanNodes) -> { ++ if (seriesScanNodes.size() == 1) { ++ root.addChild(seriesScanNodes.get(0)); ++ } else { ++ // We clone a TimeJoinNode from root to make the params to be consistent ++ TimeJoinNode parentOfGroup = (TimeJoinNode) root.clone(); ++ seriesScanNodes.forEach(parentOfGroup::addChild); ++ root.addChild(parentOfGroup); ++ } ++ }); + - public PlanNode visit(PlanNode node, DistributionPlanContext context) { - return node.accept(this, context); - } ++ return root; + } + - private class DistributionPlanContext { - ++ public PlanNode visit(PlanNode node, DistributionPlanContext context) { ++ return node.accept(this, context); + } ++ } ++ ++ private class DistributionPlanContext {} } diff --cc server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/LogicalPlanner.java index 0a45b6f,0a45b6f..ad0937d --- a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/LogicalPlanner.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/LogicalPlanner.java @@@ -19,17 -19,17 +19,18 @@@ package org.apache.iotdb.db.mpp.sql.planner.plan; import org.apache.iotdb.db.mpp.common.Analysis; --import org.apache.iotdb.db.mpp.common.QueryContext; ++import org.apache.iotdb.db.mpp.common.MPPQueryContext; import org.apache.iotdb.db.mpp.sql.planner.optimization.PlanOptimizer; import java.util.List; public class LogicalPlanner { private Analysis analysis; -- private QueryContext context; ++ private MPPQueryContext context; private List<PlanOptimizer> optimizers; -- public LogicalPlanner(Analysis analysis, QueryContext context, List<PlanOptimizer> optimizers) { ++ public LogicalPlanner( ++ Analysis analysis, MPPQueryContext context, List<PlanOptimizer> optimizers) { this.analysis = analysis; this.context = context; this.optimizers = optimizers; diff --cc server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/LogicalQueryPlan.java index 788746a,666fbf3..664497d --- a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/LogicalQueryPlan.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/LogicalQueryPlan.java @@@ -18,7 -18,7 +18,7 @@@ */ package org.apache.iotdb.db.mpp.sql.planner.plan; --import org.apache.iotdb.db.mpp.common.QueryContext; ++import org.apache.iotdb.db.mpp.common.MPPQueryContext; import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanNode; /** @@@ -26,19 -26,6 +26,19 @@@ * plan node tree. */ public class LogicalQueryPlan { -- private QueryContext context; ++ private MPPQueryContext context; private PlanNode rootNode; + - public LogicalQueryPlan(QueryContext context, PlanNode rootNode) { ++ public LogicalQueryPlan(MPPQueryContext context, PlanNode rootNode) { + this.context = context; + this.rootNode = rootNode; + } + + public PlanNode getRootNode() { + return rootNode; + } + - public QueryContext getContext() { ++ public MPPQueryContext getContext() { + return context; + } } diff --cc server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/PlanNodeAllocator.java index edac517,7979b27..63c9a58 --- a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/PlanNodeAllocator.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/PlanNodeAllocator.java @@@ -16,13 -16,9 +16,14 @@@ * specific language governing permissions and limitations * under the License. */ -package org.apache.iotdb.db.mpp.common; -public enum FilterNullPolicy { - CONTAINS_NULL, - ALL_NULL +package org.apache.iotdb.db.mpp.sql.planner.plan.node; + +public class PlanNodeAllocator { - public static int initialId = 0; - public static synchronized PlanNodeId generateId() { - initialId++; - return new PlanNodeId(String.valueOf(initialId)); - } ++ public static int initialId = 0; ++ ++ public static synchronized PlanNodeId generateId() { ++ initialId++; ++ return new PlanNodeId(String.valueOf(initialId)); ++ } } diff --cc server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/PlanNodeUtil.java index f174b74,2fddbdd..1a597b2 --- a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/PlanNodeUtil.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/PlanNodeUtil.java @@@ -16,25 -16,31 +16,24 @@@ * specific language governing permissions and limitations * under the License. */ - -package org.apache.iotdb.db.mpp.sql.planner.plan.node.sink; +package org.apache.iotdb.db.mpp.sql.planner.plan.node; -import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanNode; -import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanNodeId; - -import java.util.List; - -public class CsvSinkNode extends SinkNode { - public CsvSinkNode(PlanNodeId id) { - super(id); +public class PlanNodeUtil { - public static void printPlanNode(PlanNode root) { - printPlanNodeWithLevel(root, 0); - } ++ public static void printPlanNode(PlanNode root) { ++ printPlanNodeWithLevel(root, 0); + } - private static void printPlanNodeWithLevel(PlanNode root, int level) { - printTab(level); - System.out.println(root.toString()); - for (PlanNode child : root.getChildren()) { - printPlanNodeWithLevel(child, level + 1); - } - @Override - public List<PlanNode> getChildren() { - return null; ++ private static void printPlanNodeWithLevel(PlanNode root, int level) { ++ printTab(level); ++ System.out.println(root.toString()); ++ for (PlanNode child : root.getChildren()) { ++ printPlanNodeWithLevel(child, level + 1); + } + } - private static void printTab(int count) { - for(int i = 0 ; i < count; i ++) { - System.out.print("\t"); - } - @Override - public List<String> getOutputColumnNames() { - return null; ++ private static void printTab(int count) { ++ for (int i = 0; i < count; i++) { ++ System.out.print("\t"); + } + } - - @Override - public void close() throws Exception {} - - @Override - public void send() {} } diff --cc server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/SimplePlanNodeRewriter.java index 639e805,bb343f4..e0ca8f6 --- a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/SimplePlanNodeRewriter.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/SimplePlanNodeRewriter.java @@@ -21,25 -23,26 +21,24 @@@ package org.apache.iotdb.db.mpp.sql.pla import java.util.List; - import static com.google.common.base.Verify.verifyNotNull; -/** not implemented in current IoTDB yet */ -public class ThriftSinkNode extends SinkNode { - - public ThriftSinkNode(PlanNodeId id) { - super(id); - } +import static com.google.common.collect.ImmutableList.toImmutableList; - public class SimplePlanNodeRewriter<C> extends PlanVisitor<PlanNode, C>{ - @Override - public PlanNode visitPlan(PlanNode node, C context) { - return defaultRewrite(node, context); - } ++public class SimplePlanNodeRewriter<C> extends PlanVisitor<PlanNode, C> { + @Override - public List<PlanNode> getChildren() { - return null; ++ public PlanNode visitPlan(PlanNode node, C context) { ++ return defaultRewrite(node, context); + } - public PlanNode defaultRewrite(PlanNode node, C context) { - List<PlanNode> children = node.getChildren().stream() - .map(child -> rewrite(child, context)) - .collect(toImmutableList()); - @Override - public List<String> getOutputColumnNames() { - return null; - } ++ public PlanNode defaultRewrite(PlanNode node, C context) { ++ List<PlanNode> children = ++ node.getChildren().stream() ++ .map(child -> rewrite(child, context)) ++ .collect(toImmutableList()); - return node.cloneWithChildren(children); - } - @Override - public void close() throws Exception {} ++ return node.cloneWithChildren(children); ++ } - public PlanNode rewrite(PlanNode node, C userContext) - { - return node.accept(this, userContext); - } - @Override - public void send() {} ++ public PlanNode rewrite(PlanNode node, C userContext) { ++ return node.accept(this, userContext); ++ } } diff --cc server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/DeviceMergeNode.java index 5073d15,54a0cf8..de1c7ea --- a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/DeviceMergeNode.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/DeviceMergeNode.java @@@ -19,10 -19,10 +19,10 @@@ package org.apache.iotdb.db.mpp.sql.planner.plan.node.process; import org.apache.iotdb.db.mpp.common.FilterNullPolicy; --import org.apache.iotdb.db.mpp.common.OrderBy; import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanNode; import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanNodeId; import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanVisitor; ++import org.apache.iotdb.db.sql.statement.component.OrderBy; import java.util.List; import java.util.Map; diff --cc server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/OffsetNode.java index bde203c,2e3fc78..f009fd4 --- a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/OffsetNode.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/OffsetNode.java @@@ -7,7 -7,7 +7,7 @@@ * "License"); you may not use this file except in compliance * with the License. You may obtain a copy of the License at * -- * http://www.apache.org/licenses/LICENSE-2.0 ++ * http://www.apache.org/licenses/LICENSE-2.0 * * Unless required by applicable law or agreed to in writing, * software distributed under the License is distributed on an diff --cc server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/SortNode.java index b0752f0,19464a2..7e2b1c4 --- a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/SortNode.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/SortNode.java @@@ -18,10 -18,10 +18,10 @@@ */ package org.apache.iotdb.db.mpp.sql.planner.plan.node.process; --import org.apache.iotdb.db.mpp.common.OrderBy; import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanNode; import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanNodeId; import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanVisitor; ++import org.apache.iotdb.db.sql.statement.component.OrderBy; import com.google.common.collect.ImmutableList; diff --cc server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/TimeJoinNode.java index 907c91a,74fac14..7be2140 --- a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/TimeJoinNode.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/process/TimeJoinNode.java @@@ -19,13 -19,11 +19,13 @@@ package org.apache.iotdb.db.mpp.sql.planner.plan.node.process; import org.apache.iotdb.db.mpp.common.FilterNullPolicy; --import org.apache.iotdb.db.mpp.common.OrderBy; import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanNode; +import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanNodeAllocator; import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanNodeId; import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanVisitor; ++import org.apache.iotdb.db.sql.statement.component.OrderBy; +import java.util.ArrayList; import java.util.List; import java.util.stream.Collectors; @@@ -48,16 -46,6 +48,13 @@@ public class TimeJoinNode extends Proce private List<PlanNode> children; - public TimeJoinNode( - PlanNodeId id, - OrderBy mergeOrder, - FilterNullPolicy filterNullPolicy) { ++ public TimeJoinNode(PlanNodeId id, OrderBy mergeOrder, FilterNullPolicy filterNullPolicy) { + super(id); + this.mergeOrder = mergeOrder; + this.filterNullPolicy = filterNullPolicy; + this.children = new ArrayList<>(); + } + public TimeJoinNode( PlanNodeId id, OrderBy mergeOrder, @@@ -111,9 -89,4 +108,8 @@@ public void setWithoutPolicy(FilterNullPolicy filterNullPolicy) { this.filterNullPolicy = filterNullPolicy; } + + public String toString() { + return "TimeJoinNode-" + this.getId(); + } - } diff --cc server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/source/SeriesScanNode.java index f6ab1c6,8858074..628d4da --- a/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/source/SeriesScanNode.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/sql/planner/plan/node/source/SeriesScanNode.java @@@ -19,12 -19,10 +19,12 @@@ package org.apache.iotdb.db.mpp.sql.planner.plan.node.source; import org.apache.iotdb.db.metadata.path.PartialPath; -import org.apache.iotdb.db.mpp.common.OrderBy; +import org.apache.iotdb.db.mpp.common.DataRegion; - import org.apache.iotdb.db.mpp.common.OrderBy; import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanNode; +import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanNodeAllocator; import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanNodeId; import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanVisitor; ++import org.apache.iotdb.db.sql.statement.component.OrderBy; import org.apache.iotdb.tsfile.read.filter.basic.Filter; import com.google.common.collect.ImmutableList; @@@ -120,25 -105,4 +120,26 @@@ public class SeriesScanNode extends Sou public <R, C> R accept(PlanVisitor<R, C> visitor, C context) { return visitor.visitSeriesScan(this, context); } + + public PartialPath getSeriesPath() { + return seriesPath; + } + + public Filter getTimeFilter() { + return timeFilter; + } + + public void setDataRegion(DataRegion dataRegion) { + this.dataRegion = dataRegion; + } + + public DataRegion getDataRegion() { + return dataRegion; + } + + public String toString() { - return String.format("SeriesScanNode-%s:[SeriesPath: %s, DataRegion: %s]", - this.getId(), this.getSeriesPath(), this.getDataRegion()); ++ return String.format( ++ "SeriesScanNode-%s:[SeriesPath: %s, DataRegion: %s]", ++ this.getId(), this.getSeriesPath(), this.getDataRegion()); + } } diff --cc server/src/test/java/org/apache/iotdb/db/mpp/sql/plan/DistributionPlannerTest.java index 3a222fa,0000000..0d9c798 mode 100644,000000..100644 --- a/server/src/test/java/org/apache/iotdb/db/mpp/sql/plan/DistributionPlannerTest.java +++ b/server/src/test/java/org/apache/iotdb/db/mpp/sql/plan/DistributionPlannerTest.java @@@ -1,91 -1,0 +1,98 @@@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package org.apache.iotdb.db.mpp.sql.plan; + +import org.apache.iotdb.db.exception.metadata.IllegalPathException; +import org.apache.iotdb.db.metadata.path.PartialPath; +import org.apache.iotdb.db.mpp.common.*; +import org.apache.iotdb.db.mpp.sql.planner.plan.DistributionPlanner; +import org.apache.iotdb.db.mpp.sql.planner.plan.LogicalQueryPlan; +import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanNode; +import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanNodeAllocator; - import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanNodeId; +import org.apache.iotdb.db.mpp.sql.planner.plan.node.PlanNodeUtil; +import org.apache.iotdb.db.mpp.sql.planner.plan.node.process.LimitNode; +import org.apache.iotdb.db.mpp.sql.planner.plan.node.process.TimeJoinNode; +import org.apache.iotdb.db.mpp.sql.planner.plan.node.source.SeriesScanNode; ++import org.apache.iotdb.db.sql.statement.component.OrderBy; ++ +import org.junit.Test; + +import java.util.ArrayList; +import java.util.HashMap; +import java.util.List; +import java.util.Map; + +import static org.junit.Assert.assertEquals; + +public class DistributionPlannerTest { + - @Test - public void TestRewriteSourceNode() throws IllegalPathException { - TimeJoinNode timeJoinNode = new TimeJoinNode(PlanNodeAllocator.generateId(), OrderBy.TIMESTAMP_ASC, FilterNullPolicy.NO_FILTER); - - timeJoinNode.addChild(new SeriesScanNode(PlanNodeAllocator.generateId(), new PartialPath("root.sg.d1.s1"))); - timeJoinNode.addChild(new SeriesScanNode(PlanNodeAllocator.generateId(), new PartialPath("root.sg.d1.s2"))); - timeJoinNode.addChild(new SeriesScanNode(PlanNodeAllocator.generateId(), new PartialPath("root.sg.d2.s1"))); - - LimitNode root = new LimitNode(PlanNodeAllocator.generateId(), 10, timeJoinNode); - - Analysis analysis = constructAnalysis(); - - DistributionPlanner planner = new DistributionPlanner(analysis, new LogicalQueryPlan(new QueryContext(), root)); - PlanNode newRoot = planner.rewriteSource(); - - System.out.println("\nLogical-Plan:"); - System.out.println("------------------"); - PlanNodeUtil.printPlanNode(root); - System.out.println("\nDistributed-Plan:"); - System.out.println("------------------"); - PlanNodeUtil.printPlanNode(newRoot); - assertEquals(newRoot.getChildren().get(0).getChildren().size(), 3); - assertEquals(newRoot.getChildren().get(0).getChildren().get(0).getChildren().size(), 2); - assertEquals(newRoot.getChildren().get(0).getChildren().get(1).getChildren().size(), 2); - } - - private Analysis constructAnalysis() { - Analysis analysis = new Analysis(); - Map<String, Map<DataRegionTimeSlice, List<DataRegion>>> dataPartitionInfo = new HashMap<>(); - List<DataRegion> d1DataRegions = new ArrayList<>(); - d1DataRegions.add(new DataRegion(1, "192.0.0.1")); - d1DataRegions.add(new DataRegion(2, "192.0.0.1")); - Map<DataRegionTimeSlice, List<DataRegion>> d1DataRegionMap = new HashMap<>(); - d1DataRegionMap.put(new DataRegionTimeSlice(), d1DataRegions); - - List<DataRegion> d2DataRegions = new ArrayList<>(); - d2DataRegions.add(new DataRegion(3, "192.0.0.1")); - Map<DataRegionTimeSlice, List<DataRegion>> d2DataRegionMap = new HashMap<>(); - d2DataRegionMap.put(new DataRegionTimeSlice(), d2DataRegions); - - dataPartitionInfo.put("root.sg.d1", d1DataRegionMap); - dataPartitionInfo.put("root.sg.d2", d2DataRegionMap); - - analysis.setDataPartitionInfo(dataPartitionInfo); - return analysis; - } ++ @Test ++ public void TestRewriteSourceNode() throws IllegalPathException { ++ TimeJoinNode timeJoinNode = ++ new TimeJoinNode( ++ PlanNodeAllocator.generateId(), OrderBy.TIMESTAMP_ASC, FilterNullPolicy.NO_FILTER); ++ ++ timeJoinNode.addChild( ++ new SeriesScanNode(PlanNodeAllocator.generateId(), new PartialPath("root.sg.d1.s1"))); ++ timeJoinNode.addChild( ++ new SeriesScanNode(PlanNodeAllocator.generateId(), new PartialPath("root.sg.d1.s2"))); ++ timeJoinNode.addChild( ++ new SeriesScanNode(PlanNodeAllocator.generateId(), new PartialPath("root.sg.d2.s1"))); ++ ++ LimitNode root = new LimitNode(PlanNodeAllocator.generateId(), 10, timeJoinNode); ++ ++ Analysis analysis = constructAnalysis(); ++ ++ DistributionPlanner planner = ++ new DistributionPlanner(analysis, new LogicalQueryPlan(new MPPQueryContext(), root)); ++ PlanNode newRoot = planner.rewriteSource(); ++ ++ System.out.println("\nLogical-Plan:"); ++ System.out.println("------------------"); ++ PlanNodeUtil.printPlanNode(root); ++ System.out.println("\nDistributed-Plan:"); ++ System.out.println("------------------"); ++ PlanNodeUtil.printPlanNode(newRoot); ++ assertEquals(newRoot.getChildren().get(0).getChildren().size(), 3); ++ assertEquals(newRoot.getChildren().get(0).getChildren().get(0).getChildren().size(), 2); ++ assertEquals(newRoot.getChildren().get(0).getChildren().get(1).getChildren().size(), 2); ++ } ++ ++ private Analysis constructAnalysis() { ++ Analysis analysis = new Analysis(); ++ Map<String, Map<DataRegionTimeSlice, List<DataRegion>>> dataPartitionInfo = new HashMap<>(); ++ List<DataRegion> d1DataRegions = new ArrayList<>(); ++ d1DataRegions.add(new DataRegion(1, "192.0.0.1")); ++ d1DataRegions.add(new DataRegion(2, "192.0.0.1")); ++ Map<DataRegionTimeSlice, List<DataRegion>> d1DataRegionMap = new HashMap<>(); ++ d1DataRegionMap.put(new DataRegionTimeSlice(), d1DataRegions); ++ ++ List<DataRegion> d2DataRegions = new ArrayList<>(); ++ d2DataRegions.add(new DataRegion(3, "192.0.0.1")); ++ Map<DataRegionTimeSlice, List<DataRegion>> d2DataRegionMap = new HashMap<>(); ++ d2DataRegionMap.put(new DataRegionTimeSlice(), d2DataRegions); ++ ++ dataPartitionInfo.put("root.sg.d1", d1DataRegionMap); ++ dataPartitionInfo.put("root.sg.d2", d2DataRegionMap); ++ ++ analysis.setDataPartitionInfo(dataPartitionInfo); ++ return analysis; ++ } +}
