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 {