This is an automated email from the ASF dual-hosted git repository.
rong 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 dbb99bc88de Load: Parallelly load files into different target data
partitions (#13893)
dbb99bc88de is described below
commit dbb99bc88dea50c6effffae4081abbba1ea5f76a
Author: Steve Yurong Su <[email protected]>
AuthorDate: Thu Oct 24 12:29:42 2024 +0800
Load: Parallelly load files into different target data partitions (#13893)
---
.../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 775dad1df46..f785d0992eb 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());
}
}