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

Reply via email to