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

Reply via email to