This is an automated email from the ASF dual-hosted git repository.

rong pushed a commit to branch load-parallel-region-load
in repository https://gitbox.apache.org/repos/asf/iotdb.git

commit 6c7c110c6198d413af00ee72abf4a96436f8124f
Author: Steve Yurong Su <[email protected]>
AuthorDate: Wed Oct 23 19:15:34 2024 +0800

    Load: Parallelly load files into different target data partitions
---
 .../db/storageengine/load/LoadTsFileManager.java   | 57 +++++++++++++++-------
 1 file changed, 40 insertions(+), 17 deletions(-)

diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/LoadTsFileManager.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/LoadTsFileManager.java
index 75c82a9a1d3..b19bf841190 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/LoadTsFileManager.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/LoadTsFileManager.java
@@ -61,8 +61,11 @@ import java.io.IOException;
 import java.nio.file.DirectoryNotEmptyException;
 import java.nio.file.Files;
 import java.nio.file.Path;
+import java.util.ArrayList;
 import java.util.Arrays;
+import java.util.Collections;
 import java.util.HashMap;
+import java.util.List;
 import java.util.Map;
 import java.util.Objects;
 import java.util.Optional;
@@ -452,27 +455,47 @@ public class LoadTsFileManager {
       if (isClosed) {
         throw new 
IOException(String.format(MESSAGE_WRITER_MANAGER_HAS_BEEN_CLOSED, taskDir));
       }
+
       for (Map.Entry<DataPartitionInfo, ModificationFile> entry :
           dataPartition2ModificationFile.entrySet()) {
         entry.getValue().close();
       }
-      for (Map.Entry<DataPartitionInfo, TsFileIOWriter> entry : 
dataPartition2Writer.entrySet()) {
-        TsFileIOWriter writer = entry.getValue();
-        if (writer.isWritingChunkGroup()) {
-          writer.endChunkGroup();
-        }
-        writer.endFile();
-
-        DataRegion dataRegion = entry.getKey().getDataRegion();
-        dataRegion.loadNewTsFile(generateResource(writer, progressIndex), 
true, isGeneratedByPipe);
-
-        // Metrics
-        dataRegion
-            .getNonSystemDatabaseName()
-            .ifPresent(
-                databaseName ->
-                    updateWritePointCountMetrics(
-                        dataRegion, databaseName, 
getTsFileWritePointCount(writer), false));
+
+      final List<Map.Entry<DataPartitionInfo, TsFileIOWriter>> 
dataPartition2WriterList =
+          new ArrayList<>(dataPartition2Writer.entrySet());
+      Collections.shuffle(dataPartition2WriterList);
+
+      final AtomicReference<Exception> exception = new AtomicReference<>();
+      dataPartition2WriterList.parallelStream()
+          .forEach(
+              entry -> {
+                try {
+                  final TsFileIOWriter writer = entry.getValue();
+                  if (writer.isWritingChunkGroup()) {
+                    writer.endChunkGroup();
+                  }
+                  writer.endFile();
+
+                  final DataRegion dataRegion = entry.getKey().getDataRegion();
+                  dataRegion.loadNewTsFile(
+                      generateResource(writer, progressIndex), true, 
isGeneratedByPipe);
+
+                  // Metrics
+                  dataRegion
+                      .getNonSystemDatabaseName()
+                      .ifPresent(
+                          databaseName ->
+                              updateWritePointCountMetrics(
+                                  dataRegion,
+                                  databaseName,
+                                  getTsFileWritePointCount(writer),
+                                  false));
+                } catch (final Exception e) {
+                  exception.set(e);
+                }
+              });
+      if (exception.get() != null) {
+        throw new LoadFileException(exception.get());
       }
     }
 

Reply via email to