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);
