This is an automated email from the ASF dual-hosted git repository. xingtanzjr pushed a commit to branch xingtanzjr/mpp_issues in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit 946a69fd516b24f7578e63171fb4487890879fd3 Author: Jinrui.Zhang <[email protected]> AuthorDate: Tue Apr 19 15:13:30 2022 +0800 spotless --- .../org/apache/iotdb/db/mpp/buffer/DataBlockManager.java | 3 ++- .../java/org/apache/iotdb/db/mpp/buffer/SourceHandle.java | 15 ++++++++++----- .../org/apache/iotdb/db/mpp/execution/QueryExecution.java | 13 ++++++++++--- .../db/mpp/execution/scheduler/ClusterScheduler.java | 5 ++++- .../execution/scheduler/SimpleFragInstanceDispatcher.java | 1 + .../iotdb/db/mpp/operator/source/SeriesScanOperator.java | 4 +++- .../java/org/apache/iotdb/db/utils/QueryDataSetUtils.java | 2 +- 7 files changed, 31 insertions(+), 12 deletions(-) diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/buffer/DataBlockManager.java b/server/src/main/java/org/apache/iotdb/db/mpp/buffer/DataBlockManager.java index b5bf3f65a0..573771fe81 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/buffer/DataBlockManager.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/buffer/DataBlockManager.java @@ -180,7 +180,8 @@ public class DataBlockManager implements IDataBlockManager { .get(sourceHandle.getLocalFragmentInstanceId()) .remove(sourceHandle.getLocalPlanNodeId()); } - if (sourceHandles.containsKey(sourceHandle.getLocalFragmentInstanceId()) && sourceHandles.get(sourceHandle.getLocalFragmentInstanceId()).isEmpty()) { + if (sourceHandles.containsKey(sourceHandle.getLocalFragmentInstanceId()) + && sourceHandles.get(sourceHandle.getLocalFragmentInstanceId()).isEmpty()) { sourceHandles.remove(sourceHandle.getLocalFragmentInstanceId()); } } diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/buffer/SourceHandle.java b/server/src/main/java/org/apache/iotdb/db/mpp/buffer/SourceHandle.java index d044be172e..9a38876954 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/buffer/SourceHandle.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/buffer/SourceHandle.java @@ -211,13 +211,19 @@ public class SourceHandle implements ISourceHandle { } synchronized void setNoMoreTsBlocks(int lastSequenceId) { - logger.info("[SourceHandle {}]: No more TsBlock. {} ", localPlanNodeId, remoteFragmentInstanceId); + logger.info( + "[SourceHandle {}]: No more TsBlock. {} ", localPlanNodeId, remoteFragmentInstanceId); this.lastSequenceId = lastSequenceId; if (!blocked.isDone() && currSequenceId - 1 == lastSequenceId) { - logger.info("[SourceHandle {}]: all blocks are consumed. set blocked to null.", localPlanNodeId); + logger.info( + "[SourceHandle {}]: all blocks are consumed. set blocked to null.", localPlanNodeId); blocked.set(null); } else { - logger.info("[SourceHandle {}]: No need to set blocked. Blocked: {}, Consumed: {} ", localPlanNodeId, blocked.isDone(), currSequenceId - 1 == lastSequenceId); + logger.info( + "[SourceHandle {}]: No need to set blocked. Blocked: {}, Consumed: {} ", + localPlanNodeId, + blocked.isDone(), + currSequenceId - 1 == lastSequenceId); } } @@ -251,8 +257,7 @@ public class SourceHandle implements ISourceHandle { @Override public boolean isFinished() { - return throwable == null - && currSequenceId - 1 == lastSequenceId; + return throwable == null && currSequenceId - 1 == lastSequenceId; } String getRemoteHostname() { diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/execution/QueryExecution.java b/server/src/main/java/org/apache/iotdb/db/mpp/execution/QueryExecution.java index 380361c3fe..2b2ea1f762 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/execution/QueryExecution.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/execution/QueryExecution.java @@ -110,8 +110,11 @@ public class QueryExecution implements IQueryExecution { return; } this.stop(); - // TODO: (xingtanzjr) If the query is in abnormal state, the releaseResource() should be invoked - if (state == QueryState.FAILED || state == QueryState.ABORTED || state == QueryState.CANCELED) { + // TODO: (xingtanzjr) If the query is in abnormal state, the releaseResource() should be + // invoked + if (state == QueryState.FAILED + || state == QueryState.ABORTED + || state == QueryState.CANCELED) { releaseResource(); } }); @@ -206,7 +209,11 @@ public class QueryExecution implements IQueryExecution { LOG.info("[QueryExecution {}]: try to get result.", context.getQueryId()); ListenableFuture<Void> blocked = resultHandle.isBlocked(); blocked.get(); - LOG.info("[QueryExecution {}]: unblock. Cancelled: {}, Done: {}", context.getQueryId(), blocked.isCancelled(), blocked.isDone()); + LOG.info( + "[QueryExecution {}]: unblock. Cancelled: {}, Done: {}", + context.getQueryId(), + blocked.isCancelled(), + blocked.isDone()); if (resultHandle.isFinished()) { LOG.info("[QueryExecution {}]: result is null", context.getQueryId()); releaseResource(); diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/execution/scheduler/ClusterScheduler.java b/server/src/main/java/org/apache/iotdb/db/mpp/execution/scheduler/ClusterScheduler.java index 47c57c873d..694a82738f 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/execution/scheduler/ClusterScheduler.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/execution/scheduler/ClusterScheduler.java @@ -83,7 +83,10 @@ public class ClusterScheduler implements IScheduler { @Override public void start() { - LOGGER.info("[{}] start to dispatch fragment instance. size: {}", queryContext.getQueryId(), instances.size()); + LOGGER.info( + "[{}] start to dispatch fragment instance. size: {}", + queryContext.getQueryId(), + instances.size()); stateMachine.transitionToDispatching(); Future<FragInstanceDispatchResult> dispatchResultFuture = dispatcher.dispatch(instances); diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/execution/scheduler/SimpleFragInstanceDispatcher.java b/server/src/main/java/org/apache/iotdb/db/mpp/execution/scheduler/SimpleFragInstanceDispatcher.java index 40b4df4a63..8745dfe95b 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/execution/scheduler/SimpleFragInstanceDispatcher.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/execution/scheduler/SimpleFragInstanceDispatcher.java @@ -27,6 +27,7 @@ import org.apache.iotdb.mpp.rpc.thrift.TConsensusGroupId; import org.apache.iotdb.mpp.rpc.thrift.TFragmentInstance; import org.apache.iotdb.mpp.rpc.thrift.TSendFragmentInstanceReq; import org.apache.iotdb.mpp.rpc.thrift.TSendFragmentInstanceResp; + import org.slf4j.Logger; import org.slf4j.LoggerFactory; diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/operator/source/SeriesScanOperator.java b/server/src/main/java/org/apache/iotdb/db/mpp/operator/source/SeriesScanOperator.java index b11f643810..2927cf0d77 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/operator/source/SeriesScanOperator.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/operator/source/SeriesScanOperator.java @@ -107,7 +107,9 @@ public class SeriesScanOperator implements DataSourceOperator { return true; } } - System.out.println(String.format("[SeriesScanOperator-%s]: hasNext returned: %s", sourceId, hasCachedTsBlock)); + System.out.println( + String.format( + "[SeriesScanOperator-%s]: hasNext returned: %s", sourceId, hasCachedTsBlock)); return hasCachedTsBlock; } catch (IOException e) { throw new RuntimeException("Error happened while scanning the file", e); diff --git a/server/src/main/java/org/apache/iotdb/db/utils/QueryDataSetUtils.java b/server/src/main/java/org/apache/iotdb/db/utils/QueryDataSetUtils.java index aca0b37906..d255f03ddd 100644 --- a/server/src/main/java/org/apache/iotdb/db/utils/QueryDataSetUtils.java +++ b/server/src/main/java/org/apache/iotdb/db/utils/QueryDataSetUtils.java @@ -18,7 +18,6 @@ */ package org.apache.iotdb.db.utils; -import org.apache.iotdb.db.mpp.execution.Coordinator; import org.apache.iotdb.db.mpp.execution.IQueryExecution; import org.apache.iotdb.db.tools.watermark.WatermarkEncoder; import org.apache.iotdb.service.rpc.thrift.TSQueryDataSet; @@ -32,6 +31,7 @@ import org.apache.iotdb.tsfile.read.query.dataset.QueryDataSet; import org.apache.iotdb.tsfile.utils.Binary; import org.apache.iotdb.tsfile.utils.BitMap; import org.apache.iotdb.tsfile.utils.BytesUtils; + import org.slf4j.Logger; import org.slf4j.LoggerFactory;
