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;

Reply via email to