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(