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());
         }

Reply via email to