This is an automated email from the ASF dual-hosted git repository. jackietien pushed a commit to branch OptimizeLogPrint in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit 0bfaede0bc99046fa49bf9cda2f983915aaf76ea Author: JackieTien97 <[email protected]> AuthorDate: Tue Jun 14 12:09:46 2022 +0800 IOTDB-3481 Optimize Log Print --- .../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(
