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

jackietien 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 22f2866325 [IOTDB-3777]Avoid holding sg write lock in the whole 
processin of deletion (#6743)
22f2866325 is described below

commit 22f28663259ebbdee74332e84dad1a3073ccc65a
Author: 周沛辰 <[email protected]>
AuthorDate: Tue Jul 26 16:39:11 2022 +0800

    [IOTDB-3777]Avoid holding sg write lock in the whole processin of deletion 
(#6743)
---
 .../iotdb/db/engine/storagegroup/DataRegion.java   | 154 +++++++++++++--------
 .../db/engine/storagegroup/DataRegionTest.java     |  85 +++++++++++-
 2 files changed, 178 insertions(+), 61 deletions(-)

diff --git 
a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/DataRegion.java 
b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/DataRegion.java
index b1543d71c3..3c08cb50cb 100755
--- 
a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/DataRegion.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/DataRegion.java
@@ -2101,7 +2101,7 @@ public class DataRegion {
 
     // record files which are updated so that we can roll back them in case of 
exception
     List<ModificationFile> updatedModFiles = new ArrayList<>();
-
+    boolean hasReleasedLock = false;
     try {
       Set<PartialPath> devicePaths = 
IoTDB.schemaProcessor.getBelongedDevices(path);
       for (PartialPath device : devicePaths) {
@@ -2122,18 +2122,18 @@ public class DataRegion {
 
       Deletion deletion = new Deletion(path, MERGE_MOD_START_VERSION_NUM, 
startTime, endTime);
 
+      List<TsFileResource> sealedTsFileResource = new ArrayList<>();
+      List<TsFileResource> unsealedTsFileResource = new ArrayList<>();
+      separateTsFile(sealedTsFileResource, unsealedTsFileResource);
+
       deleteDataInFiles(
-          tsFileManager.getTsFileList(true),
-          deletion,
-          devicePaths,
-          updatedModFiles,
-          timePartitionFilter);
+          unsealedTsFileResource, deletion, devicePaths, updatedModFiles, 
timePartitionFilter);
+
+      writeUnlock();
+      hasReleasedLock = true;
+
       deleteDataInFiles(
-          tsFileManager.getTsFileList(false),
-          deletion,
-          devicePaths,
-          updatedModFiles,
-          timePartitionFilter);
+          sealedTsFileResource, deletion, devicePaths, updatedModFiles, 
timePartitionFilter);
 
     } catch (Exception e) {
       // roll back
@@ -2144,10 +2144,37 @@ public class DataRegion {
       }
       throw new IOException(e);
     } finally {
-      writeUnlock();
+      if (!hasReleasedLock) {
+        writeUnlock();
+      }
     }
   }
 
+  /** Seperate tsfiles in TsFileManager to sealedList and unsealedList. */
+  private void separateTsFile(
+      List<TsFileResource> sealedResource, List<TsFileResource> 
unsealedResource) {
+    tsFileManager
+        .getTsFileList(true)
+        .forEach(
+            tsFileResource -> {
+              if (tsFileResource.isClosed()) {
+                sealedResource.add(tsFileResource);
+              } else {
+                unsealedResource.add(tsFileResource);
+              }
+            });
+    tsFileManager
+        .getTsFileList(false)
+        .forEach(
+            tsFileResource -> {
+              if (tsFileResource.isClosed()) {
+                sealedResource.add(tsFileResource);
+              } else {
+                unsealedResource.add(tsFileResource);
+              }
+            });
+  }
+
   /**
    * @param pattern Must be a pattern start with a precise device path
    * @param startTime
@@ -2179,6 +2206,7 @@ public class DataRegion {
 
     // record files which are updated so that we can roll back them in case of 
exception
     List<ModificationFile> updatedModFiles = new ArrayList<>();
+    boolean hasReleasedLock = false;
 
     try {
 
@@ -2202,18 +2230,17 @@ public class DataRegion {
 
       Deletion deletion = new Deletion(pattern, MERGE_MOD_START_VERSION_NUM, 
startTime, endTime);
 
+      List<TsFileResource> sealedTsFileResource = new ArrayList<>();
+      List<TsFileResource> unsealedTsFileResource = new ArrayList<>();
+      separateTsFile(sealedTsFileResource, unsealedTsFileResource);
+
       deleteDataInFiles(
-          tsFileManager.getTsFileList(true),
-          deletion,
-          devicePaths,
-          updatedModFiles,
-          timePartitionFilter);
+          unsealedTsFileResource, deletion, devicePaths, updatedModFiles, 
timePartitionFilter);
+      writeUnlock();
+      hasReleasedLock = true;
+
       deleteDataInFiles(
-          tsFileManager.getTsFileList(false),
-          deletion,
-          devicePaths,
-          updatedModFiles,
-          timePartitionFilter);
+          sealedTsFileResource, deletion, devicePaths, updatedModFiles, 
timePartitionFilter);
 
     } catch (Exception e) {
       // roll back
@@ -2224,7 +2251,9 @@ public class DataRegion {
       }
       throw new IOException(e);
     } finally {
-      writeUnlock();
+      if (!hasReleasedLock) {
+        writeUnlock();
+      }
     }
   }
 
@@ -2269,10 +2298,7 @@ public class DataRegion {
             logicalStorageGroupName, tsFileResource.getTimePartition())) {
       return true;
     }
-    if (!tsFileResource.isClosed()) {
-      // tsfile is not closed
-      return false;
-    }
+
     for (PartialPath device : devicePaths) {
       String deviceId = device.getFullPath();
       if (!tsFileResource.mayContainsDevice(deviceId)) {
@@ -2280,10 +2306,18 @@ public class DataRegion {
         continue;
       }
 
-      if (deleteEnd >= tsFileResource.getStartTime(deviceId)
-          && deleteStart <= tsFileResource.getEndTime(deviceId)) {
-        // time range of device has overlap with the deletion
-        return false;
+      long deviceEndTime = tsFileResource.getEndTime(deviceId);
+      if (!tsFileResource.isClosed() && deviceEndTime == Long.MIN_VALUE) {
+        // unsealed seq file
+        if (deleteEnd >= tsFileResource.getStartTime(deviceId)) {
+          return false;
+        }
+      } else {
+        // sealed file or unsealed unseq file
+        if (deleteEnd >= tsFileResource.getStartTime(deviceId) && deleteStart 
<= deviceEndTime) {
+          // time range of device has overlap with the deletion
+          return false;
+        }
       }
     }
     return true;
@@ -2306,36 +2340,38 @@ public class DataRegion {
         continue;
       }
 
-      if (tsFileResource.isCompacting()) {
-        // we have to set modification offset to MAX_VALUE, as the offset of 
source chunk may
-        // change after compaction
-        deletion.setFileOffset(Long.MAX_VALUE);
-        // write deletion into compaction modification file
-        tsFileResource.getCompactionModFile().write(deletion);
-        // write deletion into modification file to enable query during 
compaction
-        tsFileResource.getModFile().write(deletion);
-        // remember to close mod file
-        tsFileResource.getCompactionModFile().close();
-        tsFileResource.getModFile().close();
-      } else if (tsFileResource.isClosed()) {
-        deletion.setFileOffset(tsFileResource.getTsFileSize());
-        // write deletion into modification file
-        tsFileResource.getModFile().write(deletion);
-        // remember to close mod file
-        tsFileResource.getModFile().close();
+      if (tsFileResource.isClosed()) {
+        // delete data in sealed file
+        if (tsFileResource.isCompacting()) {
+          // we have to set modification offset to MAX_VALUE, as the offset of 
source chunk may
+          // change after compaction
+          deletion.setFileOffset(Long.MAX_VALUE);
+          // write deletion into compaction modification file
+          tsFileResource.getCompactionModFile().write(deletion);
+          // write deletion into modification file to enable query during 
compaction
+          tsFileResource.getModFile().write(deletion);
+          // remember to close mod file
+          tsFileResource.getCompactionModFile().close();
+          tsFileResource.getModFile().close();
+        } else {
+          deletion.setFileOffset(tsFileResource.getTsFileSize());
+          // write deletion into modification file
+          tsFileResource.getModFile().write(deletion);
+          // remember to close mod file
+          tsFileResource.getModFile().close();
+        }
+        logger.info(
+            "[Deletion] Deletion with path:{}, time:{}-{} written into mods 
file:{}.",
+            deletion.getPath(),
+            deletion.getStartTime(),
+            deletion.getEndTime(),
+            tsFileResource.getModFile().getFilePath());
+      } else {
+        // delete data in memory of unsealed file
+        tsFileResource.getProcessor().deleteDataInMemory(deletion, 
devicePaths);
       }
-      logger.info(
-          "[Deletion] Deletion with path:{}, time:{}-{} written into mods 
file:{}.",
-          deletion.getPath(),
-          deletion.getStartTime(),
-          deletion.getEndTime(),
-          tsFileResource.getModFile().getFilePath());
 
-      // delete data in memory of unsealed file
-      if (!tsFileResource.isClosed()) {
-        TsFileProcessor tsfileProcessor = tsFileResource.getProcessor();
-        tsfileProcessor.deleteDataInMemory(deletion, devicePaths);
-      } else if (tsFileSyncManager.isEnableSync()) {
+      if (tsFileSyncManager.isEnableSync()) {
         tsFileSyncManager.collectRealTimeDeletion(deletion);
       }
 
diff --git 
a/server/src/test/java/org/apache/iotdb/db/engine/storagegroup/DataRegionTest.java
 
b/server/src/test/java/org/apache/iotdb/db/engine/storagegroup/DataRegionTest.java
index f51303f07e..e1ca36db14 100644
--- 
a/server/src/test/java/org/apache/iotdb/db/engine/storagegroup/DataRegionTest.java
+++ 
b/server/src/test/java/org/apache/iotdb/db/engine/storagegroup/DataRegionTest.java
@@ -979,7 +979,7 @@ public class DataRegionTest {
   }
 
   @Test
-  public void testDeleteDataInFlushingMemtable()
+  public void testDeleteDataNotInFlushingMemtable()
       throws IllegalPathException, WriteProcessException, 
TriggerExecutionException, IOException {
     for (int j = 0; j < 100; j++) {
       TSRecord record = new TSRecord(j, deviceId);
@@ -996,9 +996,90 @@ public class DataRegionTest {
     // delete data which is not in memtable
     dataRegion.delete(new PartialPath("root.vehicle.d200.s0"), 50, 70, 0, 
null);
 
+    dataRegion.syncCloseAllWorkingTsFileProcessors();
+    Assert.assertFalse(tsFileResource.getModFile().exists());
+  }
+
+  @Test
+  public void testDeleteDataInSeqFlushingMemtable()
+      throws IllegalPathException, WriteProcessException, 
TriggerExecutionException, IOException {
+    for (int j = 100; j < 200; j++) {
+      TSRecord record = new TSRecord(j, deviceId);
+      record.addTuple(DataPoint.getDataPoint(TSDataType.INT32, measurementId, 
String.valueOf(j)));
+      dataRegion.insert(buildInsertRowNodeByTSRecord(record));
+    }
+    TsFileResource tsFileResource = 
dataRegion.getTsFileManager().getTsFileList(true).get(0);
+    TsFileProcessor tsFileProcessor = tsFileResource.getProcessor();
+    
tsFileProcessor.getFlushingMemTable().addLast(tsFileProcessor.getWorkMemTable());
+
+    // delete data which is not in flushing memtable
+    dataRegion.delete(new PartialPath("root.vehicle.d0.s0"), 50, 99, 0, null);
+    dataRegion.delete(new PartialPath("root.vehicle.d200.s0"), 50, 70, 0, 
null);
+
+    // delete data which is in flushing memtable
+    dataRegion.delete(new PartialPath("root.vehicle.d0.s0"), 50, 100, 0, null);
+    dataRegion.delete(new PartialPath("root.vehicle.d0.s0"), 50, 150, 0, null);
+    dataRegion.delete(new PartialPath("root.vehicle.d0.s0"), 100, 300, 0, 
null);
+
+    dataRegion.syncCloseAllWorkingTsFileProcessors();
+    Assert.assertTrue(tsFileResource.getModFile().exists());
+    Assert.assertEquals(3, 
tsFileResource.getModFile().getModifications().size());
+  }
+
+  @Test
+  public void testDeleteDataInUnSeqFlushingMemtable()
+      throws IllegalPathException, WriteProcessException, 
TriggerExecutionException, IOException {
+    for (int j = 100; j < 200; j++) {
+      TSRecord record = new TSRecord(j, deviceId);
+      record.addTuple(DataPoint.getDataPoint(TSDataType.INT32, measurementId, 
String.valueOf(j)));
+      dataRegion.insert(buildInsertRowNodeByTSRecord(record));
+    }
+    TsFileResource tsFileResource = 
dataRegion.getTsFileManager().getTsFileList(true).get(0);
+
+    // delete data which is not in work memtable
+    dataRegion.delete(new PartialPath("root.vehicle.d0.s0"), 50, 99, 0, null);
+    dataRegion.delete(new PartialPath("root.vehicle.d200.s0"), 50, 70, 0, 
null);
+
+    // delete data which is in work memtable
+    dataRegion.delete(new PartialPath("root.vehicle.d0.s0"), 50, 100, 0, null);
+    dataRegion.delete(new PartialPath("root.vehicle.d0.s0"), 50, 150, 0, null);
+    dataRegion.delete(new PartialPath("root.vehicle.d0.s0"), 100, 300, 0, 
null);
+
+    dataRegion.syncCloseAllWorkingTsFileProcessors();
+    Assert.assertFalse(tsFileResource.getModFile().exists());
+
+    // insert unseq data points
+    for (int j = 50; j < 100; j++) {
+      TSRecord record = new TSRecord(j, deviceId);
+      record.addTuple(DataPoint.getDataPoint(TSDataType.INT32, measurementId, 
String.valueOf(j)));
+      dataRegion.insert(buildInsertRowNodeByTSRecord(record));
+    }
+    // delete data which is not in work memtable
+    dataRegion.delete(new PartialPath("root.vehicle.d0.s0"), 200, 299, 0, 
null);
+    dataRegion.delete(new PartialPath("root.vehicle.d200.s0"), 50, 70, 0, 
null);
+
+    // delete data which is in work memtable
+    dataRegion.delete(new PartialPath("root.vehicle.d0.s0"), 80, 85, 0, null);
+
+    Assert.assertFalse(tsFileResource.getModFile().exists());
+
+    tsFileResource = dataRegion.getTsFileManager().getTsFileList(false).get(0);
+    TsFileProcessor tsFileProcessor = tsFileResource.getProcessor();
+    
tsFileProcessor.getFlushingMemTable().addLast(tsFileProcessor.getWorkMemTable());
+
+    // delete data which is not in flushing memtable
+    dataRegion.delete(new PartialPath("root.vehicle.d0.s0"), 0, 49, 0, null);
+    dataRegion.delete(new PartialPath("root.vehicle.d0.s0"), 100, 200, 0, 
null);
+    dataRegion.delete(new PartialPath("root.vehicle.d200.s0"), 50, 70, 0, 
null);
+
+    // delete data which is in flushing memtable
+    dataRegion.delete(new PartialPath("root.vehicle.d0.s0"), 25, 50, 0, null);
+    dataRegion.delete(new PartialPath("root.vehicle.d0.s0"), 50, 80, 0, null);
+    dataRegion.delete(new PartialPath("root.vehicle.d0.s0"), 99, 150, 0, null);
+
     dataRegion.syncCloseAllWorkingTsFileProcessors();
     Assert.assertTrue(tsFileResource.getModFile().exists());
-    Assert.assertEquals(2, 
tsFileResource.getModFile().getModifications().size());
+    Assert.assertEquals(3, 
tsFileResource.getModFile().getModifications().size());
   }
 
   static class DummyDataRegion extends DataRegion {

Reply via email to