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

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

commit f7d02745ff6a966f929ade3394b756c1a2468d76
Merge: c4536e4ce9 561a4aaf64
Author: JackieTien97 <[email protected]>
AuthorDate: Sat Oct 29 20:46:49 2022 +0800

    resolve conflicts

 .../org/apache/iotdb/db/qp/sql/IoTDBSqlParser.g4   |   46 +-
 .../antlr4/org/apache/iotdb/db/qp/sql/SqlLexer.g4  |   24 +-
 .../java/org/apache/iotdb/cli/AbstractCli.java     |    9 +-
 confignode/src/assembly/confignode.xml             |    4 +
 .../resources/conf/iotdb-confignode.properties     |  301 +----
 .../confignode/client/DataNodeRequestType.java     |    8 +-
 .../client/async/AsyncDataNodeClientPool.java      |   26 +
 .../client/async/handlers/AsyncClientHandler.java  |   24 +-
 .../heartbeat/DataNodeHeartbeatHandler.java        |   24 +-
 .../rpc/CountPathsUsingTemplateRPCHandler.java     |   87 ++
 .../async/handlers/rpc/OperatePipeRPCHandler.java  |   60 +
 .../iotdb/confignode/conf/ConfigNodeConfig.java    |   21 +
 .../confignode/conf/ConfigNodeDescriptor.java      |  349 +++---
 .../consensus/request/ConfigPhysicalPlan.java      |   40 +
 .../consensus/request/ConfigPhysicalPlanType.java  |   12 +-
 .../consensus/request/read/GetUDFJarPlan.java      |   68 ++
 .../consensus/request/write/cq/ActiveCQPlan.java   |   85 ++
 .../consensus/request/write/cq/AddCQPlan.java      |   98 ++
 .../consensus/request/write/cq/DropCQPlan.java     |   93 ++
 .../consensus/request/write/cq/ShowCQPlan.java     |   30 +-
 .../request/write/cq/UpdateCQLastExecTimePlan.java |   95 ++
 .../write/statistics/UpdateLoadStatisticsPlan.java |   81 +-
 .../write/template/DropSchemaTemplatePlan.java     |   49 +-
 .../write/template/PreUnsetSchemaTemplatePlan.java |   67 ++
 .../RollbackPreUnsetSchemaTemplatePlan.java        |   67 ++
 .../write/template/UnsetSchemaTemplatePlan.java    |   67 ++
 .../consensus/response/DataNodeRegisterResp.java   |   14 +
 .../response/{TriggerJarResp.java => JarResp.java} |   10 +-
 .../{TriggerJarResp.java => ShowCQResp.java}       |   39 +-
 .../statemachine/PartitionRegionStateMachine.java  |   25 +-
 .../confignode/manager/ClusterSchemaManager.java   |  118 +-
 .../iotdb/confignode/manager/ConfigManager.java    |   95 +-
 .../apache/iotdb/confignode/manager/IManager.java  |   32 +-
 .../iotdb/confignode/manager/ProcedureManager.java |   55 +
 .../iotdb/confignode/manager/SyncManager.java      |   87 +-
 .../iotdb/confignode/manager/TriggerManager.java   |   37 +-
 .../iotdb/confignode/manager/UDFManager.java       |   37 +-
 .../iotdb/confignode/manager/cq/CQManager.java     |  186 ++++
 .../confignode/manager/cq/CQScheduleTask.java      |  276 +++++
 .../iotdb/confignode/manager/load/LoadManager.java |   55 +-
 .../manager/load/LoadManagerMetrics.java           |    8 +-
 .../manager/load/balancer/RouteBalancer.java       |  260 ++++-
 .../manager/load/balancer/router/IRouter.java      |    4 +-
 .../load/balancer/router/LazyGreedyRouter.java     |  159 ---
 .../manager/load/balancer/router/LeaderRouter.java |   25 +-
 .../balancer/router/LoadScoreGreedyRouter.java     |   21 +-
 .../load/balancer/router/RegionRouteMap.java       |  172 +++
 .../confignode/manager/node/BaseNodeCache.java     |   20 +-
 .../manager/node/ConfigNodeHeartbeatCache.java     |   18 +-
 .../manager/node/DataNodeHeartbeatCache.java       |   16 +-
 .../iotdb/confignode/manager/node/NodeManager.java |  185 ++--
 .../manager/partition/PartitionManager.java        |  117 +-
 .../confignode/manager/partition/RegionCache.java  |    2 +-
 .../manager/partition/RegionGroupCache.java        |   67 +-
 .../manager/partition/RegionHeartbeatSample.java   |   11 +-
 .../iotdb/confignode/persistence/TriggerInfo.java  |   26 +-
 .../iotdb/confignode/persistence/UDFInfo.java      |   34 +-
 .../iotdb/confignode/persistence/cq/CQInfo.java    |  501 +++++++++
 .../persistence/executor/ConfigPlanExecutor.java   |   44 +-
 .../confignode/persistence/node/NodeInfo.java      |   22 +-
 .../persistence/node/NodeStatistics.java           |   24 +-
 .../persistence/partition/PartitionInfo.java       |  109 +-
 .../partition/StorageGroupPartitionTable.java      |   17 +
 .../statistics/RegionGroupStatistics.java          |   36 +-
 .../partition/statistics/RegionStatistics.java     |   48 +-
 .../persistence/schema/ClusterSchemaInfo.java      |   74 +-
 .../persistence/schema/TemplateTable.java          |   15 +
 .../iotdb/confignode/procedure/Procedure.java      |    2 +-
 .../procedure/env/ConfigNodeProcedureEnv.java      |   73 +-
 .../procedure/impl/CreateTriggerProcedure.java     |   11 +-
 .../procedure/impl/cq/CreateCQProcedure.java       |  263 +++++
 .../impl/node/RemoveDataNodeProcedure.java         |    2 +-
 .../procedure/impl/schema/DataNodeRegionTask.java  |    4 +-
 .../impl/schema/DeleteTimeSeriesProcedure.java     |    2 +-
 .../impl/schema/UnsetTemplateProcedure.java        |  426 +++++++
 .../statemachine/CreateRegionGroupsProcedure.java  |   11 +-
 .../procedure/impl/sync/CreatePipeProcedure.java   |   67 +-
 .../procedure/impl/sync/DropPipeProcedure.java     |   39 +-
 .../OperatePipeProcedureRollbackProcessor.java     |  122 ++
 .../procedure/impl/sync/StartPipeProcedure.java    |   97 +-
 .../procedure/impl/sync/StopPipeProcedure.java     |   97 +-
 .../procedure/state/CreateRegionGroupsState.java   |   13 +-
 .../procedure/state/cq/CreateCQState.java          |   15 +-
 .../procedure/state/schema/UnsetTemplateState.java |   17 +-
 .../procedure/store/ProcedureFactory.java          |   18 +-
 .../iotdb/confignode/service/ConfigNode.java       |    6 +-
 .../thrift/ConfigNodeRPCServiceProcessor.java      |   42 +-
 .../request/ConfigPhysicalPlanSerDeTest.java       |  166 ++-
 .../iotdb/confignode/cq/CQScheduleTaskTest.java    |   32 +-
 .../load/balancer/router/LazyGreedyRouterTest.java |  166 ---
 .../load/balancer/router/LeaderRouterTest.java     |  147 +--
 .../balancer/router/LoadScoreGreedyRouterTest.java |   32 +-
 .../load/balancer/router/RegionRouteMapTest.java   |   82 ++
 .../confignode/manager/node/NodeCacheTest.java     |   54 +
 .../manager/partition/RegionGroupCacheTest.java    |   53 +-
 .../iotdb/confignode/persistence/CQInfoTest.java   |  102 ++
 .../confignode/persistence/PartitionInfoTest.java  |   29 +-
 .../confignode/persistence/TriggerInfoTest.java    |    2 +
 .../statistics/RegionGroupStatisticsTest.java      |    4 +-
 .../partition/statistics/RegionStatisticsTest.java |    2 +-
 .../procedure/impl/CreateTriggerProcedureTest.java |    2 +
 .../procedure/impl/OperatePipeProcedureTest.java   |   52 +
 .../procedure/impl/UnsetTemplateProcedureTest.java |   75 ++
 .../confignode1conf/iotdb-confignode.properties    |    1 +
 .../confignode2conf/iotdb-confignode.properties    |    1 +
 .../confignode3conf/iotdb-confignode.properties    |    1 +
 ...java => ConsensusGroupModifyPeerException.java} |    8 +-
 .../multileader/MultiLeaderConsensus.java          |   22 +-
 .../multileader/MultiLeaderServerImpl.java         |  118 +-
 .../service/MultiLeaderRPCServiceProcessor.java    |   34 +-
 distribution/src/assembly/all.xml                  |    4 +
 distribution/src/assembly/confignode.xml           |    4 +
 distribution/src/assembly/datanode.xml             |    4 +
 docker/src/main/Dockerfile-1c1d-influxdb           |    2 +-
 docs/UserGuide/Alert/Alerting.md                   |    2 +-
 docs/UserGuide/Alert/Triggers.md                   |   16 +-
 docs/UserGuide/Process-Data/Continuous-Query.md    |  678 ++++++++----
 docs/UserGuide/Process-Data/Select-Into.md         |  425 ++++---
 .../Process-Data/UDF-User-Defined-Function.md      |    9 +-
 docs/zh/UserGuide/Alert/Alerting.md                |    2 +-
 docs/zh/UserGuide/Alert/Triggers.md                |   16 +-
 docs/zh/UserGuide/Process-Data/Continuous-Query.md |  681 ++++++++----
 docs/zh/UserGuide/Process-Data/Select-Into.md      |  425 ++++---
 .../Process-Data/UDF-User-Defined-Function.md      |    8 +-
 .../server/CustomizedJsonPayloadFormatter.java     |    2 +-
 integration-test/import-control.xml                |    1 +
 integration-test/src/assembly/mpp-test.xml         |   18 +-
 .../java/org/apache/iotdb/it/env/AbstractEnv.java  |    5 +-
 .../apache/iotdb/it/env/AbstractNodeWrapper.java   |   14 +-
 .../org/apache/iotdb/it/env/ConfigNodeWrapper.java |   12 +-
 .../org/apache/iotdb/it/env/DataNodeWrapper.java   |    9 +-
 .../org/apache/iotdb/it/env/RemoteServerEnv.java   |   13 +-
 ...figNodeIT.java => IoTDBClusterAuthorizeIT.java} |  344 +-----
 .../iotdb/confignode/it/IoTDBClusterNodeIT.java    |  308 ++++++
 .../confignode/it/IoTDBClusterPartitionIT.java     |    1 +
 .../it/IoTDBClusterRegionLeaderBalancingIT.java    |  152 +++
 .../confignode/it/IoTDBConfigNodeSnapshotIT.java   |   88 +-
 .../it/IoTDBConfigNodeSwitchLeaderIT.java          |   76 +-
 .../iotdb/confignode/it/IoTDBStorageGroupIT.java   |    3 +-
 .../org/apache/iotdb/db/it/cq/IoTDBCQExecIT.java   |  466 ++++++++
 .../java/org/apache/iotdb/db/it/cq/IoTDBCQIT.java  |  553 ++++++++++
 .../iotdb/db/it/schema/IoTDBSchemaTemplateIT.java  |   63 +-
 .../org/apache/iotdb/db/it/sync/IoTDBPipeIT.java   |  100 +-
 .../db/it/trigger/IoTDBTriggerExecutionIT.java     |   21 +-
 .../db/it/trigger/IoTDBTriggerManagementIT.java    |   49 +-
 .../IoTDBSessionInsertWithTriggerExecutionIT.java  |   21 +-
 .../apache/iotdb/integration/env/ClusterNode.java  |    2 +-
 .../iotdb/db/integration/IoTDBTracingIT.java       |    4 +-
 .../apache/iotdb/jdbc/IoTDBDatabaseMetadata.java   |  843 ++++++++------
 .../org/apache/iotdb/jdbc/IoTDBJDBCResultSet.java  | 1163 +++++++++++++++++++-
 .../java/org/apache/iotdb/jdbc/IoTDBStatement.java |   26 +-
 .../iotdb/jdbc/IoTDBDatabaseMetadataTest.java      |    4 +-
 .../apache/iotdb/jdbc/IoTDBJDBCResultSetTest.java  |   77 +-
 .../iotdb/jdbc/IoTDBPreparedStatementTest.java     |   34 +-
 .../dropwizard/DropwizardMetricManager.java        |    2 +-
 .../resources/conf/iotdb-confignode-metric.yml     |    2 +-
 .../resources/conf/iotdb-datanode-metric.yml       |    2 +-
 .../iotdb/metrics/AbstractMetricManager.java       |   47 +-
 .../apache/iotdb/metrics/config/MetricConfig.java  |    2 +-
 .../iotdb/metrics/impl/DoNothingMetricManager.java |    2 +-
 .../iotdb/metrics/utils/IoTDBMetricsUtils.java     |    2 +-
 .../org/apache/iotdb/metrics/utils/MetricInfo.java |    2 +-
 .../micrometer/MicrometerMetricManager.java        |    2 +-
 .../resources/conf/iotdb-common.properties         |  993 +++++++++--------
 .../iotdb/commons/cluster/RegionRoleType.java      |   10 +-
 .../apache/iotdb/commons/cluster/RegionStatus.java |    5 +
 .../apache/iotdb/commons/conf/CommonConfig.java    |    2 +
 .../java/org/apache/iotdb/commons/cq/CQState.java  |   27 +-
 .../org/apache/iotdb/commons/cq/TimeoutPolicy.java |   27 +-
 .../exception/sync/PipeSinkBeingUsedException.java |    2 +-
 .../commons/executable/ExecutableManager.java      |   51 +-
 .../apache/iotdb/commons/sync/pipe/PipeInfo.java   |    6 +-
 .../apache/iotdb/commons/sync/pipe/PipeStatus.java |   44 +-
 .../iotdb/commons/trigger/TriggerInformation.java  |   27 +-
 .../trigger/service/TriggerExecutableManager.java  |    2 +
 .../apache/iotdb/commons/udf/UDFInformation.java   |   37 +-
 .../org/apache/iotdb/commons/udf/UDFTable.java     |    3 +-
 .../commons/udf/service/UDFClassLoaderManager.java |    5 +-
 .../commons/udf/service/UDFExecutableManager.java  |   36 +-
 .../commons/udf/service/UDFManagementService.java  |    8 +-
 .../commons/utils/ThriftCommonsSerDeUtils.java     |   19 +
 .../resources/conf/schema-rocksdb.properties       |    8 +-
 .../schemaregion/rocksdb/RSchemaRegion.java        |    6 +
 .../schemaregion/rocksdb/mnode/RMNode.java         |   20 +
 .../assembly/resources/conf/schema-tag.properties  |    2 +-
 .../metadata/tagSchemaRegion/TagSchemaRegion.java  |    6 +
 .../resources/conf/iotdb-datanode.properties       | 1065 +-----------------
 server/src/assembly/server.xml                     |    4 +
 .../apache/iotdb/db/client/ConfigNodeClient.java   |  103 +-
 .../java/org/apache/iotdb/db/conf/IoTDBConfig.java |   15 +-
 .../org/apache/iotdb/db/conf/IoTDBDescriptor.java  |   58 +-
 .../exception/query/PathNumOverLimitException.java |    2 +-
 .../exception/sql/PathNumOverLimitException.java   |    2 +-
 .../iotdb/db/localconfignode/LocalConfigNode.java  |    2 +-
 .../idtable/entry/InsertMeasurementMNode.java      |   20 +
 .../org/apache/iotdb/db/metadata/mnode/IMNode.java |    8 +
 .../iotdb/db/metadata/mnode/InternalMNode.java     |   37 +-
 .../iotdb/db/metadata/mnode/MeasurementMNode.java  |   14 +
 .../iotdb/db/metadata/mtree/ConfigMTree.java       |   36 +-
 .../db/metadata/mtree/MTreeBelowSGMemoryImpl.java  |   23 +
 .../db/metadata/schemaregion/ISchemaRegion.java    |    2 +
 .../schemaregion/SchemaRegionMemoryImpl.java       |   10 +
 .../schemaregion/SchemaRegionSchemaFileImpl.java   |    6 +
 .../apache/iotdb/db/metadata/tag/TagLogFile.java   |    2 +-
 .../metadata/template/ClusterTemplateManager.java  |   28 +-
 .../metadata/template/TemplateInternalRPCUtil.java |   99 ++
 .../iotdb/db/mpp/common/MPPQueryContext.java       |   19 +-
 .../org/apache/iotdb/db/mpp/common/QueryId.java    |    2 +
 .../db/mpp/common/header/ColumnHeaderConstant.java |   11 +
 .../db/mpp/common/header/DatasetHeaderFactory.java |    4 +
 .../db/mpp/execution/exchange/ISourceHandle.java   |   10 +
 .../mpp/execution/exchange/LocalSourceHandle.java  |   21 +
 .../db/mpp/execution/exchange/SourceHandle.java    |   22 +-
 .../execution/executor/RegionWriteExecutor.java    |   31 +
 .../iotdb/db/mpp/plan/analyze/AnalyzeVisitor.java  |   12 +-
 .../apache/iotdb/db/mpp/plan/analyze/Analyzer.java |   22 +-
 .../db/mpp/plan/analyze/ExpressionAnalyzer.java    |   26 +
 .../db/mpp/plan/execution/IQueryExecution.java     |    3 +
 .../db/mpp/plan/execution/QueryExecution.java      |   33 +-
 .../mpp/plan/execution/config/ConfigExecution.java |   23 +-
 .../plan/execution/config/ConfigTaskVisitor.java   |   57 +-
 .../config/executor/ClusterConfigTaskExecutor.java |  368 +++++--
 .../config/executor/IConfigTaskExecutor.java       |   16 +
 .../executor/StandaloneConfigTaskExecutor.java     |   70 ++
 .../config/metadata/CreateContinuousQueryTask.java |   49 +
 .../config/metadata/DropContinuousQueryTask.java   |   42 +
 .../config/metadata/ShowContinuousQueriesTask.java |   75 ++
 .../metadata/template/DropSchemaTemplateTask.java} |   39 +-
 .../metadata/template/UnsetSchemaTemplateTask.java |   45 +
 .../plan/execution/memory/MemorySourceHandle.java  |   22 +
 .../iotdb/db/mpp/plan/parser/ASTVisitor.java       |  222 +++-
 .../db/mpp/plan/statement/StatementVisitor.java    |   30 +
 .../plan/statement/component/FillComponent.java    |   12 +
 .../plan/statement/component/FromComponent.java    |   12 +
 .../statement/component/GroupByLevelComponent.java |   17 +
 .../statement/component/GroupByTimeComponent.java  |   35 +
 .../plan/statement/component/HavingCondition.java  |    4 +
 .../plan/statement/component/IntoComponent.java    |   12 +
 .../db/mpp/plan/statement/component/IntoItem.java  |   15 +
 .../plan/statement/component/OrderByComponent.java |   12 +
 .../plan/statement/component/SelectComponent.java  |   21 +-
 .../db/mpp/plan/statement/component/SortItem.java  |    4 +
 .../plan/statement/component/WhereCondition.java   |    4 +
 .../db/mpp/plan/statement/crud/QueryStatement.java |   97 +-
 .../mpp/plan/statement/literal/BooleanLiteral.java |    5 +
 .../mpp/plan/statement/literal/DoubleLiteral.java  |    5 +
 .../db/mpp/plan/statement/literal/LongLiteral.java |    5 +
 .../db/mpp/plan/statement/literal/NullLiteral.java |    5 +
 .../mpp/plan/statement/literal/StringLiteral.java  |    5 +
 .../metadata/CreateContinuousQueryStatement.java   |  215 ++++
 .../metadata/CreateFunctionStatement.java          |   19 +-
 .../statement/metadata/CreateTriggerStatement.java |   19 +-
 ...ment.java => DropContinuousQueryStatement.java} |   46 +-
 ...nt.java => ShowContinuousQueriesStatement.java} |   48 +-
 .../DropSchemaTemplateStatement.java}              |   51 +-
 .../UnsetSchemaTemplateStatement.java}             |   44 +-
 .../apache/iotdb/db/qp/sql/IoTDBSqlVisitor.java    |   59 +-
 .../iotdb/db/query/control/SessionManager.java     |    1 -
 .../java/org/apache/iotdb/db/service/DataNode.java |   61 +-
 .../db/service/metrics/DataNodeMetricsHelper.java  |    2 +-
 .../iotdb/db/service/metrics/SystemMetrics.java    |   22 +-
 .../service/thrift/impl/ClientRPCServiceImpl.java  |  450 ++++----
 .../impl/DataNodeInternalRPCServiceImpl.java       |  208 +++-
 .../db/service/thrift/impl/TSServiceImpl.java      |   32 +
 .../java/org/apache/iotdb/db/sync/SyncService.java |   84 +-
 .../db/sync/transport/server/ReceiverManager.java  |    2 +-
 .../trigger/service/TriggerManagementService.java  |    4 +-
 .../apache/iotdb/db/utils/QueryDataSetUtils.java   |   27 +
 .../apache/iotdb/db/utils/sync/SyncPipeUtil.java   |    8 +-
 .../apache/iotdb/db/conf/IoTDBDescriptorTest.java  |    6 +-
 .../iotdb/db/metadata/mtree/ConfigMTreeTest.java   |    2 +-
 .../apache/iotdb/db/metric/MetricServiceTest.java  |   19 +
 .../db/mpp/plan/StandaloneCoordinatorTest.java     |    3 +-
 service-rpc/pom.xml                                |    8 +
 .../java/org/apache/iotdb/rpc/IoTDBRpcDataSet.java |  345 +++---
 .../java/org/apache/iotdb/rpc/TSStatusCode.java    |   10 +-
 .../apache/iotdb/session/SessionConnection.java    |   22 +-
 .../org/apache/iotdb/session/SessionDataSet.java   |   29 +-
 .../src/main/thrift/confignode.thrift              |  121 +-
 .../src/main/thrift/mutlileader.thrift             |   11 +
 thrift/src/main/thrift/client.thrift               |   20 +-
 thrift/src/main/thrift/datanode.thrift             |   43 +-
 .../iotdb/tsfile/common/conf/TSFileConfig.java     |    2 +-
 .../iotdb/tsfile/utils/ReadWriteIOUtils.java       |   28 +
 284 files changed, 14729 insertions(+), 6027 deletions(-)

