This is an automated email from the ASF dual-hosted git repository.
jt2594838 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 34349e3ad8b Fix snapshot omission of closing TsFiles (#18682)
34349e3ad8b is described below
commit 34349e3ad8b626e4441ffb48d8786c73892fd4b0
Author: Hongzhi Gao <[email protected]>
AuthorDate: Mon Sep 21 11:53:15 2026 +0800
Fix snapshot omission of closing TsFiles (#18682)
---
.../iotdb/db/i18n/StorageEngineMessages.java | 2 +
.../iotdb/db/i18n/StorageEngineMessages.java | 2 +
.../db/storageengine/dataregion/DataRegion.java | 20 +++++++++-
.../dataregion/snapshot/SnapshotTaker.java | 6 +--
.../storageengine/dataregion/DataRegionTest.java | 43 ++++++++++++++++++++++
.../dataregion/snapshot/IoTDBSnapshotTest.java | 10 ++---
6 files changed, 72 insertions(+), 11 deletions(-)
diff --git
a/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/StorageEngineMessages.java
b/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/StorageEngineMessages.java
index 3d66a888696..5d946aed646 100644
---
a/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/StorageEngineMessages.java
+++
b/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/StorageEngineMessages.java
@@ -411,6 +411,8 @@ public final class StorageEngineMessages {
public static final String FAILED_TO_CLOSE_SNAPSHOT_LOGGER = "Failed to
close snapshot logger";
public static final String SNAPSHOTTING_COMPRESSION_RATIO = "Snapshotting
compression ratio {}.";
public static final String CATCH_IO_EXCEPTION_CREATING_SNAPSHOT = "Catch
IOException when creating snapshot";
+ public static final String CANNOT_SNAPSHOT_UNCLOSED_TSFILE =
+ "Cannot create snapshot because TsFile {} is not closed";
public static final String HARD_LINK_TARGET_DIR_NOT_EXIST = "Hard link
target dir {} doesn't exist";
public static final String HARD_LINK_SOURCE_FILE_NOT_EXIST = "Hard link
source file {} doesn't exist, this file will be ignored.";
public static final String COPY_TARGET_DIR_NOT_EXIST = "Copy target dir {}
doesn't exist";
diff --git
a/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/StorageEngineMessages.java
b/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/StorageEngineMessages.java
index dfc7cea9871..8324777fc85 100644
---
a/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/StorageEngineMessages.java
+++
b/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/StorageEngineMessages.java
@@ -411,6 +411,8 @@ public final class StorageEngineMessages {
public static final String FAILED_TO_CLOSE_SNAPSHOT_LOGGER = "关闭快照日志器失败";
public static final String SNAPSHOTTING_COMPRESSION_RATIO = "正在快照压缩率文件 {}。";
public static final String CATCH_IO_EXCEPTION_CREATING_SNAPSHOT = "创建快照时捕获到
IOException";
+ public static final String CANNOT_SNAPSHOT_UNCLOSED_TSFILE =
+ "无法创建快照,因为 TsFile {} 尚未关闭";
public static final String HARD_LINK_TARGET_DIR_NOT_EXIST = "硬链接目标目录 {} 不存在";
public static final String HARD_LINK_SOURCE_FILE_NOT_EXIST = "硬链接源文件 {}
不存在,该文件将被忽略。";
public static final String COPY_TARGET_DIR_NOT_EXIST = "复制目标目录 {} 不存在";
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/DataRegion.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/DataRegion.java
index 6b14b43edf5..45c3e4d560f 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/DataRegion.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/DataRegion.java
@@ -2520,7 +2520,25 @@ public class DataRegion implements IDataRegionForQuery {
/** This method will be blocked until all tsfile processors are closed. */
public void syncCloseAllWorkingTsFileProcessors() {
try {
- List<Future<?>> tsFileProcessorsClosingFutures =
asyncCloseAllWorkingTsFileProcessors();
+ List<Future<?>> tsFileProcessorsClosingFutures = new ArrayList<>();
+ writeLock("syncCloseAllWorkingTsFileProcessors");
+ try {
+ for (TsFileProcessor tsFileProcessor : closingSequenceTsFileProcessor)
{
+ Future<?> closeFuture = tsFileProcessor.getCloseFuture();
+ if (closeFuture != null) {
+ tsFileProcessorsClosingFutures.add(closeFuture);
+ }
+ }
+ for (TsFileProcessor tsFileProcessor :
closingUnSequenceTsFileProcessor) {
+ Future<?> closeFuture = tsFileProcessor.getCloseFuture();
+ if (closeFuture != null) {
+ tsFileProcessorsClosingFutures.add(closeFuture);
+ }
+ }
+
tsFileProcessorsClosingFutures.addAll(asyncCloseAllWorkingTsFileProcessors());
+ } finally {
+ writeUnlock();
+ }
for (Future<?> f : tsFileProcessorsClosingFutures) {
if (f != null) {
f.get();
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/snapshot/SnapshotTaker.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/snapshot/SnapshotTaker.java
index 6707c3f6cc9..c471047e2f1 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/snapshot/SnapshotTaker.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/snapshot/SnapshotTaker.java
@@ -226,12 +226,10 @@ public class SnapshotTaker {
try {
for (TsFileResource resource : resources) {
if (!resource.isClosed()) {
- continue;
+ LOGGER.error(StorageEngineMessages.CANNOT_SNAPSHOT_UNCLOSED_TSFILE,
resource);
+ return false;
}
File tsFile = resource.getTsFile();
- if (!resource.isClosed()) {
- continue;
- }
File snapshotTsFile = getSnapshotFilePathForTsFile(tsFile, snapshotId);
File snapshotResourceFile =
new File(snapshotTsFile.getAbsolutePath() +
TsFileResource.RESOURCE_SUFFIX);
diff --git
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/DataRegionTest.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/DataRegionTest.java
index 1000ad31cc4..4e2ca24a750 100644
---
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/DataRegionTest.java
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/DataRegionTest.java
@@ -106,14 +106,18 @@ import org.slf4j.LoggerFactory;
import java.io.File;
import java.io.IOException;
+import java.lang.reflect.Field;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collections;
import java.util.List;
import java.util.Map;
+import java.util.Set;
import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.CountDownLatch;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.Future;
+import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
import static
org.apache.iotdb.db.queryengine.plan.statement.StatementTestUtils.genInsertRowNode;
@@ -2262,4 +2266,43 @@ public class DataRegionTest {
future.get();
assertTrue(tsFileResourceSeq.isClosed());
}
+
+ @Test
+ @SuppressWarnings("unchecked")
+ public void testSyncCloseWaitsForAlreadyClosingProcessor() throws Exception {
+ TsFileProcessor processor = Mockito.mock(TsFileProcessor.class);
+ Future<?> closeFuture = Mockito.mock(Future.class);
+ CountDownLatch waitStarted = new CountDownLatch(1);
+ CountDownLatch allowClose = new CountDownLatch(1);
+ Mockito.doReturn(closeFuture).when(processor).getCloseFuture();
+ Mockito.when(closeFuture.get())
+ .thenAnswer(
+ invocation -> {
+ waitStarted.countDown();
+ assertTrue(allowClose.await(10, TimeUnit.SECONDS));
+ return null;
+ });
+
+ Field closingProcessorsField =
+ DataRegion.class.getDeclaredField("closingSequenceTsFileProcessor");
+ closingProcessorsField.setAccessible(true);
+ Set<TsFileProcessor> closingProcessors =
+ (Set<TsFileProcessor>) closingProcessorsField.get(dataRegion);
+ closingProcessors.add(processor);
+
+ CompletableFuture<Void> syncCloseTask = null;
+ try {
+ syncCloseTask =
CompletableFuture.runAsync(dataRegion::syncCloseAllWorkingTsFileProcessors);
+ assertTrue(waitStarted.await(10, TimeUnit.SECONDS));
+ Assert.assertFalse(syncCloseTask.isDone());
+ allowClose.countDown();
+ syncCloseTask.get(10, TimeUnit.SECONDS);
+ } finally {
+ allowClose.countDown();
+ closingProcessors.remove(processor);
+ if (syncCloseTask != null) {
+ syncCloseTask.get(10, TimeUnit.SECONDS);
+ }
+ }
+ }
}
diff --git
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/snapshot/IoTDBSnapshotTest.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/snapshot/IoTDBSnapshotTest.java
index 2eedca1807f..84d8deac3d7 100644
---
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/snapshot/IoTDBSnapshotTest.java
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/snapshot/IoTDBSnapshotTest.java
@@ -149,7 +149,7 @@ public class IoTDBSnapshotTest {
}
@Test
- public void testCreateSnapshotWithUnclosedTsFile()
+ public void testRejectSnapshotWithUnclosedTsFile()
throws IOException, WriteProcessException, DirectoryNotLegalException {
String[][] originDataDirs =
IoTDBDescriptor.getInstance().getConfig().getTierDataDirs();
IoTDBDescriptor.getInstance().getConfig().setTierDataDirs(testDataDirs);
@@ -163,16 +163,14 @@ public class IoTDBSnapshotTest {
File snapshotDir = new File("target" + File.separator + "snapshot");
Assert.assertTrue(snapshotDir.exists() || snapshotDir.mkdirs());
try {
- new
SnapshotTaker(region).takeFullSnapshot(snapshotDir.getAbsolutePath(), true);
+ Assert.assertFalse(
+ new
SnapshotTaker(region).takeFullSnapshot(snapshotDir.getAbsolutePath(), true));
File[] files =
snapshotDir.listFiles((dir, name) ->
name.equals(SnapshotLogger.SNAPSHOT_LOG_NAME));
assertEquals(1, files.length);
SnapshotLogAnalyzer analyzer = new SnapshotLogAnalyzer(files[0]);
- int cnt = 0;
- Assert.assertTrue(analyzer.isSnapshotComplete());
- cnt = analyzer.getTotalFileCountInSnapshot();
+ Assert.assertFalse(analyzer.isSnapshotComplete());
analyzer.close();
- assertEquals(100, cnt);
for (TsFileResource resource : resources) {
Assert.assertTrue(resource.tryWriteLock());
}