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();

Reply via email to