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

qiaojialin pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iotdb.git


The following commit(s) were added to refs/heads/master by this push:
     new a93728770f IOTDB-3481 Optimize Log Print (#6273)
a93728770f is described below

commit a93728770fea3fdac340c143f1a915a0c9a5570b
Author: Jackie Tien <[email protected]>
AuthorDate: Tue Jun 14 14:49:00 2022 +0800

    IOTDB-3481 Optimize Log Print (#6273)
---
 .../execution/datatransfer/DataBlockManager.java   | 108 +++++++----------
 .../execution/datatransfer/LocalSinkHandle.java    |   8 +-
 .../execution/datatransfer/LocalSourceHandle.java  |  46 ++++---
 .../db/mpp/execution/datatransfer/SinkHandle.java  |  36 ++----
 .../mpp/execution/datatransfer/SourceHandle.java   | 134 ++++++++++-----------
 .../execution/schedule/AbstractDriverThread.java   |  13 +-
 .../db/mpp/execution/schedule/DriverScheduler.java |   4 +-
 .../mpp/execution/schedule/DriverTaskThread.java   |  71 ++++++-----
 .../db/mpp/plan/analyze/ClusterSchemaFetcher.java  |  36 +++---
 .../db/mpp/plan/execution/QueryExecution.java      |  24 ++--
 .../db/mpp/plan/scheduler/ClusterScheduler.java    |  12 +-
 .../thrift/impl/DataNodeTSIServiceImpl.java        |  41 ++++---
 12 files changed, 252 insertions(+), 281 deletions(-)

diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/execution/datatransfer/DataBlockManager.java
 
b/server/src/main/java/org/apache/iotdb/db/mpp/execution/datatransfer/DataBlockManager.java
index e6bb1bd3c0..05748f7519 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/execution/datatransfer/DataBlockManager.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/execution/datatransfer/DataBlockManager.java
@@ -78,7 +78,7 @@ public class DataBlockManager implements IDataBlockManager {
     @Override
     public TGetDataBlockResponse getDataBlock(TGetDataBlockRequest req) throws 
TException {
       try (SetThreadName fragmentInstanceName =
-          new SetThreadName(createFullIdFrom(req.sourceFragmentInstanceId))) {
+          new SetThreadName(createFullIdFrom(req.sourceFragmentInstanceId, 
"SinkHandle"))) {
         logger.debug(
             "Get data block request received, for data blocks whose sequence 
ID in [{}, {}) from {}.",
             req.getStartSequenceId(),
@@ -107,7 +107,7 @@ public class DataBlockManager implements IDataBlockManager {
     @Override
     public void onAcknowledgeDataBlockEvent(TAcknowledgeDataBlockEvent e) 
throws TException {
       try (SetThreadName fragmentInstanceName =
-          new SetThreadName(createFullIdFrom(e.sourceFragmentInstanceId))) {
+          new SetThreadName(createFullIdFrom(e.sourceFragmentInstanceId, 
"SinkHandle"))) {
         logger.debug(
             "Acknowledge data block event received, for data blocks whose 
sequence ID in [{}, {}) from {}.",
             e.getStartSequenceId(),
@@ -127,7 +127,8 @@ public class DataBlockManager implements IDataBlockManager {
     @Override
     public void onNewDataBlockEvent(TNewDataBlockEvent e) throws TException {
       try (SetThreadName fragmentInstanceName =
-          new SetThreadName(createFullIdFrom(e.sourceFragmentInstanceId))) {
+          new SetThreadName(
+              createFullIdFrom(e.targetFragmentInstanceId, e.targetPlanNodeId 
+ ".SourceHandle"))) {
         logger.debug(
             "New data block event received, for plan node {} of {} from {}.",
             e.getTargetPlanNodeId(),
@@ -161,7 +162,8 @@ public class DataBlockManager implements IDataBlockManager {
     @Override
     public void onEndOfDataBlockEvent(TEndOfDataBlockEvent e) throws 
TException {
       try (SetThreadName fragmentInstanceName =
-          new SetThreadName(createFullIdFrom(e.sourceFragmentInstanceId))) {
+          new SetThreadName(
+              createFullIdFrom(e.targetFragmentInstanceId, e.targetPlanNodeId 
+ ".SourceHandle"))) {
         logger.debug(
             "End of data block event received, for plan node {} of {} from 
{}.",
             e.getTargetPlanNodeId(),
@@ -201,43 +203,34 @@ public class DataBlockManager implements 
IDataBlockManager {
 
     @Override
     public void onFinished(ISourceHandle sourceHandle) {
-      try (SetThreadName fragmentInstanceName =
-          new 
SetThreadName(createFullIdFrom(sourceHandle.getLocalFragmentInstanceId()))) {
-        logger.info("{} finished and release resources", sourceHandle);
-        if 
(!sourceHandles.containsKey(sourceHandle.getLocalFragmentInstanceId())
-            || !sourceHandles
-                .get(sourceHandle.getLocalFragmentInstanceId())
-                .containsKey(sourceHandle.getLocalPlanNodeId())) {
-          logger.info("{} resources has already been released", sourceHandle);
-        } else {
-          sourceHandles
+      logger.info("finished and release resources");
+      if (!sourceHandles.containsKey(sourceHandle.getLocalFragmentInstanceId())
+          || !sourceHandles
               .get(sourceHandle.getLocalFragmentInstanceId())
-              .remove(sourceHandle.getLocalPlanNodeId());
-        }
-        if 
(sourceHandles.containsKey(sourceHandle.getLocalFragmentInstanceId())
-            && 
sourceHandles.get(sourceHandle.getLocalFragmentInstanceId()).isEmpty()) {
-          sourceHandles.remove(sourceHandle.getLocalFragmentInstanceId());
-        }
+              .containsKey(sourceHandle.getLocalPlanNodeId())) {
+        logger.info("resources has already been released");
+      } else {
+        sourceHandles
+            .get(sourceHandle.getLocalFragmentInstanceId())
+            .remove(sourceHandle.getLocalPlanNodeId());
+      }
+      if (sourceHandles.containsKey(sourceHandle.getLocalFragmentInstanceId())
+          && 
sourceHandles.get(sourceHandle.getLocalFragmentInstanceId()).isEmpty()) {
+        sourceHandles.remove(sourceHandle.getLocalFragmentInstanceId());
       }
     }
 
     @Override
     public void onAborted(ISourceHandle sourceHandle) {
-      try (SetThreadName fragmentInstanceName =
-          new 
SetThreadName(createFullIdFrom(sourceHandle.getLocalFragmentInstanceId()))) {
-        logger.info("{}: onAborted is invoked", sourceHandle);
-        onFinished(sourceHandle);
-      }
+      logger.info("onAborted is invoked");
+      onFinished(sourceHandle);
     }
 
     @Override
     public void onFailure(ISourceHandle sourceHandle, Throwable t) {
-      try (SetThreadName fragmentInstanceName =
-          new 
SetThreadName(createFullIdFrom(sourceHandle.getLocalFragmentInstanceId()))) {
-        logger.error("Source handle {} failed due to {}", sourceHandle, t);
-        if (onFailureCallback != null) {
-          onFailureCallback.call(t);
-        }
+      logger.error("Source handle failed due to: ", t);
+      if (onFailureCallback != null) {
+        onFailureCallback.call(t);
       }
     }
   }
@@ -256,50 +249,35 @@ public class DataBlockManager implements 
IDataBlockManager {
 
     @Override
     public void onFinish(ISinkHandle sinkHandle) {
-      try (SetThreadName fragmentInstanceName =
-          new 
SetThreadName(createFullIdFrom(sinkHandle.getLocalFragmentInstanceId()))) {
-        removeFromDataBlockManager(sinkHandle);
-        context.finished();
-      }
+      removeFromDataBlockManager(sinkHandle);
+      context.finished();
     }
 
     @Override
     public void onEndOfBlocks(ISinkHandle sinkHandle) {
-      try (SetThreadName fragmentInstanceName =
-          new 
SetThreadName(createFullIdFrom(sinkHandle.getLocalFragmentInstanceId()))) {
-        context.transitionToFlushing();
-      }
+      context.transitionToFlushing();
     }
 
     @Override
     public void onAborted(ISinkHandle sinkHandle) {
-      try (SetThreadName fragmentInstanceName =
-          new 
SetThreadName(createFullIdFrom(sinkHandle.getLocalFragmentInstanceId()))) {
-        logger.info("{} onAborted is invoked", sinkHandle);
-        removeFromDataBlockManager(sinkHandle);
-      }
+      logger.info("onAborted is invoked");
+      removeFromDataBlockManager(sinkHandle);
     }
 
     private void removeFromDataBlockManager(ISinkHandle sinkHandle) {
-      try (SetThreadName fragmentInstanceName =
-          new 
SetThreadName(createFullIdFrom(sinkHandle.getLocalFragmentInstanceId()))) {
-        logger.info("{} release resources of finished sink handle", 
sinkHandle);
-        if (!sinkHandles.containsKey(sinkHandle.getLocalFragmentInstanceId())) 
{
-          logger.info("{} resources already been released", sinkHandle);
-        }
-        sinkHandles.remove(sinkHandle.getLocalFragmentInstanceId());
+      logger.info("{} release resources of finished sink handle", sinkHandle);
+      if (!sinkHandles.containsKey(sinkHandle.getLocalFragmentInstanceId())) {
+        logger.info("{} resources already been released", sinkHandle);
       }
+      sinkHandles.remove(sinkHandle.getLocalFragmentInstanceId());
     }
 
     @Override
     public void onFailure(ISinkHandle sinkHandle, Throwable t) {
-      try (SetThreadName fragmentInstanceName =
-          new 
SetThreadName(createFullIdFrom(sinkHandle.getLocalFragmentInstanceId()))) {
-        // TODO: (xingtanzjr) should we remove the sinkHandle from 
DataBlockManager ?
-        logger.error("Sink handle {} failed due to {}", sinkHandle, t);
-        if (onFailureCallback != null) {
-          onFailureCallback.call(t);
-        }
+      // TODO: (xingtanzjr) should we remove the sinkHandle from 
DataBlockManager ?
+      logger.error("Sink handle {} failed due to {}", sinkHandle, t);
+      if (onFailureCallback != null) {
+        onFailureCallback.call(t);
       }
     }
   }
@@ -497,10 +475,9 @@ public class DataBlockManager implements IDataBlockManager 
{
    * <p>This method should be called when a fragment instance finished in an 
abnormal state.
    */
   public void forceDeregisterFragmentInstance(TFragmentInstanceId 
fragmentInstanceId) {
-    logger.info("Force deregister fragment instance {}", fragmentInstanceId);
+    logger.info("Force deregister fragment instance");
     if (sinkHandles.containsKey(fragmentInstanceId)) {
       ISinkHandle sinkHandle = sinkHandles.get(fragmentInstanceId);
-      logger.info("Abort sink handle {}", sinkHandle);
       sinkHandle.abort();
       sinkHandles.remove(fragmentInstanceId);
     }
@@ -514,8 +491,13 @@ public class DataBlockManager implements IDataBlockManager 
{
     }
   }
 
-  public static String createFullIdFrom(TFragmentInstanceId 
fragmentInstanceId) {
+  /** @param suffix should be like [PlanNodeId].SourceHandle/SinHandle */
+  public static String createFullIdFrom(TFragmentInstanceId 
fragmentInstanceId, String suffix) {
     return createFullId(
-        fragmentInstanceId.queryId, fragmentInstanceId.fragmentId, 
fragmentInstanceId.instanceId);
+            fragmentInstanceId.queryId,
+            fragmentInstanceId.fragmentId,
+            fragmentInstanceId.instanceId)
+        + "."
+        + suffix;
   }
 }
diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/execution/datatransfer/LocalSinkHandle.java
 
b/server/src/main/java/org/apache/iotdb/db/mpp/execution/datatransfer/LocalSinkHandle.java
index 9f51017261..d1a8ac33f1 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/execution/datatransfer/LocalSinkHandle.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/execution/datatransfer/LocalSinkHandle.java
@@ -111,7 +111,7 @@ public class LocalSinkHandle implements ISinkHandle {
 
   @Override
   public synchronized void setNoMoreTsBlocks() {
-    logger.info("Set no-more-tsblocks to {}.", this);
+    logger.info("Set no-more-tsblocks.");
     if (aborted) {
       return;
     }
@@ -120,16 +120,16 @@ public class LocalSinkHandle implements ISinkHandle {
     if (isFinished()) {
       sinkHandleListener.onFinish(this);
     }
-    logger.info("No-more-tsblocks has been set to {}.", this);
+    logger.info("No-more-tsblocks has been set.");
   }
 
   @Override
   public synchronized void abort() {
-    logger.info("Sink handle {} is being aborted.", this);
+    logger.info("Sink handle is being aborted.");
     aborted = true;
     queue.destroy();
     sinkHandleListener.onAborted(this);
-    logger.info("Sink handle {} is aborted", this);
+    logger.info("Sink handle is aborted");
   }
 
   public TFragmentInstanceId getRemoteFragmentInstanceId() {
diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/execution/datatransfer/LocalSourceHandle.java
 
b/server/src/main/java/org/apache/iotdb/db/mpp/execution/datatransfer/LocalSourceHandle.java
index e4691b715a..14a0ebc664 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/execution/datatransfer/LocalSourceHandle.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/execution/datatransfer/LocalSourceHandle.java
@@ -24,11 +24,13 @@ import org.apache.iotdb.mpp.rpc.thrift.TFragmentInstanceId;
 import org.apache.iotdb.tsfile.read.common.block.TsBlock;
 
 import com.google.common.util.concurrent.ListenableFuture;
+import io.airlift.concurrent.SetThreadName;
 import org.apache.commons.lang3.Validate;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
 import static 
com.google.common.util.concurrent.Futures.nonCancellationPropagating;
+import static 
org.apache.iotdb.db.mpp.execution.datatransfer.DataBlockManager.createFullIdFrom;
 
 public class LocalSourceHandle implements ISourceHandle {
 
@@ -41,6 +43,8 @@ public class LocalSourceHandle implements ISourceHandle {
   private final SharedTsBlockQueue queue;
   private boolean aborted = false;
 
+  private final String threadName;
+
   public LocalSourceHandle(
       TFragmentInstanceId remoteFragmentInstanceId,
       TFragmentInstanceId localFragmentInstanceId,
@@ -52,6 +56,8 @@ public class LocalSourceHandle implements ISourceHandle {
     this.localPlanNodeId = Validate.notNull(localPlanNodeId);
     this.queue = Validate.notNull(queue);
     this.sourceHandleListener = Validate.notNull(sourceHandleListener);
+    this.threadName =
+        createFullIdFrom(localFragmentInstanceId, localPlanNodeId + "." + 
"SourceHandle");
   }
 
   @Override
@@ -71,20 +77,22 @@ public class LocalSourceHandle implements ISourceHandle {
 
   @Override
   public TsBlock receive() {
-    if (aborted) {
-      throw new IllegalStateException("Source handle is aborted.");
-    }
-    if (!queue.isBlocked().isDone()) {
-      throw new IllegalStateException("Source handle is blocked.");
+    try (SetThreadName sourceHandleName = new SetThreadName(threadName)) {
+      if (aborted) {
+        throw new IllegalStateException("Source handle is aborted.");
+      }
+      if (!queue.isBlocked().isDone()) {
+        throw new IllegalStateException("Source handle is blocked.");
+      }
+      TsBlock tsBlock;
+      synchronized (this) {
+        tsBlock = queue.remove();
+      }
+      if (isFinished()) {
+        sourceHandleListener.onFinished(this);
+      }
+      return tsBlock;
     }
-    TsBlock tsBlock;
-    synchronized (this) {
-      tsBlock = queue.remove();
-    }
-    if (isFinished()) {
-      sourceHandleListener.onFinished(this);
-    }
-    return tsBlock;
   }
 
   @Override
@@ -107,12 +115,14 @@ public class LocalSourceHandle implements ISourceHandle {
 
   @Override
   public synchronized void abort() {
-    if (aborted) {
-      return;
+    try (SetThreadName sourceHandleName = new SetThreadName(threadName)) {
+      if (aborted) {
+        return;
+      }
+      queue.destroy();
+      aborted = true;
+      sourceHandleListener.onAborted(this);
     }
-    queue.destroy();
-    aborted = true;
-    sourceHandleListener.onAborted(this);
   }
 
   public TFragmentInstanceId getRemoteFragmentInstanceId() {
diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/execution/datatransfer/SinkHandle.java
 
b/server/src/main/java/org/apache/iotdb/db/mpp/execution/datatransfer/SinkHandle.java
index 7342b93d72..b7a131614d 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/execution/datatransfer/SinkHandle.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/execution/datatransfer/SinkHandle.java
@@ -65,6 +65,7 @@ public class SinkHandle implements ISinkHandle {
   private final ExecutorService executorService;
   private final TsBlockSerde serde;
   private final SinkHandleListener sinkHandleListener;
+  private final String threadName;
   private long retryIntervalInMs;
 
   // Use LinkedHashMap to meet 2 needs,
@@ -103,6 +104,7 @@ public class SinkHandle implements ISinkHandle {
     this.sinkHandleListener = Validate.notNull(sinkHandleListener);
     this.dataBlockServiceClientManager = dataBlockServiceClientManager;
     this.retryIntervalInMs = DEFAULT_RETRY_INTERVAL_IN_MS;
+    this.threadName = createFullIdFrom(localFragmentInstanceId, "SinkHandle");
   }
 
   @Override
@@ -159,7 +161,7 @@ public class SinkHandle implements ISinkHandle {
   }
 
   private void sendEndOfDataBlockEvent() throws Exception {
-    logger.info("{} send end of data block event", this);
+    logger.info("send end of data block event");
     int attempt = 0;
     TEndOfDataBlockEvent endOfDataBlockEvent =
         new TEndOfDataBlockEvent(
@@ -174,12 +176,7 @@ public class SinkHandle implements ISinkHandle {
         client.onEndOfDataBlockEvent(endOfDataBlockEvent);
         break;
       } catch (Throwable e) {
-        logger.error(
-            "{} Failed to send end of data block event due to {}, attempt 
times: {}",
-            this,
-            e.getMessage(),
-            attempt,
-            e);
+        logger.error("Failed to send end of data block event, attempt times: 
{}", attempt, e);
         if (attempt == MAX_ATTEMPT_TIMES) {
           throw e;
         }
@@ -190,7 +187,7 @@ public class SinkHandle implements ISinkHandle {
 
   @Override
   public synchronized void setNoMoreTsBlocks() {
-    logger.info("{} start to set no-more-tsblocks", this);
+    logger.info("start to set no-more-tsblocks");
     if (aborted) {
       return;
     }
@@ -199,19 +196,19 @@ public class SinkHandle implements ISinkHandle {
     } catch (Exception e) {
       throw new RuntimeException("Send EndOfDataBlockEvent failed", e);
     }
-    logger.info("{} set noMoreTsBlocks to true", this);
+    logger.info("set noMoreTsBlocks to true");
     noMoreTsBlocks = true;
     if (isFinished()) {
-      logger.info("{} revoke onFinish() of sinkHandleListener", this);
+      logger.info("revoke onFinish() of sinkHandleListener");
       sinkHandleListener.onFinish(this);
     }
-    logger.info("{} revoke onEndOfBlocks() of sinkHandleListener", this);
+    logger.info("revoke onEndOfBlocks() of sinkHandleListener");
     sinkHandleListener.onEndOfBlocks(this);
   }
 
   @Override
   public synchronized void abort() {
-    logger.info("{} is being aborted.", this);
+    logger.info("SinkHandle is being aborted.");
     sequenceIdToTsBlock.clear();
     aborted = true;
     bufferRetainedSizeInBytes -= 
localMemoryManager.getQueryPool().tryCancel(blocked);
@@ -222,7 +219,7 @@ public class SinkHandle implements ISinkHandle {
       bufferRetainedSizeInBytes = 0;
     }
     sinkHandleListener.onAborted(this);
-    logger.info("{} is aborted", this);
+    logger.info("SinkHandle is aborted");
   }
 
   @Override
@@ -334,11 +331,9 @@ public class SinkHandle implements ISinkHandle {
 
     @Override
     public void run() {
-      try (SetThreadName fragmentInstanceName =
-          new SetThreadName(createFullIdFrom(getLocalFragmentInstanceId()))) {
+      try (SetThreadName sinkHandleName = new SetThreadName(threadName)) {
         logger.info(
-            "{} send new data block event [{}, {})",
-            SinkHandle.this,
+            "Send new data block event [{}, {})",
             startSequenceId,
             startSequenceId + blockSizes.size());
         int attempt = 0;
@@ -356,12 +351,7 @@ public class SinkHandle implements ISinkHandle {
             client.onNewDataBlockEvent(newDataBlockEvent);
             break;
           } catch (Throwable e) {
-            logger.error(
-                "{} failed to send new data block event due to {}, attempt 
times: {}",
-                SinkHandle.this,
-                e.getMessage(),
-                attempt,
-                e);
+            logger.error("Failed to send new data block event, attempt times: 
{}", attempt, e);
             if (attempt == MAX_ATTEMPT_TIMES) {
               sinkHandleListener.onFailure(SinkHandle.this, e);
             }
diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/execution/datatransfer/SourceHandle.java
 
b/server/src/main/java/org/apache/iotdb/db/mpp/execution/datatransfer/SourceHandle.java
index 60d1c9c85e..a27244e2c8 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/execution/datatransfer/SourceHandle.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/execution/datatransfer/SourceHandle.java
@@ -67,6 +67,8 @@ public class SourceHandle implements ISourceHandle {
 
   private final Map<Integer, TsBlock> sequenceIdToTsBlock = new HashMap<>();
   private final Map<Integer, Long> sequenceIdToDataBlockSize = new HashMap<>();
+
+  private final String threadName;
   private long retryIntervalInMs;
 
   private final IClientManager<TEndPoint, SyncDataNodeDataBlockServiceClient>
@@ -102,44 +104,48 @@ public class SourceHandle implements ISourceHandle {
     this.executorService = Validate.notNull(executorService);
     this.serde = Validate.notNull(serde);
     this.sourceHandleListener = Validate.notNull(sourceHandleListener);
-    bufferRetainedSizeInBytes = 0L;
+    this.bufferRetainedSizeInBytes = 0L;
     this.dataBlockServiceClientManager = dataBlockServiceClientManager;
     this.retryIntervalInMs = DEFAULT_RETRY_INTERVAL_IN_MS;
+    this.threadName =
+        createFullIdFrom(localFragmentInstanceId, localPlanNodeId + "." + 
"SourceHandle");
   }
 
   @Override
   public synchronized TsBlock receive() {
-    if (aborted) {
-      throw new IllegalStateException("Source handle is aborted.");
-    }
-    if (!blocked.isDone()) {
-      throw new IllegalStateException("Source handle is blocked.");
-    }
+    try (SetThreadName sourceHandleName = new SetThreadName(threadName)) {
 
-    TsBlock tsBlock;
-    tsBlock = sequenceIdToTsBlock.remove(currSequenceId);
-    logger.info(
-        "Receive {} TsdBlock, size is {}", currSequenceId, 
tsBlock.getRetainedSizeInBytes());
-    currSequenceId += 1;
-    bufferRetainedSizeInBytes -= tsBlock.getRetainedSizeInBytes();
-    localMemoryManager
-        .getQueryPool()
-        .free(localFragmentInstanceId.getQueryId(), 
tsBlock.getRetainedSizeInBytes());
-
-    if (sequenceIdToTsBlock.isEmpty() && !isFinished()) {
-      logger.info("{}: no buffered TsBlock, blocked", this);
-      blocked = SettableFuture.create();
-    }
-    if (isFinished()) {
-      sourceHandleListener.onFinished(this);
+      if (aborted) {
+        throw new IllegalStateException("Source handle is aborted.");
+      }
+      if (!blocked.isDone()) {
+        throw new IllegalStateException("Source handle is blocked.");
+      }
+
+      TsBlock tsBlock;
+      tsBlock = sequenceIdToTsBlock.remove(currSequenceId);
+      logger.info(
+          "Receive {} TsdBlock, size is {}", currSequenceId, 
tsBlock.getRetainedSizeInBytes());
+      currSequenceId += 1;
+      bufferRetainedSizeInBytes -= tsBlock.getRetainedSizeInBytes();
+      localMemoryManager
+          .getQueryPool()
+          .free(localFragmentInstanceId.getQueryId(), 
tsBlock.getRetainedSizeInBytes());
+
+      if (sequenceIdToTsBlock.isEmpty() && !isFinished()) {
+        logger.info("no buffered TsBlock, blocked");
+        blocked = SettableFuture.create();
+      }
+      if (isFinished()) {
+        sourceHandleListener.onFinished(this);
+      }
+      trySubmitGetDataBlocksTask();
+      return tsBlock;
     }
-    trySubmitGetDataBlocksTask();
-    return tsBlock;
   }
 
   private synchronized void trySubmitGetDataBlocksTask() {
-    try (SetThreadName fragmentInstanceName =
-        new SetThreadName(createFullIdFrom(getLocalFragmentInstanceId()))) {
+    try (SetThreadName sourceHandleName = new SetThreadName(threadName)) {
       if (aborted) {
         return;
       }
@@ -200,7 +206,7 @@ public class SourceHandle implements ISourceHandle {
   }
 
   synchronized void setNoMoreTsBlocks(int lastSequenceId) {
-    logger.info("{}: receive NoMoreTsBlock event. ", this);
+    logger.info("receive NoMoreTsBlock event. ");
     this.lastSequenceId = lastSequenceId;
     if (!blocked.isDone() && remoteTsBlockedConsumedUp()) {
       blocked.set(null);
@@ -212,8 +218,7 @@ public class SourceHandle implements ISourceHandle {
 
   synchronized void updatePendingDataBlockInfo(int startSequenceId, List<Long> 
dataBlockSizes) {
     logger.info(
-        "{}: receive newDataBlockEvent. [{}, {}), each size is: {}",
-        this,
+        "receive newDataBlockEvent. [{}, {}), each size is: {}",
         startSequenceId,
         startSequenceId + dataBlockSizes.size(),
         dataBlockSizes);
@@ -225,24 +230,26 @@ public class SourceHandle implements ISourceHandle {
 
   @Override
   public synchronized void abort() {
-    if (aborted) {
-      return;
-    }
-    if (blocked != null && !blocked.isDone()) {
-      blocked.cancel(true);
-    }
-    if (blockedOnMemory != null) {
-      bufferRetainedSizeInBytes -= 
localMemoryManager.getQueryPool().tryCancel(blockedOnMemory);
-    }
-    sequenceIdToDataBlockSize.clear();
-    if (bufferRetainedSizeInBytes > 0) {
-      localMemoryManager
-          .getQueryPool()
-          .free(localFragmentInstanceId.getQueryId(), 
bufferRetainedSizeInBytes);
-      bufferRetainedSizeInBytes = 0;
+    try (SetThreadName sourceHandleName = new SetThreadName(threadName)) {
+      if (aborted) {
+        return;
+      }
+      if (blocked != null && !blocked.isDone()) {
+        blocked.cancel(true);
+      }
+      if (blockedOnMemory != null) {
+        bufferRetainedSizeInBytes -= 
localMemoryManager.getQueryPool().tryCancel(blockedOnMemory);
+      }
+      sequenceIdToDataBlockSize.clear();
+      if (bufferRetainedSizeInBytes > 0) {
+        localMemoryManager
+            .getQueryPool()
+            .free(localFragmentInstanceId.getQueryId(), 
bufferRetainedSizeInBytes);
+        bufferRetainedSizeInBytes = 0;
+      }
+      aborted = true;
+      sourceHandleListener.onAborted(this);
     }
-    aborted = true;
-    sourceHandleListener.onAborted(this);
   }
 
   @Override
@@ -323,13 +330,8 @@ public class SourceHandle implements ISourceHandle {
 
     @Override
     public void run() {
-      try (SetThreadName fragmentInstanceName =
-          new SetThreadName(createFullIdFrom(getLocalFragmentInstanceId()))) {
-        logger.info(
-            "{}: try to get data blocks [{}, {}) ",
-            SourceHandle.this,
-            startSequenceId,
-            endSequenceId);
+      try (SetThreadName sourceHandleName = new SetThreadName(threadName)) {
+        logger.info("try to get data blocks [{}, {}) ", startSequenceId, 
endSequenceId);
         TGetDataBlockRequest req =
             new TGetDataBlockRequest(remoteFragmentInstanceId, 
startSequenceId, endSequenceId);
         int attempt = 0;
@@ -343,7 +345,7 @@ public class SourceHandle implements ISourceHandle {
               TsBlock tsBlock = serde.deserialize(byteBuffer);
               tsBlocks.add(tsBlock);
             }
-            logger.info("{}: got data blocks. count: {}", SourceHandle.this, 
tsBlocks.size());
+            logger.info("got data blocks. count: {}", tsBlocks.size());
             executorService.submit(
                 new SendAcknowledgeDataBlockEventTask(startSequenceId, 
endSequenceId));
             synchronized (SourceHandle.this) {
@@ -359,11 +361,7 @@ public class SourceHandle implements ISourceHandle {
             }
             break;
           } catch (Throwable e) {
-            logger.error(
-                "{}: failed to get data block {}, attempt times: {}",
-                SourceHandle.this,
-                e.getMessage(),
-                attempt);
+            logger.error("failed to get data block, attempt times: {}", 
attempt, e);
             if (attempt == MAX_ATTEMPT_TIMES) {
               synchronized (SourceHandle.this) {
                 bufferRetainedSizeInBytes -= reservedBytes;
@@ -399,13 +397,8 @@ public class SourceHandle implements ISourceHandle {
 
     @Override
     public void run() {
-      try (SetThreadName fragmentInstanceName =
-          new SetThreadName(createFullIdFrom(getLocalFragmentInstanceId()))) {
-        logger.info(
-            "{}: send ack data block event [{}, {}).",
-            SourceHandle.this,
-            startSequenceId,
-            endSequenceId);
+      try (SetThreadName sourceHandleName = new SetThreadName(threadName)) {
+        logger.info("send ack data block event [{}, {}).", startSequenceId, 
endSequenceId);
         int attempt = 0;
         TAcknowledgeDataBlockEvent acknowledgeDataBlockEvent =
             new TAcknowledgeDataBlockEvent(
@@ -418,12 +411,11 @@ public class SourceHandle implements ISourceHandle {
             break;
           } catch (Throwable e) {
             logger.error(
-                "{}: failed to send ack data block event [{}, {}) due to {}, 
attempt times: {}",
-                SourceHandle.this,
+                "failed to send ack data block event [{}, {}), attempt times: 
{}",
                 startSequenceId,
                 endSequenceId,
-                e.getMessage(),
-                attempt);
+                attempt,
+                e);
             if (attempt == MAX_ATTEMPT_TIMES) {
               synchronized (SourceHandle.this) {
                 sourceHandleListener.onFailure(SourceHandle.this, e);
diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/execution/schedule/AbstractDriverThread.java
 
b/server/src/main/java/org/apache/iotdb/db/mpp/execution/schedule/AbstractDriverThread.java
index 8dea0529e2..c5fc410f4c 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/execution/schedule/AbstractDriverThread.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/execution/schedule/AbstractDriverThread.java
@@ -21,6 +21,7 @@ package org.apache.iotdb.db.mpp.execution.schedule;
 import org.apache.iotdb.db.mpp.execution.schedule.queue.IndexedBlockingQueue;
 import org.apache.iotdb.db.mpp.execution.schedule.task.DriverTask;
 
+import io.airlift.concurrent.SetThreadName;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
@@ -51,14 +52,18 @@ public abstract class AbstractDriverThread extends Thread 
implements Closeable {
   public void run() {
     DriverTask next;
     while (!closed && !Thread.currentThread().isInterrupted()) {
-      next = null;
       try {
         next = queue.poll();
-        execute(next);
       } catch (InterruptedException e) {
+        logger.error("Executor " + this.getName() + "failed to poll driver 
task from queue");
+        Thread.currentThread().interrupt();
         break;
-      } catch (Exception e) {
-        logger.error("Executor " + this.getName() + " processes failed", e);
+      }
+      try (SetThreadName fragmentInstanceName =
+          new SetThreadName(next.getFragmentInstance().getInfo().getFullId())) 
{
+        execute(next);
+      } catch (Throwable t) {
+        logger.error("execute failed", t);
         if (next != null) {
           
next.setAbortCause(FragmentInstanceAbortedException.BY_INTERNAL_ERROR_SCHEDULED);
           scheduler.toAborted(next);
diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/execution/schedule/DriverScheduler.java
 
b/server/src/main/java/org/apache/iotdb/db/mpp/execution/schedule/DriverScheduler.java
index f58ccee753..25675f04af 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/execution/schedule/DriverScheduler.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/execution/schedule/DriverScheduler.java
@@ -204,7 +204,7 @@ public class DriverScheduler implements IDriverScheduler, 
IService {
                   new FragmentInstanceAbortedException(
                       task.getFragmentInstance().getInfo(), 
task.getAbortCause()));
         } catch (Exception e) {
-          logger.error("Clear DriverTask {} failed", task.getId().toString(), 
e);
+          logger.error("Clear DriverTask failed", e);
         }
       }
       if (task.getStatus() == DriverTaskStatus.ABORTED) {
@@ -215,7 +215,7 @@ public class DriverScheduler implements IDriverScheduler, 
IService {
                   task.getId().getFragmentId().getId(),
                   task.getId().getInstanceId()));
         } catch (Exception e) {
-          logger.error("Clear DriverTask {} failed", task.getId().toString(), 
e);
+          logger.error("Clear DriverTask failed", e);
         }
       }
     }
diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/execution/schedule/DriverTaskThread.java
 
b/server/src/main/java/org/apache/iotdb/db/mpp/execution/schedule/DriverTaskThread.java
index e3b8909895..b013377a13 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/execution/schedule/DriverTaskThread.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/execution/schedule/DriverTaskThread.java
@@ -49,44 +49,41 @@ public class DriverTaskThread extends AbstractDriverThread {
 
   @Override
   public void execute(DriverTask task) throws InterruptedException {
-    try (SetThreadName fragmentInstanceName =
-        new SetThreadName(task.getFragmentInstance().getInfo().getFullId())) {
-      // try to switch it to RUNNING
-      if (!scheduler.readyToRunning(task)) {
-        return;
-      }
-      IDriver instance = task.getFragmentInstance();
-      CpuTimer timer = new CpuTimer();
-      ListenableFuture<Void> future = 
instance.processFor(EXECUTION_TIME_SLICE);
-      CpuTimer.CpuDuration duration = timer.elapsedTime();
-      // long cost = System.nanoTime() - startTime;
-      // If the future is cancelled, the task is in an error and should be 
thrown.
-      if (future.isCancelled()) {
-        
task.setAbortCause(FragmentInstanceAbortedException.BY_ALREADY_BEING_CANCELLED);
-        scheduler.toAborted(task);
-        return;
-      }
-      ExecutionContext context = new ExecutionContext();
-      context.setCpuDuration(duration);
-      context.setTimeSlice(EXECUTION_TIME_SLICE);
-      if (instance.isFinished()) {
-        scheduler.runningToFinished(task, context);
-        return;
-      }
+    // try to switch it to RUNNING
+    if (!scheduler.readyToRunning(task)) {
+      return;
+    }
+    IDriver instance = task.getFragmentInstance();
+    CpuTimer timer = new CpuTimer();
+    ListenableFuture<Void> future = instance.processFor(EXECUTION_TIME_SLICE);
+    CpuTimer.CpuDuration duration = timer.elapsedTime();
+    // long cost = System.nanoTime() - startTime;
+    // If the future is cancelled, the task is in an error and should be 
thrown.
+    if (future.isCancelled()) {
+      
task.setAbortCause(FragmentInstanceAbortedException.BY_ALREADY_BEING_CANCELLED);
+      scheduler.toAborted(task);
+      return;
+    }
+    ExecutionContext context = new ExecutionContext();
+    context.setCpuDuration(duration);
+    context.setTimeSlice(EXECUTION_TIME_SLICE);
+    if (instance.isFinished()) {
+      scheduler.runningToFinished(task, context);
+      return;
+    }
 
-      if (future.isDone()) {
-        scheduler.runningToReady(task, context);
-      } else {
-        scheduler.runningToBlocked(task, context);
-        future.addListener(
-            () -> {
-              try (SetThreadName fragmentInstanceName2 =
-                  new 
SetThreadName(task.getFragmentInstance().getInfo().getFullId())) {
-                scheduler.blockedToReady(task);
-              }
-            },
-            listeningExecutor);
-      }
+    if (future.isDone()) {
+      scheduler.runningToReady(task, context);
+    } else {
+      scheduler.runningToBlocked(task, context);
+      future.addListener(
+          () -> {
+            try (SetThreadName fragmentInstanceName2 =
+                new 
SetThreadName(task.getFragmentInstance().getInfo().getFullId())) {
+              scheduler.blockedToReady(task);
+            }
+          },
+          listeningExecutor);
     }
   }
 }
diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/analyze/ClusterSchemaFetcher.java
 
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/analyze/ClusterSchemaFetcher.java
index cf7c6bb9af..3ed55d9995 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/analyze/ClusterSchemaFetcher.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/analyze/ClusterSchemaFetcher.java
@@ -47,6 +47,8 @@ import org.apache.iotdb.tsfile.utils.Binary;
 import org.apache.iotdb.tsfile.utils.Pair;
 import org.apache.iotdb.tsfile.write.schema.MeasurementSchema;
 
+import io.airlift.concurrent.SetThreadName;
+
 import java.nio.ByteBuffer;
 import java.util.ArrayList;
 import java.util.Arrays;
@@ -106,24 +108,26 @@ public class ClusterSchemaFetcher implements 
ISchemaFetcher {
               "cannot fetch schema, status is: %s, msg is: %s",
               executionResult.status.getCode(), 
executionResult.status.getMessage()));
     }
-    SchemaTree result = new SchemaTree();
-    while (coordinator.getQueryExecution(queryId).hasNextResult()) {
-      // The query will be transited to FINISHED when invoking 
getBatchResult() at the last time
-      // So we don't need to clean up it manually
-      Optional<TsBlock> tsBlock = 
coordinator.getQueryExecution(queryId).getBatchResult();
-      if (!tsBlock.isPresent() || tsBlock.get().isEmpty()) {
-        break;
-      }
-      Binary binary;
-      SchemaTree fetchedSchemaTree;
-      Column column = tsBlock.get().getColumn(0);
-      for (int i = 0; i < column.getPositionCount(); i++) {
-        binary = column.getBinary(i);
-        fetchedSchemaTree = 
SchemaTree.deserialize(ByteBuffer.wrap(binary.getValues()));
-        result.mergeSchemaTree(fetchedSchemaTree);
+    try (SetThreadName threadName = new 
SetThreadName(executionResult.queryId.getId())) {
+      SchemaTree result = new SchemaTree();
+      while (coordinator.getQueryExecution(queryId).hasNextResult()) {
+        // The query will be transited to FINISHED when invoking 
getBatchResult() at the last time
+        // So we don't need to clean up it manually
+        Optional<TsBlock> tsBlock = 
coordinator.getQueryExecution(queryId).getBatchResult();
+        if (!tsBlock.isPresent() || tsBlock.get().isEmpty()) {
+          break;
+        }
+        Binary binary;
+        SchemaTree fetchedSchemaTree;
+        Column column = tsBlock.get().getColumn(0);
+        for (int i = 0; i < column.getPositionCount(); i++) {
+          binary = column.getBinary(i);
+          fetchedSchemaTree = 
SchemaTree.deserialize(ByteBuffer.wrap(binary.getValues()));
+          result.mergeSchemaTree(fetchedSchemaTree);
+        }
       }
+      return result;
     }
-    return result;
   }
 
   @Override
diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/execution/QueryExecution.java
 
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/execution/QueryExecution.java
index 9453ce73a1..1462033971 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/execution/QueryExecution.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/execution/QueryExecution.java
@@ -85,7 +85,7 @@ import static 
org.apache.iotdb.db.mpp.plan.constant.DataNodeEndPoints.isSameNode
 public class QueryExecution implements IQueryExecution {
   private static final Logger logger = 
LoggerFactory.getLogger(QueryExecution.class);
 
-  private static IoTDBConfig config = 
IoTDBDescriptor.getInstance().getConfig();
+  private static final IoTDBConfig config = 
IoTDBDescriptor.getInstance().getConfig();
 
   private final MPPQueryContext context;
   private IScheduler scheduler;
@@ -155,8 +155,7 @@ public class QueryExecution implements IQueryExecution {
 
   public void start() {
     if (skipExecute()) {
-      logger.info(
-          "{} execution of query will be skipped. Transit to RUNNING 
immediately.", getLogHeader());
+      logger.info("execution of query will be skipped. Transit to RUNNING 
immediately.");
       constructResultForMemorySource();
       stateMachine.transitionToRunning();
       return;
@@ -189,7 +188,7 @@ public class QueryExecution implements IQueryExecution {
       IPartitionFetcher partitionFetcher,
       ISchemaFetcher schemaFetcher) {
     // initialize the variable `analysis`
-    logger.info("{} start to analyze query", getLogHeader());
+    logger.info("start to analyze query");
     return new Analyzer(context, partitionFetcher, 
schemaFetcher).analyze(statement);
   }
 
@@ -219,23 +218,20 @@ public class QueryExecution implements IQueryExecution {
 
   // Use LogicalPlanner to do the logical query plan and logical optimization
   public void doLogicalPlan() {
-    logger.info("{} do logical plan...", getLogHeader());
+    logger.info("do logical plan...");
     LogicalPlanner planner = new LogicalPlanner(this.context, 
this.planOptimizers);
     this.logicalPlan = planner.plan(this.analysis);
     logger.info(
-        "{} logical plan is: \n {}",
-        getLogHeader(),
-        PlanNodeUtil.nodeToString(this.logicalPlan.getRootNode()));
+        "logical plan is: \n {}", 
PlanNodeUtil.nodeToString(this.logicalPlan.getRootNode()));
   }
 
   // Generate the distributed plan and split it into fragments
   public void doDistributedPlan() {
-    logger.info("{} do distribution plan...", getLogHeader());
+    logger.info("do distribution plan...");
     DistributionPlanner planner = new DistributionPlanner(this.analysis, 
this.logicalPlan);
     this.distributedPlan = planner.planFragments();
     logger.info(
-        "{} distribution plan done. Fragment instance count is {}, details is: 
\n {}",
-        getLogHeader(),
+        "distribution plan done. Fragment instance count is {}, details is: \n 
{}",
         distributedPlan.getInstances().size(),
         distributedPlan.getInstances());
   }
@@ -280,7 +276,7 @@ public class QueryExecution implements IQueryExecution {
       if (resultHandle == null || resultHandle.isAborted() || 
resultHandle.isFinished()) {
         // Once the resultHandle is finished, we should transit the state of 
this query to FINISHED.
         // So that the corresponding cleanup work could be triggered.
-        logger.info("{} resultHandle for client is finished", getLogHeader());
+        logger.info("resultHandle for client is finished");
         stateMachine.transitionToFinished();
         return Optional.empty();
       }
@@ -439,8 +435,4 @@ public class QueryExecution implements IQueryExecution {
   public String toString() {
     return String.format("QueryExecution[%s]", context.getQueryId());
   }
-
-  private String getLogHeader() {
-    return String.format("Query[%s]:", context.getQueryId());
-  }
 }
diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/scheduler/ClusterScheduler.java
 
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/scheduler/ClusterScheduler.java
index dd2ca46271..398f0f8272 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/scheduler/ClusterScheduler.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/scheduler/ClusterScheduler.java
@@ -94,7 +94,7 @@ public class ClusterScheduler implements IScheduler {
   @Override
   public void start() {
     stateMachine.transitionToDispatching();
-    logger.info("{} transit to DISPATCHING", getLogHeader());
+    logger.info("transit to DISPATCHING");
     Future<FragInstanceDispatchResult> dispatchResultFuture = 
dispatcher.dispatch(instances);
 
     // NOTICE: the FragmentInstance may be dispatched to another Host due to 
consensus redirect.
@@ -102,7 +102,7 @@ public class ClusterScheduler implements IScheduler {
     try {
       FragInstanceDispatchResult result = dispatchResultFuture.get();
       if (!result.isSuccessful()) {
-        logger.error("{} dispatch failed.", getLogHeader());
+        logger.error("dispatch failed.");
         stateMachine.transitionToFailed(new IllegalStateException("Fragment 
cannot be dispatched"));
         return;
       }
@@ -124,7 +124,7 @@ public class ClusterScheduler implements IScheduler {
     // The FragmentInstances has been dispatched successfully to corresponding 
host, we mark the
     // QueryState to Running
     stateMachine.transitionToRunning();
-    logger.info("{} transit to RUNNING", getLogHeader());
+    logger.info("transit to RUNNING");
     instances.forEach(
         instance -> {
           stateMachine.initialFragInstanceState(instance.getId(), 
FragmentInstanceState.RUNNING);
@@ -132,7 +132,7 @@ public class ClusterScheduler implements IScheduler {
 
     // TODO: (xingtanzjr) start the stateFetcher/heartbeat for each fragment 
instance
     this.stateTracker.start();
-    logger.info("{} state tracker starts", getLogHeader());
+    logger.info("state tracker starts");
   }
 
   @Override
@@ -166,8 +166,4 @@ public class ClusterScheduler implements IScheduler {
 
   // After sending, start to collect the states of these fragment instances
   private void startMonitorInstances() {}
-
-  private String getLogHeader() {
-    return String.format("Query[%s]", queryContext.getQueryId());
-  }
 }
diff --git 
a/server/src/main/java/org/apache/iotdb/db/service/thrift/impl/DataNodeTSIServiceImpl.java
 
b/server/src/main/java/org/apache/iotdb/db/service/thrift/impl/DataNodeTSIServiceImpl.java
index fc3bc9a303..6e5ac4349e 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/service/thrift/impl/DataNodeTSIServiceImpl.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/service/thrift/impl/DataNodeTSIServiceImpl.java
@@ -1077,17 +1077,18 @@ public class DataNodeTSIServiceImpl implements 
TSIEventHandler {
 
       IQueryExecution queryExecution = COORDINATOR.getQueryExecution(queryId);
 
-      TSExecuteStatementResp resp;
-      if (queryExecution.isQuery()) {
-        resp = createResponse(queryExecution.getDatasetHeader(), queryId);
-        resp.setStatus(result.status);
-        resp.setQueryDataSet(
-            QueryDataSetUtils.convertTsBlockByFetchSize(queryExecution, 
req.fetchSize));
-      } else {
-        resp = RpcUtils.getTSExecuteStatementResp(result.status);
+      try (SetThreadName threadName = new 
SetThreadName(result.queryId.getId())) {
+        TSExecuteStatementResp resp;
+        if (queryExecution.isQuery()) {
+          resp = createResponse(queryExecution.getDatasetHeader(), queryId);
+          resp.setStatus(result.status);
+          resp.setQueryDataSet(
+              QueryDataSetUtils.convertTsBlockByFetchSize(queryExecution, 
req.fetchSize));
+        } else {
+          resp = RpcUtils.getTSExecuteStatementResp(result.status);
+        }
+        return resp;
       }
-
-      return resp;
     } catch (Exception e) {
       // TODO call the coordinator to release query resource
       return RpcUtils.getTSExecuteStatementResp(
@@ -1136,17 +1137,19 @@ public class DataNodeTSIServiceImpl implements 
TSIEventHandler {
 
       IQueryExecution queryExecution = COORDINATOR.getQueryExecution(queryId);
 
-      TSExecuteStatementResp resp;
-      if (queryExecution.isQuery()) {
-        resp = createResponse(queryExecution.getDatasetHeader(), queryId);
-        resp.setStatus(result.status);
-        resp.setQueryDataSet(
-            QueryDataSetUtils.convertTsBlockByFetchSize(queryExecution, 
req.fetchSize));
-      } else {
-        resp = RpcUtils.getTSExecuteStatementResp(result.status);
+      try (SetThreadName threadName = new 
SetThreadName(result.queryId.getId())) {
+        TSExecuteStatementResp resp;
+        if (queryExecution.isQuery()) {
+          resp = createResponse(queryExecution.getDatasetHeader(), queryId);
+          resp.setStatus(result.status);
+          resp.setQueryDataSet(
+              QueryDataSetUtils.convertTsBlockByFetchSize(queryExecution, 
req.fetchSize));
+        } else {
+          resp = RpcUtils.getTSExecuteStatementResp(result.status);
+        }
+        return resp;
       }
 
-      return resp;
     } catch (Exception e) {
       // TODO call the coordinator to release query resource
       return RpcUtils.getTSExecuteStatementResp(

Reply via email to