This is an automated email from the ASF dual-hosted git repository.
haonan 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 78954ea fix tsfilemanage concurrent bug (#1849)
78954ea is described below
commit 78954ea5f96d6d2cf1210ad8b4fc03374173db61
Author: zhanglingzhe0820 <[email protected]>
AuthorDate: Mon Oct 26 09:48:17 2020 +0800
fix tsfilemanage concurrent bug (#1849)
---
.../level/LevelTsFileManagement.java | 49 ++++++++++++----------
1 file changed, 27 insertions(+), 22 deletions(-)
diff --git
a/server/src/main/java/org/apache/iotdb/db/engine/tsfilemanagement/level/LevelTsFileManagement.java
b/server/src/main/java/org/apache/iotdb/db/engine/tsfilemanagement/level/LevelTsFileManagement.java
index cc42102..9e95c63 100644
---
a/server/src/main/java/org/apache/iotdb/db/engine/tsfilemanagement/level/LevelTsFileManagement.java
+++
b/server/src/main/java/org/apache/iotdb/db/engine/tsfilemanagement/level/LevelTsFileManagement.java
@@ -37,6 +37,7 @@ import java.util.HashSet;
import java.util.Iterator;
import java.util.List;
import java.util.Map;
+import java.util.Map.Entry;
import java.util.Set;
import java.util.TreeSet;
import java.util.concurrent.ConcurrentSkipListMap;
@@ -437,7 +438,8 @@ public class LevelTsFileManagement extends TsFileManagement
{
Pair<Double, Double> unSeqStatisticsPair = forkTsFileList(
forkedUnSequenceTsFileResources,
unSequenceTsFileResources
- .computeIfAbsent(timePartition,
this::newUnSequenceTsFileResources), maxUnseqLevelNum);
+ .computeIfAbsent(timePartition,
this::newUnSequenceTsFileResources),
+ maxUnseqLevelNum);
forkedUnSeqListPointNum = unSeqStatisticsPair.left;
forkedUnSeqListMeasurementSize = unSeqStatisticsPair.right;
}
@@ -454,43 +456,46 @@ public class LevelTsFileManagement extends
TsFileManagement {
List<TsFileResource> forkedLevelTsFileResources = new ArrayList<>();
Collection<TsFileResource> levelRawTsFileResources =
(Collection<TsFileResource>) rawTsFileResources
.get(i);
- synchronized (levelRawTsFileResources) {
- for (TsFileResource tsFileResource : levelRawTsFileResources) {
- if (tsFileResource.isClosed()) {
- String path = tsFileResource.getTsFile().getAbsolutePath();
- if (tsFileResource.getTsFile().exists()) {
- try (TsFileSequenceReader reader = new
TsFileSequenceReader(path)) {
- List<Path> pathList = reader.getAllPaths();
- for (Path sensorPath : pathList) {
+ List<TsFileResource> allCurrLevelTsFileResources = new
ArrayList<>(levelRawTsFileResources);
+ for (TsFileResource tsFileResource : allCurrLevelTsFileResources) {
+ if (tsFileResource.isClosed()) {
+ String path = tsFileResource.getTsFile().getAbsolutePath();
+ if (tsFileResource.getTsFile().exists()) {
+ try (TsFileSequenceReader reader = new TsFileSequenceReader(path))
{
+ List<String> devices = reader.getAllDevices();
+ for (String device : devices) {
+ Map<String, List<ChunkMetadata>>
measurementChunkMetadataListMap = reader
+ .readChunkMetadataInDevice(device);
+ for (Entry<String, List<ChunkMetadata>>
measurementChunkMetadataList : measurementChunkMetadataListMap
+ .entrySet()) {
+ Path sensorPath = new Path(device,
measurementChunkMetadataList.getKey());
measurementSet.offer(sensorPath.getFullPath());
List<ChunkMetadata> chunkMetadataList =
reader.getChunkMetadataList(sensorPath);
for (ChunkMetadata chunkMetadata : chunkMetadataList) {
pointNum += chunkMetadata.getNumOfPoints();
}
}
- } catch (IOException e) {
- logger.error(
- "{} tsfile reader creates error", path, e);
}
- } else {
- logger.info("{} tsfile does not exist", path);
+ } catch (IOException e) {
+ logger.error(
+ "{} tsfile reader creates error", path, e);
}
+ } else {
+ logger.info("{} tsfile does not exist", path);
}
- if (measurementSet.cardinality() > 0
- && pointNum / measurementSet.cardinality() >= maxChunkPointNum) {
- forkedLevelTsFileResources.add(tsFileResource);
- break;
- }
- forkedLevelTsFileResources.add(tsFileResource);
+ }
+ forkedLevelTsFileResources.add(tsFileResource);
+ if (measurementSet.cardinality() > 0
+ && pointNum / measurementSet.cardinality() >= maxChunkPointNum) {
+ break;
}
}
+ forkedTsFileResources.add(forkedLevelTsFileResources);
if (measurementSet.cardinality() > 0
&& pointNum / measurementSet.cardinality() >= maxChunkPointNum) {
- forkedTsFileResources.add(forkedLevelTsFileResources);
break;
}
- forkedTsFileResources.add(forkedLevelTsFileResources);
}
// fill in empty file