This is an automated email from the ASF dual-hosted git repository. hui pushed a commit to branch lmh/groupByTest in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit 7b15b0a49ac832263481abf8158d91fc177a85fc Merge: 8d22dc2ac1 2ae6ae9c48 Author: Minghui Liu <[email protected]> AuthorDate: Fri Apr 14 10:10:13 2023 +0800 Merge remote-tracking branch 'origin/master' into lmh/groupByTest # Conflicts: # server/src/main/java/org/apache/iotdb/db/mpp/execution/exchange/SharedTsBlockQueue.java # server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/LocalExecutionPlanner.java Jenkinsfile | 6 +- LICENSE-binary | 7 +- .../org/apache/iotdb/db/qp/sql/IdentifierParser.g4 | 13 +- .../org/apache/iotdb/db/qp/sql/IoTDBSqlParser.g4 | 579 ++++++++------ .../antlr4/org/apache/iotdb/db/qp/sql/SqlLexer.g4 | 43 + client-cpp/src/main/Session.cpp | 486 +++++++++--- client-cpp/src/main/Session.h | 121 ++- client-cpp/src/test/cpp/sessionIT.cpp | 220 +++++- client-py/SessionExample.py | 3 +- client-py/iotdb/Session.py | 580 ++++++++++---- client-py/iotdb/utils/IoTDBConstants.py | 3 + client-py/pom.xml | 3 + compile-tools/README.md | 2 +- .../confignode/client/DataNodeRequestType.java | 6 +- .../client/async/AsyncDataNodeClientPool.java | 14 + .../heartbeat/DataNodeHeartbeatHandler.java | 20 +- .../consensus/request/ConfigPhysicalPlan.java | 8 + .../consensus/request/ConfigPhysicalPlanType.java | 6 +- .../request/write/quota/SetSpaceQuotaPlan.java | 101 +++ .../request/write/quota/SetThrottleQuotaPlan.java | 113 +++ .../confignode/manager/ClusterQuotaManager.java | 281 +++++++ .../iotdb/confignode/manager/ConfigManager.java | 62 +- .../apache/iotdb/confignode/manager/IManager.java | 11 + .../iotdb/confignode/manager/node/NodeManager.java | 16 +- .../manager/partition/PartitionManager.java | 9 + .../iotdb/confignode/persistence/ModelInfo.java | 14 +- .../persistence/executor/ConfigPlanExecutor.java | 15 +- .../partition/DatabasePartitionTable.java | 20 + .../persistence/partition/PartitionInfo.java | 17 + .../confignode/persistence/quota/QuotaInfo.java | 260 ++++++ .../procedure/impl/model/CreateModelProcedure.java | 2 +- .../procedure/impl/model/DropModelProcedure.java | 31 +- .../procedure/state/model/DropModelState.java | 1 - .../procedure/store/ProcedureFactory.java | 16 +- .../thrift/ConfigNodeRPCServiceProcessor.java | 35 + .../request/ConfigPhysicalPlanSerDeTest.java | 40 + .../confignode/persistence/QuotaInfoTest.java | 103 +++ consensus/pom.xml | 2 +- .../org/apache/iotdb/consensus/common/Utils.java | 32 - .../iot/logdispatcher/IndexController.java | 2 +- .../ratis/ApplicationStateMachineProxy.java | 1 + .../iotdb/consensus/ratis/RatisConsensus.java | 35 +- .../iotdb/consensus/ratis/ResponseMessage.java | 1 + .../iotdb/consensus/ratis/SnapshotStorage.java | 11 +- .../ratis/metrics/IoTDBMetricRegistry.java | 2 +- .../consensus/ratis/utils/RatisLogMonitor.java | 87 ++ .../iotdb/consensus/ratis/{ => utils}/Utils.java | 4 +- .../iot/logdispatcher/IndexControllerTest.java | 2 +- .../apache/iotdb/consensus/ratis/SnapshotTest.java | 54 +- .../apache/iotdb/consensus/ratis/UtilsTest.java | 1 + .../DockerCompose/docker-compose-cluster-1c2d.yml | 6 +- .../DockerCompose/docker-compose-host-3c3d.yml | 4 +- .../DockerCompose/docker-compose-standalone.yml | 3 +- docker/src/main/Dockerfile-1.0.0-datanode | 3 +- docs/Community/Materials.md | 98 +-- docs/Download/README.md | 22 +- docs/UserGuide/API/InfluxDB-Protocol.md | 10 +- docs/UserGuide/API/Programming-Java-Native-API.md | 93 +-- docs/UserGuide/API/Programming-MQTT.md | 4 +- .../UserGuide/API/Programming-Python-Native-API.md | 8 +- .../API/{RestService.md => RestServiceV1.md} | 46 +- .../API/{RestService.md => RestServiceV2.md} | 50 +- docs/UserGuide/Cluster/Cluster-Concept.md | 4 +- docs/UserGuide/Cluster/Cluster-Maintenance.md | 2 +- docs/UserGuide/Data-Concept/Compression.md | 2 + .../Data-Concept/Data-Model-and-Terminology.md | 4 +- docs/UserGuide/Data-Concept/Encoding.md | 24 +- docs/UserGuide/Data-Concept/Schema-Template.md | 6 +- docs/UserGuide/Data-Concept/Time-Partition.md | 2 +- docs/UserGuide/Ecosystem-Integration/DBeaver.md | 16 +- .../Ecosystem-Integration/Grafana-Connector.md | 6 +- .../Ecosystem-Integration/Grafana-Plugin.md | 58 +- docs/UserGuide/Ecosystem-Integration/NiFi-IoTDB.md | 4 +- .../UserGuide/Ecosystem-Integration/Spark-IoTDB.md | 2 +- .../Ecosystem-Integration/Spark-TsFile.md | 6 +- .../Ecosystem-Integration/Writing-Data-on-HDFS.md | 2 +- .../Ecosystem-Integration/Zeppelin-IoTDB.md | 8 +- .../Edge-Cloud-Collaboration/Sync-Tool.md | 2 +- docs/UserGuide/IoTDB-Introduction/Architecture.md | 2 +- docs/UserGuide/IoTDB-Introduction/Publication.md | 2 +- docs/UserGuide/IoTDB-Introduction/Scenario.md | 14 +- docs/UserGuide/Maintenance-Tools/JMX-Tool.md | 4 +- docs/UserGuide/Maintenance-Tools/Log-Tool.md | 6 +- docs/UserGuide/Monitor-Alert/Alerting.md | 2 +- docs/UserGuide/Monitor-Alert/Metric-Tool.md | 10 +- .../Operate-Metadata/Auto-Create-MetaData.md | 2 +- docs/UserGuide/Operate-Metadata/Template.md | 6 +- docs/UserGuide/Operate-Metadata/Timeseries.md | 2 +- docs/UserGuide/Operators-Functions/Aggregation.md | 30 +- docs/UserGuide/Operators-Functions/Conditional.md | 351 +++++++++ docs/UserGuide/Operators-Functions/Conversion.md | 2 +- docs/UserGuide/Operators-Functions/Sample.md | 6 +- .../Operators-Functions/User-Defined-Function.md | 10 +- docs/UserGuide/Query-Data/Continuous-Query.md | 8 +- docs/UserGuide/Query-Data/Group-By.md | 6 +- docs/UserGuide/Query-Data/Overview.md | 2 +- .../UserGuide/QuickStart/Command-Line-Interface.md | 24 +- docs/UserGuide/QuickStart/WayToGetIoTDB.md | 13 +- docs/UserGuide/Reference/Common-Config-Manual.md | 31 +- docs/UserGuide/Reference/Keywords.md | 1 + docs/UserGuide/Reference/TSDB-Comparison.md | 16 +- docs/UserGuide/Write-Data/REST-API.md | 2 +- docs/zh/Download/README.md | 29 +- docs/zh/UserGuide/API/InfluxDB-Protocol.md | 10 +- .../UserGuide/API/Programming-Java-Native-API.md | 83 +- docs/zh/UserGuide/API/Programming-MQTT.md | 4 +- .../UserGuide/API/Programming-Python-Native-API.md | 8 +- .../API/{RestService.md => RestServiceV1.md} | 46 +- .../API/{RestService.md => RestServiceV2.md} | 50 +- docs/zh/UserGuide/Cluster/Cluster-Concept.md | 4 +- docs/zh/UserGuide/Cluster/IoTDB-Deploy.md | 361 --------- docs/zh/UserGuide/Data-Concept/Compression.md | 1 + .../Data-Concept/Data-Model-and-Terminology.md | 4 +- docs/zh/UserGuide/Data-Concept/Encoding.md | 29 +- docs/zh/UserGuide/Data-Concept/Schema-Template.md | 6 +- docs/zh/UserGuide/Data-Concept/Time-Partition.md | 2 +- docs/zh/UserGuide/Ecosystem-Integration/DBeaver.md | 16 +- .../Ecosystem-Integration/Grafana-Connector.md | 6 +- .../Ecosystem-Integration/Grafana-Plugin.md | 58 +- .../UserGuide/Ecosystem-Integration/NiFi-IoTDB.md | 4 +- .../Ecosystem-Integration/Spark-TsFile.md | 24 +- .../UserGuide/Ecosystem-Integration/Workbench.md | 82 +- .../Ecosystem-Integration/Writing-Data-on-HDFS.md | 2 +- .../Ecosystem-Integration/Zeppelin-IoTDB.md | 8 +- .../Edge-Cloud-Collaboration/Sync-Tool.md | 2 +- .../UserGuide/IoTDB-Introduction/Architecture.md | 2 +- .../zh/UserGuide/IoTDB-Introduction/Publication.md | 2 +- docs/zh/UserGuide/IoTDB-Introduction/Scenario.md | 14 +- docs/zh/UserGuide/Maintenance-Tools/JMX-Tool.md | 4 +- docs/zh/UserGuide/Maintenance-Tools/Log-Tool.md | 6 +- docs/zh/UserGuide/Monitor-Alert/Alerting.md | 2 +- docs/zh/UserGuide/Monitor-Alert/Metric-Tool.md | 6 +- docs/zh/UserGuide/Operate-Metadata/Template.md | 6 +- docs/zh/UserGuide/Operate-Metadata/Timeseries.md | 2 +- .../UserGuide/Operators-Functions/Aggregation.md | 30 +- .../UserGuide/Operators-Functions/Conditional.md | 347 ++++++++ .../zh/UserGuide/Operators-Functions/Conversion.md | 2 +- docs/zh/UserGuide/Operators-Functions/Overview.md | 10 +- docs/zh/UserGuide/Operators-Functions/Sample.md | 6 +- .../Operators-Functions/User-Defined-Function.md | 2 +- docs/zh/UserGuide/Query-Data/Continuous-Query.md | 8 +- docs/zh/UserGuide/Query-Data/Group-By.md | 6 +- docs/zh/UserGuide/Query-Data/Overview.md | 2 +- .../UserGuide/QuickStart/Command-Line-Interface.md | 24 +- docs/zh/UserGuide/QuickStart/WayToGetIoTDB.md | 11 +- .../zh/UserGuide/Reference/Common-Config-Manual.md | 30 +- docs/zh/UserGuide/Reference/Keywords.md | 1 + docs/zh/UserGuide/Reference/TSDB-Comparison.md | 14 +- docs/zh/UserGuide/Trigger/Implement-Trigger.md | 4 +- docs/zh/UserGuide/Write-Data/REST-API.md | 2 +- .../src/AlignedTimeseriesSessionExample.cpp | 8 +- example/client-cpp-example/src/SessionExample.cpp | 9 +- grafana-plugin/pkg/plugin/plugin.go | 8 +- .../java/org/apache/iotdb/db/it/IoTDBFilterIT.java | 5 + .../db/it/IoTDBSyntaxConventionIdentifierIT.java | 20 +- .../iotdb/db/it/aggregation/IoTDBModeIT.java | 24 +- .../db/it/aggregation/IoTDBTagAggregationIT.java | 10 + .../db/it/alignbydevice/IoTDBAlignByDeviceIT.java | 108 +++ .../scalar/IoTDBSubStringFunctionIT.java | 82 +- .../iotdb/db/it/query/IoTDBCaseWhenThenIT.java | 876 +++++++++++++++++++++ .../iotdb/db/it/query/IoTDBNullOperandIT.java | 3 + .../iotdb/db/it/schema/IoTDBSchemaTemplateIT.java | 13 + .../db/it/specialwords/IoTDBSpecialWordsIT.java | 77 ++ .../session/it/IoTDBSessionSchemaTemplateIT.java | 6 +- .../java/org/apache/iotdb/isession/ISession.java | 2 +- .../apache/iotdb/isession/pool/ISessionPool.java | 2 +- .../apache/iotdb/jdbc/IoTDBDatabaseMetadata.java | 1 + .../iotdb/library/dprofile/util/GKArray.java | 17 +- .../iotdb/metrics/AbstractMetricService.java | 10 +- mlnode/.gitignore | 6 +- mlnode/iotdb/mlnode/client.py | 107 ++- mlnode/iotdb/mlnode/config.py | 14 +- mlnode/iotdb/mlnode/constant.py | 10 + mlnode/iotdb/mlnode/handler.py | 29 +- mlnode/iotdb/mlnode/service.py | 8 +- .../iotdb/mlnode/{model_storage.py => storage.py} | 23 +- mlnode/iotdb/mlnode/util.py | 18 +- mlnode/pyproject.toml | 1 + mlnode/requirements.txt | 2 +- mlnode/requirements_dev.txt | 4 +- mlnode/test/test_model_storage.py | 37 +- .../resources/conf/iotdb-common.properties | 28 +- .../iotdb/commons/concurrent/ThreadName.java | 8 +- .../apache/iotdb/commons/conf/IoTDBConstant.java | 12 + .../commons/exception/RpcThrottlingException.java | 13 +- .../iotdb/commons/model/ModelHyperparameter.java | 10 + .../iotdb/commons/model/ModelInformation.java | 89 ++- .../iotdb/commons/model/TrailInformation.java | 7 +- .../iotdb/commons/quotas/SpaceQuotaType.java | 8 +- .../apache/iotdb/commons/service/ServiceType.java | 3 +- .../commons/utils/BasicStructureSerDeUtil.java | 16 + openapi/pom.xml | 56 +- openapi/src/main/openapi3/iotdb_rest_common.yaml | 63 ++ .../{iotdb-rest.yaml => iotdb_rest_v1.yaml} | 35 +- .../{iotdb-rest.yaml => iotdb_rest_v2.yaml} | 35 +- pom.xml | 11 +- .../schemaregion/rocksdb/RSchemaRegion.java | 10 + .../metadata/tagSchemaRegion/TagSchemaRegion.java | 10 + .../src/main/codegen/templates/ModeAccumulator.ftl | 49 +- .../apache/iotdb/db/client/ConfigNodeClient.java | 159 +++- .../org/apache/iotdb/db/client/MLNodeClient.java | 18 +- .../java/org/apache/iotdb/db/conf/IoTDBConfig.java | 86 +- .../org/apache/iotdb/db/conf/IoTDBDescriptor.java | 45 ++ .../org/apache/iotdb/db/engine/StorageEngine.java | 10 + .../impl/RewriteCrossSpaceCompactionSelector.java | 30 + .../utils/CrossSpaceCompactionCandidate.java | 15 +- .../iotdb/db/engine/flush/MemTableFlushTask.java | 2 +- .../iotdb/db/engine/storagegroup/DataRegion.java | 53 +- .../quota/ExceedQuotaException.java} | 13 +- .../runtime/MemoryLeakException.java} | 11 +- .../db/metadata/cache/DataNodeSchemaCache.java | 6 +- .../db/metadata/mtree/MTreeBelowSGMemoryImpl.java | 24 +- .../db/metadata/schemaregion/ISchemaRegion.java | 5 + .../db/metadata/schemaregion/SchemaEngine.java | 30 + .../schemaregion/SchemaRegionMemoryImpl.java | 45 ++ .../schemaregion/SchemaRegionSchemaFileImpl.java | 45 ++ .../metadata/template/ClusterTemplateManager.java | 17 + .../iotdb/db/mpp/common/FragmentInstanceId.java | 4 + .../apache/iotdb/db/mpp/common/SessionInfo.java | 14 + .../db/mpp/common/header/ColumnHeaderConstant.java | 48 ++ .../db/mpp/common/header/DatasetHeaderFactory.java | 16 + .../exception/CpuNotEnoughException.java} | 12 +- .../iotdb/db/mpp/execution/driver/DataDriver.java | 11 +- .../db/mpp/execution/driver/DriverContext.java | 28 +- .../iotdb/db/mpp/execution/driver/IDriver.java | 6 +- .../execution/exchange/MPPDataExchangeManager.java | 169 +++- .../mpp/execution/exchange/SharedTsBlockQueue.java | 63 +- .../execution/exchange/sink/LocalSinkChannel.java | 2 - .../execution/exchange/sink/ShuffleSinkHandle.java | 90 ++- .../mpp/execution/exchange/sink/SinkChannel.java | 33 +- .../exchange/source/LocalSourceHandle.java | 2 +- .../exchange/source/PipelineSourceHandle.java} | 28 +- .../execution/exchange/source/SourceHandle.java | 12 +- .../fragment/FragmentInstanceContext.java | 19 +- .../fragment/FragmentInstanceExecution.java | 21 +- .../iotdb/db/mpp/execution/memory/MemoryPool.java | 251 +++--- .../operator/process/AbstractIntoOperator.java | 64 +- .../operator/process/DeviceViewIntoOperator.java | 7 +- .../operator/process/FilterAndProjectOperator.java | 22 + .../execution/operator/process/IntoOperator.java | 7 +- .../execution/operator/process/OffsetOperator.java | 4 +- .../operator/source/ExchangeOperator.java | 13 + .../db/mpp/execution/schedule/DriverScheduler.java | 95 ++- .../mpp/execution/schedule/IDriverScheduler.java | 7 +- .../db/mpp/execution/schedule/task/DriverTask.java | 32 +- .../iotdb/db/mpp/plan/analyze/AnalyzeVisitor.java | 31 +- .../apache/iotdb/db/mpp/plan/analyze/Analyzer.java | 16 +- .../db/mpp/plan/analyze/ConcatPathRewriter.java | 8 - .../db/mpp/plan/analyze/ExpressionAnalyzer.java | 101 ++- .../mpp/plan/analyze/ExpressionTypeAnalyzer.java | 49 ++ .../iotdb/db/mpp/plan/analyze/ExpressionUtils.java | 26 + .../plan/analyze/schema/ClusterSchemaFetcher.java | 7 +- .../db/mpp/plan/analyze/schema/ISchemaFetcher.java | 3 +- .../db/mpp/plan/execution/QueryExecution.java | 7 +- .../plan/execution/config/ConfigTaskVisitor.java | 62 ++ .../config/executor/ClusterConfigTaskExecutor.java | 244 ++++++ .../config/executor/IConfigTaskExecutor.java | 29 + .../config/metadata/model/CreateModelTask.java | 42 + .../config/metadata/model/DropModelTask.java} | 28 +- .../config/metadata/model/ShowModelsTask.java | 96 +++ .../config/metadata/model/ShowTrailsTask.java | 90 +++ .../config/sys/quota/SetSpaceQuotaTask.java | 42 + .../config/sys/quota/SetThrottleQuotaTask.java | 42 + .../config/sys/quota/ShowSpaceQuotaTask.java | 130 +++ .../config/sys/quota/ShowThrottleQuotaTask.java | 189 +++++ .../iotdb/db/mpp/plan/expression/Expression.java | 10 + .../db/mpp/plan/expression/ExpressionFactory.java | 15 + .../db/mpp/plan/expression/ExpressionType.java | 4 + .../plan/expression/binary/BinaryExpression.java | 3 +- .../plan/expression/binary/WhenThenExpression.java | 73 ++ .../builtin/helper/SubStringFunctionHelper.java | 35 +- .../expression/other/CaseWhenThenExpression.java | 172 ++++ .../visitor/CartesianProductVisitor.java | 27 + .../plan/expression/visitor/CollectVisitor.java | 7 + .../visitor/ColumnTransformerVisitor.java | 44 ++ .../ConcatExpressionWithSuffixPathsVisitor.java | 3 +- .../visitor/ExpressionAnalyzeVisitor.java | 2 +- .../plan/expression/visitor/ExpressionVisitor.java | 10 + .../visitor/IntermediateLayerVisitor.java | 7 + .../expression/visitor/ReconstructVisitor.java | 9 + .../iotdb/db/mpp/plan/parser/ASTVisitor.java | 403 +++++++++- .../db/mpp/plan/parser/StatementGenerator.java | 86 ++ .../plan/planner/LocalExecutionPlanContext.java | 10 +- .../db/mpp/plan/planner/LocalExecutionPlanner.java | 15 +- .../db/mpp/plan/planner/OperatorTreeGenerator.java | 19 +- .../db/mpp/plan/planner/PipelineDriverFactory.java | 22 +- .../db/mpp/plan/planner/plan/node/PlanVisitor.java | 190 +++-- .../scheduler/FragmentInstanceDispatcherImpl.java | 30 +- .../iotdb/db/mpp/plan/statement/StatementType.java | 5 + .../db/mpp/plan/statement/StatementVisitor.java | 42 + .../db/mpp/plan/statement/crud/QueryStatement.java | 2 +- .../metadata/model/CreateModelStatement.java | 107 +++ .../metadata/model/DropModelStatement.java} | 40 +- .../metadata/model/ShowModelsStatement.java} | 32 +- .../metadata/model/ShowTrailsStatement.java | 57 ++ .../sys/quota/SetSpaceQuotaStatement.java | 100 +++ .../sys/quota/SetThrottleQuotaStatement.java | 94 +++ .../sys/quota/ShowSpaceQuotaStatement.java | 62 ++ .../sys/quota/ShowThrottleQuotaStatement.java | 63 ++ .../dag/column/CaseWhenThenColumnTransformer.java | 132 ++++ .../binary/CompareNonEqualColumnTransformer.java | 2 +- .../binary/LogicBinaryColumnTransformer.java | 4 +- .../db/pipe/agent/runtime/PipeRuntimeAgent.java | 15 + .../PipeConnectorPluginRuntimeWrapper.java | 44 +- .../PipeProcessorPluginRuntimeWrapper.java | 48 +- .../executor/PipeAssignerSubtaskExecutor.java | 12 +- .../executor/PipeConnectorSubtaskExecutor.java | 12 +- .../executor/PipeProcessorSubtaskExecutor.java | 12 +- .../execution/executor/PipeSubtaskExecutor.java | 122 ++- ...kExecutor.java => PipeTaskExecutorManager.java} | 40 +- .../scheduler/PipeProcessorSubtaskScheduler.java | 36 - .../execution/scheduler/PipeSubtaskScheduler.java | 33 - .../execution/scheduler/PipeTaskScheduler.java | 44 +- .../org/apache/iotdb/db/pipe/task/PipeTask.java | 31 +- .../DecoratingLock.java} | 26 +- .../PipeAssignerSubtask.java | 6 +- .../PipeConnectorSubtask.java | 13 +- .../PipeProcessorSubtask.java | 13 +- .../iotdb/db/pipe/task/callable/PipeSubtask.java | 135 ++++ .../db/pipe/task/stage/PipeTaskCollectorStage.java | 20 +- .../db/pipe/task/stage/PipeTaskConnectorStage.java | 20 +- .../db/pipe/task/stage/PipeTaskProcessorStage.java | 20 +- .../iotdb/db/pipe/task/stage/PipeTaskStage.java | 37 +- .../rest/handler/AuthorizationHandler.java | 8 +- .../rest/{ => v1}/handler/ExceptionHandler.java | 4 +- .../{ => v1}/handler/ExecuteStatementHandler.java | 2 +- .../rest/{ => v1}/handler/QueryDataSetHandler.java | 24 +- .../{ => v1}/handler/RequestValidationHandler.java | 22 +- .../handler/StatementConstructionHandler.java | 6 +- .../rest/{ => v1}/impl/GrafanaApiServiceImpl.java | 25 +- .../rest/{ => v1}/impl/RestApiServiceImpl.java | 20 +- .../rest/{ => v2}/handler/ExceptionHandler.java | 2 +- .../{ => v2}/handler/ExecuteStatementHandler.java | 2 +- .../rest/{ => v2}/handler/QueryDataSetHandler.java | 26 +- .../{ => v2}/handler/RequestValidationHandler.java | 8 +- .../handler/StatementConstructionHandler.java | 4 +- .../rest/{ => v2}/impl/GrafanaApiServiceImpl.java | 25 +- .../rest/{ => v2}/impl/RestApiServiceImpl.java | 20 +- .../iotdb/db/query/control/SessionManager.java | 6 +- .../db/quotas/AverageIntervalRateLimiter.java | 75 ++ .../apache/iotdb/db/quotas/DataNodeSizeStore.java | 60 ++ .../iotdb/db/quotas/DataNodeSpaceQuotaManager.java | 153 ++++ .../db/quotas/DataNodeThrottleQuotaManager.java | 153 ++++ .../iotdb/db/quotas/DefaultOperationQuota.java | 189 +++++ .../iotdb/db/quotas/FixedIntervalRateLimiter.java | 57 ++ .../NoopOperationQuota.java} | 35 +- .../org/apache/iotdb/db/quotas/OperationQuota.java | 50 ++ .../org/apache/iotdb/db/quotas/QuotaLimiter.java | 198 +++++ .../org/apache/iotdb/db/quotas/RateLimiter.java | 130 +++ .../apache/iotdb/db/quotas/ThrottleQuotaLimit.java | 76 ++ .../java/org/apache/iotdb/db/service/DataNode.java | 4 + .../apache/iotdb/db/service/MLNodeRPCService.java | 98 +++ .../MLNodeRPCServiceMBean.java} | 4 +- .../metrics/IoTDBInternalLocalReporter.java | 66 +- .../handler/MLNodeRPCServiceThriftHandler.java | 56 ++ .../service/thrift/impl/ClientRPCServiceImpl.java | 119 ++- .../impl/DataNodeInternalRPCServiceImpl.java | 67 +- .../thrift/impl/IMLNodeRPCServiceWithHandler.java | 13 +- .../service/thrift/impl/MLNodeRPCServiceImpl.java | 205 +++++ .../org/apache/iotdb/db/utils/SchemaUtils.java | 6 + .../engine/compaction/CompactionSchedulerTest.java | 3 + ...eCompactionWithFastPerformerValidationTest.java | 705 +++++++++++++++++ ...actionWithReadPointPerformerValidationTest.java | 713 ++++++++++++++++- .../db/metadata/cache/DataNodeSchemaCacheTest.java | 9 +- .../iotdb/db/mpp/execution/DataDriverTest.java | 2 +- .../db/mpp/execution/memory/MemoryPoolTest.java | 27 +- .../mpp/execution/operator/OffsetOperatorTest.java | 87 ++ .../mpp/execution/operator/OperatorMemoryTest.java | 77 ++ .../schedule/DefaultDriverSchedulerTest.java | 28 +- .../execution/schedule/DriverSchedulerTest.java | 31 +- .../DriverTaskTimeoutSentinelThreadTest.java | 18 +- .../other/CaseWhenThenExpressionTest.java | 73 ++ .../iotdb/db/mpp/plan/analyze/AnalyzeTest.java | 36 + .../mpp/plan/analyze/ExpressionAnalyzerTest.java | 3 +- .../db/mpp/plan/analyze/FakeSchemaFetcherImpl.java | 7 +- .../iotdb/db/mpp/plan/plan/distribution/Util.java | 2 +- .../executor/PipeAssignerSubtaskExecutorTest.java} | 20 +- .../PipeConnectorSubtaskExecutorTest.java} | 24 +- .../PipeProcessorSubtaskExecutorTest.java} | 24 +- .../executor/PipeSubtaskExecutorTest.java | 158 ++++ .../java/org/apache/iotdb/rpc/TSStatusCode.java | 14 +- .../java/org/apache/iotdb/session/Session.java | 8 +- .../apache/iotdb/session/SessionConnection.java | 9 +- .../org/apache/iotdb/session/pool/SessionPool.java | 4 +- site/iotdb-doap.rdf | 8 + site/src/main/.vuepress/sidebar/V1.0.x/zh.ts | 1 - site/src/main/.vuepress/sidebar/V1.1.x/en.ts | 3 +- site/src/main/.vuepress/sidebar/V1.1.x/zh.ts | 4 +- site/src/main/.vuepress/sidebar/en.ts | 4 +- site/src/main/.vuepress/sidebar/zh.ts | 5 +- spark-iotdb-connector/pom.xml | 2 +- thrift-commons/src/main/thrift/common.thrift | 43 +- .../src/main/thrift/confignode.thrift | 39 + thrift-mlnode/src/main/thrift/mlnode.thrift | 2 +- thrift/src/main/thrift/client.thrift | 4 +- thrift/src/main/thrift/datanode.thrift | 85 +- tsfile/pom.xml | 5 + .../iotdb/tsfile/common/conf/TSFileConfig.java | 4 + .../apache/iotdb/tsfile/compress/ICompressor.java | 85 ++ .../iotdb/tsfile/compress/IUnCompressor.java | 49 ++ .../iotdb/tsfile/encoding/decoder/Decoder.java | 26 + .../tsfile/encoding/decoder/DoubleRLBEDecoder.java | 197 +++++ .../encoding/decoder/DoubleSprintzDecoder.java | 139 ++++ .../tsfile/encoding/decoder/FloatRLBEDecoder.java | 197 +++++ .../encoding/decoder/FloatSprintzDecoder.java | 141 ++++ .../tsfile/encoding/decoder/IntRLBEDecoder.java | 196 +++++ .../tsfile/encoding/decoder/IntSprintzDecoder.java | 129 +++ .../tsfile/encoding/decoder/LongRLBEDecoder.java | 196 +++++ .../encoding/decoder/LongSprintzDecoder.java | 127 +++ .../tsfile/encoding/decoder/SprintzDecoder.java | 54 ++ .../iotdb/tsfile/encoding/encoder/DoubleRLBE.java | 272 +++++++ .../encoding/encoder/DoubleSprintzEncoder.java | 157 ++++ .../iotdb/tsfile/encoding/encoder/FloatRLBE.java | 273 +++++++ .../encoding/encoder/FloatSprintzEncoder.java | 156 ++++ .../iotdb/tsfile/encoding/encoder/IntRLBE.java | 263 +++++++ .../tsfile/encoding/encoder/IntSprintzEncoder.java | 153 ++++ .../iotdb/tsfile/encoding/encoder/LongRLBE.java | 257 ++++++ .../encoding/encoder/LongSprintzEncoder.java | 154 ++++ .../apache/iotdb/tsfile/encoding/encoder/RLBE.java | 61 ++ .../tsfile/encoding/encoder/SprintzEncoder.java | 70 ++ .../tsfile/encoding/encoder/TSEncodingBuilder.java | 50 ++ .../apache/iotdb/tsfile/encoding/fire/Fire.java | 56 ++ .../apache/iotdb/tsfile/encoding/fire/IntFire.java | 34 +- .../iotdb/tsfile/encoding/fire/LongFire.java | 32 +- .../file/metadata/enums/CompressionType.java | 6 +- .../tsfile/file/metadata/enums/TSEncoding.java | 9 +- .../common/block/column/ColumnEncoderFactory.java | 5 +- .../apache/iotdb/tsfile/compress/LZMA2Test.java | 104 +++ .../tsfile/encoding/decoder/RLBEDecoderTest.java | 257 ++++++ .../encoding/decoder/SprintzDecoderTest.java | 593 ++++++++++++++ 430 files changed, 19254 insertions(+), 3099 deletions(-) diff --cc server/src/main/java/org/apache/iotdb/db/mpp/execution/exchange/SharedTsBlockQueue.java index 24ef8d8a08,0878e0e4d6..9965262ef1 --- a/server/src/main/java/org/apache/iotdb/db/mpp/execution/exchange/SharedTsBlockQueue.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/execution/exchange/SharedTsBlockQueue.java @@@ -42,10 -41,6 +42,11 @@@ import java.util.Queue import static com.google.common.util.concurrent.Futures.immediateVoidFuture; import static com.google.common.util.concurrent.MoreExecutors.directExecutor; ++import static org.apache.iotdb.db.mpp.statistics.QueryStatistics.RESERVE_MEMORY; +import static org.apache.iotdb.db.mpp.statistics.QueryStatistics.FREE_MEM; +import static org.apache.iotdb.db.mpp.statistics.QueryStatistics.NOTIFY_END; +import static org.apache.iotdb.db.mpp.statistics.QueryStatistics.NOTIFY_NEW_TSBLOCK; +import static org.apache.iotdb.db.mpp.statistics.QueryStatistics.RESERVE_MEMORY; /** This is not thread safe class, the caller should ensure multi-threads safety. */ @NotThreadSafe @@@ -181,12 -173,6 +183,7 @@@ public class SharedTsBlockQueue throw new IllegalStateException("queue has been destroyed"); } TsBlock tsBlock = queue.remove(); - // Every time LocalSourceHandle consumes a TsBlock, it needs to send the event to - // corresponding LocalSinkChannel. - if (sinkChannel != null) { - sinkChannel.checkAndInvokeOnFinished(); - } + long startTime = System.nanoTime(); localMemoryManager .getQueryPool() .free( @@@ -194,8 -180,12 +191,13 @@@ fullFragmentInstanceId, localPlanNodeId, tsBlock.getRetainedSizeInBytes()); + QUERY_STATISTICS.addCost(FREE_MEM, System.nanoTime() - startTime); bufferRetainedSizeInBytes -= tsBlock.getRetainedSizeInBytes(); + // Every time LocalSourceHandle consumes a TsBlock, it needs to send the event to + // corresponding LocalSinkChannel. + if (sinkChannel != null) { + sinkChannel.checkAndInvokeOnFinished(); + } if (blocked.isDone() && queue.isEmpty() && !noMoreTsBlocks) { blocked = SettableFuture.create(); } @@@ -213,25 -203,26 +215,33 @@@ } Validate.notNull(tsBlock, "TsBlock cannot be null"); - Validate.isTrue(blockedOnMemory == null || blockedOnMemory.isDone(), "queue is full"); + Validate.isTrue( + blockedOnMemory == null || blockedOnMemory.isDone(), "SharedTsBlockQueue is full"); + if (!alreadyRegistered) { + localMemoryManager + .getQueryPool() + .registerPlanNodeIdToQueryMemoryMap( + localFragmentInstanceId.queryId, fullFragmentInstanceId, localPlanNodeId); + alreadyRegistered = true; + } - Pair<ListenableFuture<Void>, Boolean> pair = + + long startTime = System.nanoTime(); + Pair<ListenableFuture<Void>, Boolean> pair; + try { + pair = - localMemoryManager - .getQueryPool() - .reserve( - localFragmentInstanceId.getQueryId(), - fullFragmentInstanceId, - localPlanNodeId, - tsBlock.getRetainedSizeInBytes(), - maxBytesCanReserve); - blockedOnMemory = pair.left; - bufferRetainedSizeInBytes += tsBlock.getRetainedSizeInBytes(); + localMemoryManager + .getQueryPool() + .reserve( + localFragmentInstanceId.getQueryId(), + fullFragmentInstanceId, + localPlanNodeId, + tsBlock.getRetainedSizeInBytes(), + maxBytesCanReserve); + blockedOnMemory = pair.left; + bufferRetainedSizeInBytes += tsBlock.getRetainedSizeInBytes(); + } finally { + QUERY_STATISTICS.addCost(RESERVE_MEMORY, System.nanoTime() - startTime); + } // reserve memory failed, we should wait until there is enough memory if (!pair.right) { diff --cc server/src/main/java/org/apache/iotdb/db/mpp/plan/analyze/AnalyzeVisitor.java index 1494da179a,128b4e9dd1..dcb158c682 --- a/server/src/main/java/org/apache/iotdb/db/mpp/plan/analyze/AnalyzeVisitor.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/plan/analyze/AnalyzeVisitor.java @@@ -244,12 -243,10 +244,12 @@@ public class AnalyzeVisitor extends Sta if (queryStatement.isGroupByTag()) { schemaTree = schemaFetcher.fetchSchemaWithTags(patternTree); } else { - schemaTree = schemaFetcher.fetchSchema(patternTree); + schemaTree = schemaFetcher.fetchSchema(patternTree, context); } - QueryMetricsManager.getInstance() - .recordPlanCost(SCHEMA_FETCHER, System.nanoTime() - startTime); + long endTime = System.nanoTime() - startTime; + QueryMetricsManager.getInstance().recordPlanCost(SCHEMA_FETCHER, endTime); + QueryStatistics.getInstance().addCost(QueryStatistics.SCHEMA_FETCHER, endTime); + logger.debug("[EndFetchSchema]"); // If there is no leaf node in the schema tree, the query should be completed immediately diff --cc server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/LocalExecutionPlanner.java index 565dc8b159,30a982bc74..b1189dcc36 --- a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/LocalExecutionPlanner.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/LocalExecutionPlanner.java @@@ -68,19 -62,13 +68,18 @@@ public class LocalExecutionPlanner // Generate pipelines, return the last pipeline data structure // TODO Replace operator with operatorFactory to build multiple driver for one pipeline + long startTime = System.nanoTime(); Operator root = plan.accept(new OperatorTreeGenerator(), context); + long endTime = System.nanoTime(); + QUERY_STATISTICS.addCost(NODE_TO_OPERATOR, endTime - startTime); + startTime = endTime; // check whether current free memory is enough to execute current query - checkMemory(root, instanceContext.getStateMachine()); + long estimatedMemorySize = checkMemory(root, instanceContext.getStateMachine()); - + context.addPipelineDriverFactory(root, context.getDriverContext(), estimatedMemorySize); + endTime = System.nanoTime(); + QUERY_STATISTICS.addCost(CHECK_MEMORY, endTime - startTime); - context.addPipelineDriverFactory(root, context.getDriverContext()); - instanceContext.setSourcePaths(collectSourcePaths(context)); // set maxBytes one SourceHandle can reserve after visiting the whole tree diff --cc server/src/main/java/org/apache/iotdb/db/service/thrift/impl/ClientRPCServiceImpl.java index 7458e9d6d5,6c1f310195..551f790cf9 --- a/server/src/main/java/org/apache/iotdb/db/service/thrift/impl/ClientRPCServiceImpl.java +++ b/server/src/main/java/org/apache/iotdb/db/service/thrift/impl/ClientRPCServiceImpl.java @@@ -62,9 -62,10 +62,11 @@@ import org.apache.iotdb.db.mpp.plan.sta import org.apache.iotdb.db.mpp.plan.statement.metadata.template.DropSchemaTemplateStatement; import org.apache.iotdb.db.mpp.plan.statement.metadata.template.SetSchemaTemplateStatement; import org.apache.iotdb.db.mpp.plan.statement.metadata.template.UnsetSchemaTemplateStatement; +import org.apache.iotdb.db.mpp.statistics.QueryStatistics; import org.apache.iotdb.db.query.control.SessionManager; import org.apache.iotdb.db.query.control.clientsession.IClientSession; + import org.apache.iotdb.db.quotas.DataNodeThrottleQuotaManager; + import org.apache.iotdb.db.quotas.OperationQuota; import org.apache.iotdb.db.service.basic.BasicOpenSessionResp; import org.apache.iotdb.db.sync.SyncService; import org.apache.iotdb.db.utils.QueryDataSetUtils;
