This is an automated email from the ASF dual-hosted git repository. xingtanzjr pushed a commit to branch new_geely_car_0205 in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit 0d054430622b04047f4152ba6c3f638fb26bbc0d Merge: 3228b83105 ec55c45e4f Author: Jinrui.Zhang <[email protected]> AuthorDate: Mon Feb 6 18:14:20 2023 +0800 Merge branch 'lock' into new_geely_car_0205 client-cpp/src/main/Session.cpp | 2 +- .../iotdb/confignode/conf/ConfigNodeConfig.java | 11 ++ .../confignode/conf/ConfigNodeDescriptor.java | 8 + .../iotdb/confignode/manager/ConfigManager.java | 144 ++++++++++++-- .../iotdb/commons/concurrent/ThreadName.java | 5 +- .../threadpool/WrappedThreadPoolExecutor.java | 5 +- .../exception/CompactionExceptionHandler.java | 4 +- .../execute/task/CrossSpaceCompactionTask.java | 6 +- .../execute/task/InnerSpaceCompactionTask.java | 4 +- .../compaction/execute/utils/CompactionUtils.java | 2 +- .../compaction/schedule/CompactionScheduler.java | 8 +- .../db/engine/settle/SettleRequestHandler.java | 4 +- .../iotdb/db/engine/storagegroup/DataRegion.java | 40 ++-- .../engine/storagegroup/HashLastFlushTimeMap.java | 2 +- .../storagegroup/IDTableLastFlushTimeMap.java | 2 +- .../db/engine/storagegroup/TsFileManager.java | 12 +- .../db/metadata/mtree/MTreeBelowSGCachedImpl.java | 65 +++++-- .../db/metadata/mtree/MTreeBelowSGMemoryImpl.java | 67 +++++-- .../db/metadata/mtree/store/CachedMTreeStore.java | 88 ++------- .../mtree/store/disk/MTreeFlushTaskManager.java | 71 ------- .../mtree/store/disk/MTreeReleaseTaskManager.java | 73 -------- .../mtree/store/disk/cache/CacheManager.java | 8 +- .../mtree/store/disk/cache/CacheMemoryManager.java | 207 +++++++++++++++++++++ .../db/metadata/rescon/SchemaResourceManager.java | 10 +- .../db/metadata/schemaregion/ISchemaRegion.java | 4 + .../scheduler/FragmentInstanceDispatcherImpl.java | 30 ++- .../cross/CrossSpaceCompactionExceptionTest.java | 32 ++-- ...eCompactionWithFastPerformerValidationTest.java | 12 +- .../inner/InnerSpaceCompactionExceptionTest.java | 9 +- .../SizeTieredCompactionSelectorTest.java | 4 +- .../schemaRegion/SchemaRegionTestUtil.java | 12 +- 31 files changed, 611 insertions(+), 340 deletions(-) diff --cc server/src/main/java/org/apache/iotdb/db/mpp/plan/scheduler/FragmentInstanceDispatcherImpl.java index 71c55a592d,f299d862e9..a50c00184d --- a/server/src/main/java/org/apache/iotdb/db/mpp/plan/scheduler/FragmentInstanceDispatcherImpl.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/plan/scheduler/FragmentInstanceDispatcherImpl.java @@@ -130,62 -125,35 +131,79 @@@ public class FragmentInstanceDispatcher try (SetThreadName threadName = new SetThreadName(instance.getId().getFullId())) { dispatchOneInstance(instance); } catch (FragmentInstanceDispatchException e) { - return immediateFuture(new FragInstanceDispatchResult(e.getFailureStatus())); + TSStatus failureStatus = e.getFailureStatus(); + if (instances.size() == 1) { + failureStatusList.add(failureStatus); + } else { + if (failureStatus.getCode() == TSStatusCode.MULTIPLE_ERROR.getStatusCode()) { + failureStatusList.addAll(failureStatus.getSubStatus()); + } else { + failureStatusList.add(failureStatus); + } + } } catch (Throwable t) { - logger.warn("[DispatchFailed]", t); + // logger.warn("[DispatchFailed]", t); + failureStatusList.add( + RpcUtils.getStatus( + TSStatusCode.INTERNAL_SERVER_ERROR, "Unexpected errors: " + t.getMessage())); + } + } + if (failureStatusList.isEmpty()) { + return immediateFuture(new FragInstanceDispatchResult(true)); + } else { + if (instances.size() == 1) { + return immediateFuture(new FragInstanceDispatchResult(failureStatusList.get(0))); + } else { return immediateFuture( - new FragInstanceDispatchResult( - RpcUtils.getStatus( - TSStatusCode.INTERNAL_SERVER_ERROR, "Unexpected errors: " + t.getMessage()))); + new FragInstanceDispatchResult(RpcUtils.getStatus(failureStatusList))); } } - return immediateFuture(new FragInstanceDispatchResult(true)); } + private Future<FragInstanceDispatchResult> dispatchWriteAsync(List<FragmentInstance> instances) { + // split local and remote instances + List<FragmentInstance> localInstances = new ArrayList<>(); + List<FragmentInstance> remoteInstances = new ArrayList<>(); + for (FragmentInstance instance : instances) { + TEndPoint endPoint = instance.getHostDataNode().getInternalEndPoint(); + if (isDispatchedToLocal(endPoint)) { + localInstances.add(instance); + } else { + remoteInstances.add(instance); + } + } + // async dispatch to remote + AsyncPlanNodeSender asyncPlanNodeSender = + new AsyncPlanNodeSender(asyncInternalServiceClientManager, remoteInstances); + asyncPlanNodeSender.sendAll(); + // sync dispatch to local + for (FragmentInstance localInstance : localInstances) { + try (SetThreadName threadName = new SetThreadName(localInstance.getId().getFullId())) { + dispatchOneInstance(localInstance); + } catch (FragmentInstanceDispatchException e) { + return immediateFuture(new FragInstanceDispatchResult(e.getFailureStatus())); + } catch (Throwable t) { - logger.warn("[DispatchFailed]", t); ++ // logger.warn("[DispatchFailed]", t); + return immediateFuture( + new FragInstanceDispatchResult( + RpcUtils.getStatus( + TSStatusCode.INTERNAL_SERVER_ERROR, "Unexpected errors: " + t.getMessage()))); + } + } + // wait until remote dispatch done + try { + asyncPlanNodeSender.waitUntilCompleted(); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + logger.error("Interrupted when dispatching write async", e); + return immediateFuture( + new FragInstanceDispatchResult( + RpcUtils.getStatus( + TSStatusCode.INTERNAL_SERVER_ERROR, "Interrupted errors: " + e.getMessage()))); + } + return asyncPlanNodeSender.getResult(); + } + private void dispatchOneInstance(FragmentInstance instance) throws FragmentInstanceDispatchException { TEndPoint endPoint = instance.getHostDataNode().getInternalEndPoint();
