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