diff --cc 
server/src/main/java/org/apache/iotdb/db/query/control/SessionManager.java
index 1270cd9e57,6ec0348c9c..bfde5cc9d3
--- a/server/src/main/java/org/apache/iotdb/db/query/control/SessionManager.java
+++ b/server/src/main/java/org/apache/iotdb/db/query/control/SessionManager.java
@@@ -232,14 -224,29 +232,13 @@@ public class SessionManager implements 
      return isLoggedIn;
    }
  
 -  public long requestSessionId(
 -      String username, String zoneId, IoTDBConstant.ClientVersion 
clientVersion) {
 -    long sessionId = sessionIdGenerator.incrementAndGet();
 -
 -    currSessionId.set(sessionId);
 -    sessionIdToUsername.put(sessionId, username);
 -    sessionIdToZoneId.put(sessionId, ZoneId.of(zoneId));
 -    sessionIdToClientVersion.put(sessionId, clientVersion);
 -    sessionIdToSessionInfo.put(sessionId, new SessionInfo(sessionId, 
username, zoneId));
 -
 -    return sessionId;
 -  }
 -
 -  public boolean releaseSessionResource(long sessionId) {
 -    return releaseSessionResource(sessionId, 
this::releaseQueryResourceNoExceptions);
 +  public boolean releaseSessionResource(IClientSession session) {
 +    return releaseSessionResource(session, 
this::releaseQueryResourceNoExceptions);
    }
  
 -  public boolean releaseSessionResource(long sessionId, Consumer<Long> 
releaseQueryResource) {
 -    sessionIdToZoneId.remove(sessionId);
 -    sessionIdToClientVersion.remove(sessionId);
 -    sessionIdToSessionInfo.remove(sessionId);
 -
 -    Set<Long> statementIdSet = sessionIdToStatementId.remove(sessionId);
 +  public boolean releaseSessionResource(
 +      IClientSession session, Consumer<Long> releaseQueryResource) {
 +    Set<Long> statementIdSet = session.getStatementIds();
- 
      if (statementIdSet != null) {
        for (Long statementId : statementIdSet) {
          Set<Long> queryIdSet = statementIdToQueryId.remove(statementId);
diff --cc 
server/src/main/java/org/apache/iotdb/db/service/thrift/impl/ClientRPCServiceImpl.java
index cdc33165e1,6c0873f5f5..75e57e5b4c
--- 
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
@@@ -158,6 -175,246 +175,246 @@@ public class ClientRPCServiceImpl imple
      }
    }
  
+   private TSExecuteStatementResp executeStatementInternal(
+       TSExecuteStatementReq req, SelectResult setResult) {
+     String statement = req.getStatement();
 -    if (!SESSION_MANAGER.checkLogin(req.getSessionId())) {
++    if (!SESSION_MANAGER.checkLogin(SESSION_MANAGER.getCurrSession())) {
+       return RpcUtils.getTSExecuteStatementResp(getNotLoggedInStatus());
+     }
+ 
+     long startTime = System.currentTimeMillis();
+     try {
+       Statement s =
+           StatementGenerator.createStatement(
 -              statement, SESSION_MANAGER.getZoneId(req.getSessionId()));
++              statement, SESSION_MANAGER.getCurrSession().getZoneId());
+ 
+       if (s == null) {
+         return RpcUtils.getTSExecuteStatementResp(
+             RpcUtils.getStatus(
+                 TSStatusCode.SQL_PARSE_ERROR, "This operation type is not 
supported"));
+       }
+       // permission check
 -      TSStatus status = AuthorityChecker.checkAuthority(s, req.sessionId);
++      TSStatus status = AuthorityChecker.checkAuthority(s, 
SESSION_MANAGER.getCurrSession());
+       if (status.getCode() != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
+         return RpcUtils.getTSExecuteStatementResp(status);
+       }
+ 
+       QUERY_FREQUENCY_RECORDER.incrementAndGet();
+       AUDIT_LOGGER.debug("Session {} execute Query: {}", req.sessionId, 
statement);
+ 
 -      long queryId = SESSION_MANAGER.requestQueryId(req.statementId, true);
++      long queryId = SESSION_MANAGER.requestQueryId(req.statementId, false);
+       // create and cache dataset
+       ExecutionResult result =
+           COORDINATOR.execute(
+               s,
+               queryId,
 -              SESSION_MANAGER.getSessionInfo(req.sessionId),
++              
SESSION_MANAGER.getSessionInfo(SESSION_MANAGER.getCurrSession()),
+               statement,
+               PARTITION_FETCHER,
+               SCHEMA_FETCHER,
+               req.getTimeout());
+ 
+       if (result.status.code != TSStatusCode.SUCCESS_STATUS.getStatusCode()
+           && result.status.code != 
TSStatusCode.NEED_REDIRECTION.getStatusCode()) {
+         return RpcUtils.getTSExecuteStatementResp(result.status);
+       }
+ 
+       IQueryExecution queryExecution = COORDINATOR.getQueryExecution(queryId);
+ 
+       try (SetThreadName threadName = new 
SetThreadName(result.queryId.getId())) {
+         TSExecuteStatementResp resp;
+         if (queryExecution != null && queryExecution.isQuery()) {
+           resp = createResponse(queryExecution.getDatasetHeader(), queryId);
+           resp.setStatus(result.status);
+           setResult.apply(resp, queryExecution, req.fetchSize);
+         } else {
+           resp = RpcUtils.getTSExecuteStatementResp(result.status);
+         }
+         return resp;
+       }
+     } catch (Exception e) {
+       return RpcUtils.getTSExecuteStatementResp(
+           onQueryException(e, "\"" + statement + "\". " + 
OperationType.EXECUTE_STATEMENT));
+     } finally {
+       addOperationLatency(Operation.EXECUTE_QUERY, startTime);
+       long costTime = System.currentTimeMillis() - startTime;
+       if (costTime >= CONFIG.getSlowQueryThreshold()) {
+         SLOW_SQL_LOGGER.info("Cost: {} ms, sql is {}", costTime, statement);
+       }
+     }
+   }
+ 
+   private TSExecuteStatementResp executeRawDataQueryInternal(
+       TSRawDataQueryReq req, SelectResult setResult) {
 -    if (!SESSION_MANAGER.checkLogin(req.getSessionId())) {
++    if (!SESSION_MANAGER.checkLogin(SESSION_MANAGER.getCurrSession())) {
+       return RpcUtils.getTSExecuteStatementResp(getNotLoggedInStatus());
+     }
+     long startTime = System.currentTimeMillis();
+     try {
+       Statement s =
 -          StatementGenerator.createStatement(req, 
SESSION_MANAGER.getZoneId(req.getSessionId()));
++          StatementGenerator.createStatement(req, 
SESSION_MANAGER.getCurrSession().getZoneId());
+ 
+       // permission check
 -      TSStatus status = AuthorityChecker.checkAuthority(s, req.sessionId);
++      TSStatus status = AuthorityChecker.checkAuthority(s, 
SESSION_MANAGER.getCurrSession());
+       if (status.getCode() != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
+         return RpcUtils.getTSExecuteStatementResp(status);
+       }
+ 
+       QUERY_FREQUENCY_RECORDER.incrementAndGet();
+       AUDIT_LOGGER.debug("Session {} execute Raw Data Query: {}", 
req.sessionId, req);
 -      long queryId = SESSION_MANAGER.requestQueryId(req.statementId, true);
++      long queryId = SESSION_MANAGER.requestQueryId(req.statementId, false);
+       // create and cache dataset
+       ExecutionResult result =
+           COORDINATOR.execute(
+               s,
+               queryId,
 -              SESSION_MANAGER.getSessionInfo(req.sessionId),
++              
SESSION_MANAGER.getSessionInfo(SESSION_MANAGER.getCurrSession()),
+               "",
+               PARTITION_FETCHER,
+               SCHEMA_FETCHER,
+               req.getTimeout());
+ 
+       if (result.status.code != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
+         throw new RuntimeException("error code: " + result.status);
+       }
+ 
+       IQueryExecution queryExecution = COORDINATOR.getQueryExecution(queryId);
+ 
+       try (SetThreadName threadName = new 
SetThreadName(result.queryId.getId())) {
+         TSExecuteStatementResp resp;
+         if (queryExecution.isQuery()) {
+           resp = createResponse(queryExecution.getDatasetHeader(), queryId);
+           resp.setStatus(result.status);
+           setResult.apply(resp, queryExecution, req.fetchSize);
+         } else {
+           resp = RpcUtils.getTSExecuteStatementResp(result.status);
+         }
+         return resp;
+       }
+     } catch (Exception e) {
+       // TODO call the coordinator to release query resource
+       return RpcUtils.getTSExecuteStatementResp(
+           onQueryException(e, "\"" + req + "\". " + 
OperationType.EXECUTE_RAW_DATA_QUERY));
+     } finally {
+       addOperationLatency(Operation.EXECUTE_QUERY, startTime);
+       long costTime = System.currentTimeMillis() - startTime;
+       if (costTime >= CONFIG.getSlowQueryThreshold()) {
+         SLOW_SQL_LOGGER.info("Cost: {} ms, sql is {}", costTime, req);
+       }
+     }
+   }
+ 
+   private TSExecuteStatementResp executeLastDataQueryInternal(
+       TSLastDataQueryReq req, SelectResult setResult) {
 -    if (!SESSION_MANAGER.checkLogin(req.getSessionId())) {
++    if (!SESSION_MANAGER.checkLogin(SESSION_MANAGER.getCurrSession())) {
+       return RpcUtils.getTSExecuteStatementResp(getNotLoggedInStatus());
+     }
+     long startTime = System.currentTimeMillis();
+     try {
+       Statement s =
 -          StatementGenerator.createStatement(req, 
SESSION_MANAGER.getZoneId(req.getSessionId()));
++          StatementGenerator.createStatement(req, 
SESSION_MANAGER.getCurrSession().getZoneId());
+       // permission check
 -      TSStatus status = AuthorityChecker.checkAuthority(s, req.sessionId);
++      TSStatus status = AuthorityChecker.checkAuthority(s, 
SESSION_MANAGER.getCurrSession());
+       if (status.getCode() != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
+         return RpcUtils.getTSExecuteStatementResp(status);
+       }
+       QUERY_FREQUENCY_RECORDER.incrementAndGet();
+       AUDIT_LOGGER.debug("Session {} execute Last Data Query: {}", 
req.sessionId, req);
 -      long queryId = SESSION_MANAGER.requestQueryId(req.statementId, true);
++      long queryId = SESSION_MANAGER.requestQueryId(req.statementId, false);
+       // create and cache dataset
+       ExecutionResult result =
+           COORDINATOR.execute(
+               s,
+               queryId,
 -              SESSION_MANAGER.getSessionInfo(req.sessionId),
++              
SESSION_MANAGER.getSessionInfo(SESSION_MANAGER.getCurrSession()),
+               "",
+               PARTITION_FETCHER,
+               SCHEMA_FETCHER,
+               req.getTimeout());
+ 
+       if (result.status.code != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
+         throw new RuntimeException("error code: " + result.status);
+       }
+ 
+       IQueryExecution queryExecution = COORDINATOR.getQueryExecution(queryId);
+ 
+       try (SetThreadName threadName = new 
SetThreadName(result.queryId.getId())) {
+         TSExecuteStatementResp resp;
+         if (queryExecution.isQuery()) {
+           resp = createResponse(queryExecution.getDatasetHeader(), queryId);
+           resp.setStatus(result.status);
+           setResult.apply(resp, queryExecution, req.fetchSize);
+         } else {
+           resp = RpcUtils.getTSExecuteStatementResp(result.status);
+         }
+         return resp;
+       }
+ 
+     } catch (Exception e) {
+       // TODO call the coordinator to release query resource
+       return RpcUtils.getTSExecuteStatementResp(
+           onQueryException(e, "\"" + req + "\". " + 
OperationType.EXECUTE_LAST_DATA_QUERY));
+     } finally {
+       addOperationLatency(Operation.EXECUTE_QUERY, startTime);
+       long costTime = System.currentTimeMillis() - startTime;
+       if (costTime >= CONFIG.getSlowQueryThreshold()) {
+         SLOW_SQL_LOGGER.info("Cost: {} ms, sql is {}", costTime, req);
+       }
+     }
+   }
+ 
+   @Override
+   public TSExecuteStatementResp executeQueryStatementV2(TSExecuteStatementReq 
req) {
+     return executeStatementV2(req);
+   }
+ 
+   @Override
+   public TSExecuteStatementResp 
executeUpdateStatementV2(TSExecuteStatementReq req) {
+     return executeStatementV2(req);
+   }
+ 
+   @Override
+   public TSExecuteStatementResp executeStatementV2(TSExecuteStatementReq req) 
{
+     return executeStatementInternal(req, SELECT_RESULT);
+   }
+ 
+   @Override
+   public TSExecuteStatementResp executeRawDataQueryV2(TSRawDataQueryReq req) {
+     return executeRawDataQueryInternal(req, SELECT_RESULT);
+   }
+ 
+   @Override
+   public TSExecuteStatementResp executeLastDataQueryV2(TSLastDataQueryReq 
req) {
+     return executeLastDataQueryInternal(req, SELECT_RESULT);
+   }
+ 
+   @Override
+   public TSFetchResultsResp fetchResultsV2(TSFetchResultsReq req) {
+     try {
 -      if (!SESSION_MANAGER.checkLogin(req.getSessionId())) {
++      if (!SESSION_MANAGER.checkLogin(SESSION_MANAGER.getCurrSession())) {
+         return RpcUtils.getTSFetchResultsResp(getNotLoggedInStatus());
+       }
+       TSFetchResultsResp resp = 
RpcUtils.getTSFetchResultsResp(TSStatusCode.SUCCESS_STATUS);
+ 
+       IQueryExecution queryExecution = 
COORDINATOR.getQueryExecution(req.queryId);
+       try (SetThreadName queryName = new 
SetThreadName(queryExecution.getQueryId())) {
+         List<ByteBuffer> result =
+             QueryDataSetUtils.convertQueryResultByFetchSize(queryExecution, 
req.fetchSize);
+         boolean hasResultSet = !(result.size() == 0);
+         resp.setHasResultSet(hasResultSet);
+         resp.setIsAlign(true);
+         resp.setQueryResult(result);
+         QUERY_TIME_MANAGER.unRegisterQuery(req.queryId, false);
+         if (!hasResultSet) {
+           COORDINATOR.removeQueryExecution(req.queryId);
+         }
+         return resp;
+       }
+     } catch (Exception e) {
+       return RpcUtils.getTSFetchResultsResp(onQueryException(e, 
OperationType.FETCH_RESULTS));
+     }
+   }
+ 
    @Override
    public TSOpenSessionResp openSession(TSOpenSessionReq req) throws 
TException {
      IoTDBConstant.ClientVersion clientVersion = parseClientVersion(req);
diff --cc 
server/src/main/java/org/apache/iotdb/db/service/thrift/impl/DataNodeInternalRPCServiceImpl.java
index 3afa2ab07c,fbeaabd90e..bc1bc46fae
--- 
a/server/src/main/java/org/apache/iotdb/db/service/thrift/impl/DataNodeInternalRPCServiceImpl.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/service/thrift/impl/DataNodeInternalRPCServiceImpl.java
@@@ -85,6 -103,9 +103,11 @@@ import org.apache.iotdb.db.mpp.plan.pla
  import 
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write.RollbackSchemaBlackListNode;
  import org.apache.iotdb.db.mpp.plan.planner.plan.node.write.DeleteDataNode;
  import org.apache.iotdb.db.mpp.plan.scheduler.load.LoadTsFileScheduler;
+ import org.apache.iotdb.db.mpp.plan.statement.component.WhereCondition;
+ import org.apache.iotdb.db.mpp.plan.statement.crud.QueryStatement;
+ import org.apache.iotdb.db.query.control.SessionManager;
++import org.apache.iotdb.db.query.control.clientsession.IClientSession;
++import org.apache.iotdb.db.query.control.clientsession.InternalClientSession;
  import org.apache.iotdb.db.service.DataNode;
  import org.apache.iotdb.db.service.RegionMigrateService;
  import org.apache.iotdb.db.sync.SyncService;
@@@ -703,6 -785,122 +787,126 @@@ public class DataNodeInternalRPCService
      }
    }
  
+   @Override
+   public TSStatus operatePipeOnDataNodeForRollback(TOperatePipeOnDataNodeReq 
req) {
+     // Operate PIPE on DataNode for rollback, createTime in req is required.
+     switch (SyncOperation.values()[req.getOperation()]) {
+       case START_PIPE:
+         SyncService.getInstance().startPipe(req.getPipeName(), 
req.getCreateTime());
+         break;
+       case STOP_PIPE:
+         SyncService.getInstance().stopPipe(req.getPipeName(), 
req.getCreateTime());
+         break;
+       case DROP_PIPE:
+         SyncService.getInstance().dropPipe(req.getPipeName(), 
req.getCreateTime());
+         break;
+       default:
+         return new TSStatus(TSStatusCode.PIPE_ERROR.getStatusCode())
+             .setMessage("Unsupported operation.");
+     }
+     return RpcUtils.SUCCESS_STATUS;
+   }
+ 
+   @Override
+   public TSStatus executeCQ(TExecuteCQ req) {
+ 
 -    long sessionId =
 -        SESSION_MANAGER.requestSessionId(
 -            req.username, req.zoneId, IoTDBConstant.ClientVersion.V_0_13);
++    IClientSession session = new InternalClientSession(req.cqId);
++
++    SESSION_MANAGER.registerSession(session);
++
++    SESSION_MANAGER.supplySession(
++        session, req.getUsername(), req.getZoneId(), 
IoTDBConstant.ClientVersion.V_0_13);
++
+     String executedSQL = req.queryBody;
+ 
+     try {
+       QueryStatement s =
+           (QueryStatement)
+               StatementGenerator.createStatement(
 -                  req.queryBody, SESSION_MANAGER.getZoneId(sessionId));
++                  req.queryBody, 
SESSION_MANAGER.getCurrSession().getZoneId());
+       if (s == null) {
+         return RpcUtils.getStatus(
+             TSStatusCode.SQL_PARSE_ERROR, "This operation type is not 
supported");
+       }
+ 
+       // 1. add time filter in where
+       Expression timeFilter =
+           new LogicAndExpression(
+               new GreaterEqualExpression(
+                   new TimestampOperand(),
+                   new ConstantOperand(TSDataType.INT64, 
String.valueOf(req.startTime))),
+               new LessThanExpression(
+                   new TimestampOperand(),
+                   new ConstantOperand(TSDataType.INT64, 
String.valueOf(req.endTime))));
+       if (s.getWhereCondition() != null) {
+         s.getWhereCondition()
+             .setPredicate(new LogicAndExpression(timeFilter, 
s.getWhereCondition().getPredicate()));
+       } else {
+         s.setWhereCondition(new WhereCondition(timeFilter));
+       }
+ 
 -      // 2. add time rage in group by time
++      // 2. add time range in group by time
+       if (s.getGroupByTimeComponent() != null) {
+         s.getGroupByTimeComponent().setStartTime(req.startTime);
+         s.getGroupByTimeComponent().setEndTime(req.endTime);
+         s.getGroupByTimeComponent().setLeftCRightO(true);
+       }
+       executedSQL = String.join(" ", 
s.constructFormattedSQL().split("\n")).replaceAll(" +", " ");
+ 
+       QUERY_FREQUENCY_RECORDER.incrementAndGet();
+ 
+       long queryId =
 -          
SESSION_MANAGER.requestQueryId(SESSION_MANAGER.requestStatementId(sessionId), 
true);
++          
SESSION_MANAGER.requestQueryId(SESSION_MANAGER.requestStatementId(session), 
false);
+       // create and cache dataset
+       ExecutionResult result =
+           COORDINATOR.execute(
+               s,
+               queryId,
 -              SESSION_MANAGER.getSessionInfo(sessionId),
++              SESSION_MANAGER.getSessionInfo(session),
+               executedSQL,
+               PARTITION_FETCHER,
+               SCHEMA_FETCHER,
+               req.getTimeout());
+ 
+       if (result.status.code != TSStatusCode.SUCCESS_STATUS.getStatusCode()
+           && result.status.code != 
TSStatusCode.NEED_REDIRECTION.getStatusCode()) {
+         return result.status;
+       }
+ 
+       IQueryExecution queryExecution = COORDINATOR.getQueryExecution(queryId);
+ 
+       try (SetThreadName threadName = new 
SetThreadName(result.queryId.getId())) {
+         if (queryExecution != null) {
+           // consume up all the result
+           while (true) {
+             Optional<TsBlock> optionalTsBlock = 
queryExecution.getBatchResult();
+             if (!optionalTsBlock.isPresent()) {
+               break;
+             }
+           }
+         }
+         return result.status;
+       }
+     } catch (Exception e) {
+       // TODO call the coordinator to release query resource
+       return onQueryException(e, "\"" + executedSQL + "\". " + 
OperationType.EXECUTE_STATEMENT);
+     } finally {
 -      SESSION_MANAGER.releaseSessionResource(sessionId, 
this::cleanupQueryExecution);
 -      SESSION_MANAGER.closeSession(sessionId);
++      SESSION_MANAGER.releaseSessionResource(session, 
this::cleanupQueryExecution);
++      SESSION_MANAGER.closeSession(session);
+     }
+   }
+ 
+   private void cleanupQueryExecution(Long queryId) {
+     IQueryExecution queryExecution = COORDINATOR.getQueryExecution(queryId);
+     if (queryExecution != null) {
+       try (SetThreadName threadName = new 
SetThreadName(queryExecution.getQueryId())) {
+         LOGGER.info("[CleanUpQuery]]");
+         queryExecution.stopAndCleanup();
+         COORDINATOR.removeQueryExecution(queryId);
+       }
+     }
+   }
+ 
    private PathPatternTree filterPathPatternTree(PathPatternTree patternTree, 
String storageGroup) {
      PathPatternTree filteredPatternTree = new PathPatternTree();
      try {
diff --cc thrift/src/main/thrift/client.thrift
index 1c8b1e048e,dcf2171d13..61744021c8
--- a/thrift/src/main/thrift/client.thrift
+++ b/thrift/src/main/thrift/client.thrift
@@@ -429,24 -432,20 +432,37 @@@ struct TSyncTransportMetaInfo
    2:required i64 startIndex
  }
  
 +enum TSConnectionType {
 +  THRIFT_BASED
 +  MQTT_BASED
 +  INTERNAL
 +}
 +
 +struct TSConnectionInfo {
 +  1: required string userName
 +  2: required i64 logInTime
 +  3: required string connectionId // ip:port for thrift-based service and 
clientId for mqtt-based service
 +  4: required TSConnectionType type
 +}
 +
 +struct TSConnectionInfoResp {
 +  1: required list<TSConnectionInfo> connectionInfoList
 +}
 +
  service IClientRPCService {
+ 
+   TSExecuteStatementResp executeQueryStatementV2(1:TSExecuteStatementReq req);
+ 
+   TSExecuteStatementResp executeUpdateStatementV2(1:TSExecuteStatementReq 
req);
+ 
+   TSExecuteStatementResp executeStatementV2(1:TSExecuteStatementReq req);
+ 
+   TSExecuteStatementResp executeRawDataQueryV2(1:TSRawDataQueryReq req);
+ 
+   TSExecuteStatementResp executeLastDataQueryV2(1:TSLastDataQueryReq req);
+ 
+   TSFetchResultsResp fetchResultsV2(1:TSFetchResultsReq req);
+ 
    TSOpenSessionResp openSession(1:TSOpenSessionReq req);
  
    common.TSStatus closeSession(1:TSCloseSessionReq req);


Reply via email to