This is an automated email from the ASF dual-hosted git repository. rong pushed a commit to branch select-into in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit f0c5770a7856f4b0065163aa3a62f4f20fdc237d Author: Steve Yurong Su <[email protected]> AuthorDate: Wed Jul 14 11:42:54 2021 +0800 refactor internalExecuteQueryStatement --- .../apache/iotdb/db/qp/physical/PhysicalPlan.java | 4 +- .../iotdb/db/query/control/QueryTimeManager.java | 14 +++++ .../iotdb/db/query/control/TracingManager.java | 49 ++++++++++++--- .../org/apache/iotdb/db/service/TSServiceImpl.java | 71 ++++++++-------------- .../iotdb/db/query/control/TracingManagerTest.java | 19 ++++-- 5 files changed, 98 insertions(+), 59 deletions(-) diff --git a/server/src/main/java/org/apache/iotdb/db/qp/physical/PhysicalPlan.java b/server/src/main/java/org/apache/iotdb/db/qp/physical/PhysicalPlan.java index 74b4f05..af58de9 100644 --- a/server/src/main/java/org/apache/iotdb/db/qp/physical/PhysicalPlan.java +++ b/server/src/main/java/org/apache/iotdb/db/qp/physical/PhysicalPlan.java @@ -219,7 +219,9 @@ public abstract class PhysicalPlan { } public void setLoginUserName(String loginUserName) { - this.loginUserName = loginUserName; + if (this instanceof AuthorPlan) { + this.loginUserName = loginUserName; + } } public static class Factory { diff --git a/server/src/main/java/org/apache/iotdb/db/query/control/QueryTimeManager.java b/server/src/main/java/org/apache/iotdb/db/query/control/QueryTimeManager.java index 3d78a92..48bb732 100644 --- a/server/src/main/java/org/apache/iotdb/db/query/control/QueryTimeManager.java +++ b/server/src/main/java/org/apache/iotdb/db/query/control/QueryTimeManager.java @@ -22,6 +22,8 @@ import org.apache.iotdb.db.concurrent.IoTDBThreadPoolFactory; import org.apache.iotdb.db.conf.IoTDBConfig; import org.apache.iotdb.db.conf.IoTDBDescriptor; import org.apache.iotdb.db.exception.query.QueryTimeoutRuntimeException; +import org.apache.iotdb.db.qp.physical.PhysicalPlan; +import org.apache.iotdb.db.qp.physical.sys.ShowQueryProcesslistPlan; import org.apache.iotdb.db.service.IService; import org.apache.iotdb.db.service.ServiceType; @@ -76,6 +78,14 @@ public class QueryTimeManager implements IService { queryScheduledTaskMap.put(queryId, scheduledFuture); } + public void registerQuery( + long queryId, long startTime, String sql, long timeout, PhysicalPlan plan) { + if (plan instanceof ShowQueryProcesslistPlan) { + return; + } + registerQuery(queryId, startTime, sql, timeout); + } + public void killQuery(long queryId) { if (queryInfoMap.get(queryId) == null) { return; @@ -101,6 +111,10 @@ public class QueryTimeManager implements IService { return successRemoved; } + public AtomicBoolean unRegisterQuery(long queryId, PhysicalPlan plan) { + return plan instanceof ShowQueryProcesslistPlan ? null : unRegisterQuery(queryId); + } + public static void checkQueryAlive(long queryId) { QueryInfo queryInfo = getInstance().queryInfoMap.get(queryId); if (queryInfo != null && queryInfo.isInterrupted()) { diff --git a/server/src/main/java/org/apache/iotdb/db/query/control/TracingManager.java b/server/src/main/java/org/apache/iotdb/db/query/control/TracingManager.java index 450c4f4..31a071f 100644 --- a/server/src/main/java/org/apache/iotdb/db/query/control/TracingManager.java +++ b/server/src/main/java/org/apache/iotdb/db/query/control/TracingManager.java @@ -18,10 +18,16 @@ */ package org.apache.iotdb.db.query.control; +import org.apache.iotdb.db.conf.IoTDBConfig; import org.apache.iotdb.db.conf.IoTDBConstant; import org.apache.iotdb.db.conf.IoTDBDescriptor; import org.apache.iotdb.db.engine.fileSystem.SystemFileFactory; import org.apache.iotdb.db.engine.storagegroup.TsFileResource; +import org.apache.iotdb.db.qp.physical.PhysicalPlan; +import org.apache.iotdb.db.qp.physical.crud.AlignByDevicePlan; +import org.apache.iotdb.db.qp.physical.crud.QueryPlan; +import org.apache.iotdb.db.query.dataset.AlignByDeviceDataSet; +import org.apache.iotdb.tsfile.read.query.dataset.QueryDataSet; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -39,12 +45,15 @@ import java.util.concurrent.ConcurrentHashMap; public class TracingManager { private static final Logger logger = LoggerFactory.getLogger(TracingManager.class); + private static final IoTDBConfig config = IoTDBDescriptor.getInstance().getConfig(); + private static final String QUERY_ID = "Query Id: "; private static final String DATE_FORMAT = "yyyy-MM-dd HH:mm:ss.SSS"; + private BufferedWriter writer; private Map<Long, Long> queryStartTime = new ConcurrentHashMap<>(); - public TracingManager(String dirName, String logFileName) { + private TracingManager(String dirName, String logFileName) { initTracingManager(dirName, logFileName); } @@ -71,8 +80,12 @@ public class TracingManager { return TracingManagerHelper.INSTANCE; } - public void writeQueryInfo(long queryId, String statement, long startTime, int pathsNum) + public void writeQueryInfo(long queryId, String statement, long startTime, PhysicalPlan plan) throws IOException { + if (!config.isEnablePerformanceTracing() || !(plan instanceof QueryPlan)) { + return; + } + queryStartTime.put(queryId, startTime); StringBuilder builder = new StringBuilder(); builder @@ -85,15 +98,19 @@ public class TracingManager { .append(" - Start time: ") .append(new SimpleDateFormat(DATE_FORMAT).format(startTime)) .append("\n" + QUERY_ID) - .append(queryId) - .append(" - Number of series paths: ") - .append(pathsNum) - .append("\n"); + .append(queryId); + if (plan instanceof AlignByDevicePlan) { + builder.append(" - Number of series paths: ").append(plan.getPaths().size()).append("\n"); + } writer.write(builder.toString()); } // for align by device query public void writeQueryInfo(long queryId, String statement, long startTime) throws IOException { + if (!config.isEnablePerformanceTracing()) { + return; + } + queryStartTime.put(queryId, startTime); StringBuilder builder = new StringBuilder(); builder @@ -109,12 +126,16 @@ public class TracingManager { writer.write(builder.toString()); } - public void writePathsNum(long queryId, int pathsNum) throws IOException { + public void writePathsNum(long queryId, QueryDataSet dataSet) throws IOException { + if (!config.isEnablePerformanceTracing() || !(dataSet instanceof AlignByDeviceDataSet)) { + return; + } + StringBuilder builder = new StringBuilder(QUERY_ID) .append(queryId) .append(" - Number of series paths: ") - .append(pathsNum) + .append(((AlignByDeviceDataSet) dataSet).getPathsNum()) .append("\n"); writer.write(builder.toString()); } @@ -122,6 +143,10 @@ public class TracingManager { public void writeTsFileInfo( long queryId, Set<TsFileResource> seqFileResources, Set<TsFileResource> unSeqFileResources) throws IOException { + if (!config.isEnablePerformanceTracing()) { + return; + } + // to avoid the disorder info of multi query // add query id as prefix of each info StringBuilder builder = @@ -176,6 +201,10 @@ public class TracingManager { public void writeChunksInfo(long queryId, long totalChunkNum, long totalChunkSize) throws IOException { + if (!config.isEnablePerformanceTracing()) { + return; + } + StringBuilder builder = new StringBuilder(QUERY_ID) .append(queryId) @@ -190,6 +219,10 @@ public class TracingManager { } public void writeEndTime(long queryId) throws IOException { + if (!config.isEnablePerformanceTracing()) { + return; + } + long endTime = System.currentTimeMillis(); StringBuilder builder = new StringBuilder(QUERY_ID) diff --git a/server/src/main/java/org/apache/iotdb/db/service/TSServiceImpl.java b/server/src/main/java/org/apache/iotdb/db/service/TSServiceImpl.java index 0616f61..f1af737 100644 --- a/server/src/main/java/org/apache/iotdb/db/service/TSServiceImpl.java +++ b/server/src/main/java/org/apache/iotdb/db/service/TSServiceImpl.java @@ -77,7 +77,6 @@ import org.apache.iotdb.db.query.control.QueryTimeManager; import org.apache.iotdb.db.query.control.SessionManager; import org.apache.iotdb.db.query.control.SessionTimeoutManager; import org.apache.iotdb.db.query.control.TracingManager; -import org.apache.iotdb.db.query.dataset.AlignByDeviceDataSet; import org.apache.iotdb.db.query.dataset.DirectAlignByTimeDataSet; import org.apache.iotdb.db.query.dataset.DirectNonAlignDataSet; import org.apache.iotdb.db.query.expression.ResultColumn; @@ -174,10 +173,12 @@ public class TSServiceImpl implements TSIService.Iface { private static final String INFO_QUERY_PROCESS_ERROR = "Error occurred in query process: "; private static final String INFO_NOT_ALLOWED_IN_BATCH_ERROR = "The query statement is not allowed in batch: "; - private static final String INFO_INTERRUPT_ERROR = "Current Thread interrupted when dealing with request {}"; + public static final TSProtocolVersion CURRENT_RPC_VERSION = + TSProtocolVersion.IOTDB_SERVICE_PROTOCOL_V3; + private static final int MAX_SIZE = IoTDBDescriptor.getInstance().getConfig().getQueryCacheSizeInMetric(); private static final int DELETE_SIZE = 20; @@ -185,23 +186,18 @@ public class TSServiceImpl implements TSIService.Iface { private static final long MS_TO_MONTH = 30 * 86400_000L; private final IoTDBConfig config = IoTDBDescriptor.getInstance().getConfig(); - private final boolean enableMetric = config.isEnableMetricService(); - - private static final List<SqlArgument> sqlArgumentList = new ArrayList<>(MAX_SIZE); - protected Planner processor; - protected IPlanExecutor executor; - - private SessionManager sessionManager = SessionManager.getInstance(); // When the client abnormally exits, we can still know who to disconnect private final ThreadLocal<Long> currSessionId = new ThreadLocal<>(); + private final SessionManager sessionManager = SessionManager.getInstance(); - public static final TSProtocolVersion CURRENT_RPC_VERSION = - TSProtocolVersion.IOTDB_SERVICE_PROTOCOL_V3; - + private static final List<SqlArgument> sqlArgumentList = new ArrayList<>(MAX_SIZE); private static final AtomicInteger queryCount = new AtomicInteger(0); + private final QueryTimeManager queryTimeManager = QueryTimeManager.getInstance(); + private final TracingManager tracingManager = TracingManager.getInstance(); - private QueryTimeManager queryTimeManager = QueryTimeManager.getInstance(); + protected Planner processor; + protected IPlanExecutor executor; public TSServiceImpl() throws QueryProcessException { processor = new Planner(); @@ -733,33 +729,24 @@ public class TSServiceImpl implements TSIService.Iface { QueryFilterOptimizationException, MetadataException, IOException, InterruptedException, TException, AuthException { queryCount.incrementAndGet(); - AUDIT_LOGGER.debug("Session {} execute Query: {}", currSessionId.get(), statement); + AUDIT_LOGGER.debug("Session {} execute query: {}", currSessionId.get(), statement); + long startTime = System.currentTimeMillis(); long queryId = -1; - try { - - // pair.left = fetchSize, pair.right = deduplicatedNum - Pair<Integer, Integer> p = getMemoryParametersFromPhysicalPlan(plan, fetchSize); - fetchSize = p.left; + try { + Pair<Integer, Integer> fetchSizeDeduplicatedPathNumPair = + getMemoryParametersFromPhysicalPlan(plan, fetchSize); + fetchSize = fetchSizeDeduplicatedPathNumPair.left; // generate the queryId for the operation - queryId = sessionManager.requestQueryId(statementId, true, fetchSize, p.right); - // register query info to queryTimeManager - if (!(plan instanceof ShowQueryProcesslistPlan)) { - queryTimeManager.registerQuery(queryId, startTime, statement, timeout); - } - if (plan instanceof QueryPlan && config.isEnablePerformanceTracing()) { - TracingManager tracingManager = TracingManager.getInstance(); - if (!(plan instanceof AlignByDevicePlan)) { - tracingManager.writeQueryInfo(queryId, statement, startTime, plan.getPaths().size()); - } else { - tracingManager.writeQueryInfo(queryId, statement, startTime); - } - } + queryId = + sessionManager.requestQueryId( + statementId, true, fetchSize, fetchSizeDeduplicatedPathNumPair.right); - if (plan instanceof AuthorPlan) { - plan.setLoginUserName(username); - } + queryTimeManager.registerQuery(queryId, startTime, statement, timeout, plan); + tracingManager.writeQueryInfo(queryId, statement, startTime, plan); + + plan.setLoginUserName(username); TSExecuteStatementResp resp = null; // execute it before createDataSet since it may change the content of query plan @@ -823,14 +810,12 @@ public class TSServiceImpl implements TSIService.Iface { } } } - resp.setQueryId(queryId); - if (plan instanceof AlignByDevicePlan && config.isEnablePerformanceTracing()) { - TracingManager.getInstance() - .writePathsNum(queryId, ((AlignByDeviceDataSet) newDataSet).getPathsNum()); - } + resp.setQueryId(queryId); - if (enableMetric) { + tracingManager.writePathsNum(queryId, newDataSet); + queryTimeManager.unRegisterQuery(queryId, plan); + if (config.isEnableMetricService()) { long endTime = System.currentTimeMillis(); SqlArgument sqlArgument = new SqlArgument(resp, plan, statement, startTime, endTime); synchronized (sqlArgumentList) { @@ -841,10 +826,6 @@ public class TSServiceImpl implements TSIService.Iface { } } - // remove query info in QueryTimeManager - if (!(plan instanceof ShowQueryProcesslistPlan)) { - queryTimeManager.unRegisterQuery(queryId); - } return resp; } catch (Exception e) { sessionManager.releaseQueryResourceNoExceptions(queryId); diff --git a/server/src/test/java/org/apache/iotdb/db/query/control/TracingManagerTest.java b/server/src/test/java/org/apache/iotdb/db/query/control/TracingManagerTest.java index 3f1f0a9..acb40e6 100644 --- a/server/src/test/java/org/apache/iotdb/db/query/control/TracingManagerTest.java +++ b/server/src/test/java/org/apache/iotdb/db/query/control/TracingManagerTest.java @@ -23,6 +23,9 @@ import org.apache.iotdb.db.conf.IoTDBDescriptor; import org.apache.iotdb.db.constant.TestConstant; import org.apache.iotdb.db.engine.fileSystem.SystemFileFactory; import org.apache.iotdb.db.engine.storagegroup.TsFileResource; +import org.apache.iotdb.db.exception.metadata.IllegalPathException; +import org.apache.iotdb.db.qp.physical.crud.AlignByDevicePlan; +import org.apache.iotdb.db.query.dataset.AlignByDeviceDataSet; import org.apache.iotdb.db.utils.EnvironmentUtils; import org.apache.commons.io.FileUtils; @@ -63,17 +66,17 @@ public class TracingManagerTest { } @Test - public void tracingQueryTest() throws IOException { + public void tracingQueryTest() throws IOException, IllegalPathException { if (!tracingManager.getWriterStatus()) { tracingManager.openTracingWriteStream(); } String[] ans = { "Query Id: 10 - Query Statement: " + sql, "Query Id: 10 - Start time: 2020-12-", - "Query Id: 10 - Number of series paths: 3", + "Query Id: 10 - Number of series paths: 0", "Query Id: 10 - Query Statement: " + sql, "Query Id: 10 - Start time: 2020-12-", - "Query Id: 10 - Number of series paths: 3", + "Query Id: 10 - Number of series paths: 0", "Query Id: 10 - Number of sequence files: 1", "Query Id: 10 - SeqFile_1-1-0.tsfile root.sg.d1[1, 999], root.sg.d2[2, 998]", "Query Id: 10 - Number of unSequence files: 0", @@ -81,9 +84,15 @@ public class TracingManagerTest { "Query Id: 10 - Average size of chunks: 1371", "Query Id: 10 - Total cost time: " }; + + AlignByDevicePlan plan = new AlignByDevicePlan(); + plan.setPaths(Collections.emptyList()); + plan.setDevices(Collections.emptyList()); + AlignByDeviceDataSet dataSet = new AlignByDeviceDataSet(plan, null, null); + tracingManager.writeQueryInfo(queryId, sql, 1607529600000L); - tracingManager.writePathsNum(queryId, 3); - tracingManager.writeQueryInfo(queryId, sql, 1607529600000L, 3); + tracingManager.writePathsNum(queryId, dataSet); + tracingManager.writeQueryInfo(queryId, sql, 1607529600000L, plan); tracingManager.writeTsFileInfo(queryId, seqResources, Collections.EMPTY_SET); tracingManager.writeChunksInfo(queryId, 3, 4113L); tracingManager.writeEndTime(queryId);
