This is an automated email from the ASF dual-hosted git repository. haonan pushed a commit to branch cp_fix_empty_file in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit 7a4f4dc319b3a300379e4ff34a70266d4f9fb12b Author: Haonan <[email protected]> AuthorDate: Thu Aug 21 09:52:37 2025 +0800 Fix empty tsfile and resource generated when insert, load, kill -9 and restart (#16215) * Fix empty tsfile and resource generated when insert, load, kill -9 and restart * Fix more * Add IT --- .../org/apache/iotdb/db/it/IoTDBRestartIT.java | 59 ++++++++++++++++++++++ .../db/storageengine/dataregion/DataRegion.java | 35 +++++++------ 2 files changed, 78 insertions(+), 16 deletions(-) diff --git a/integration-test/src/test/java/org/apache/iotdb/db/it/IoTDBRestartIT.java b/integration-test/src/test/java/org/apache/iotdb/db/it/IoTDBRestartIT.java index f423b5928ee..aac9b7da597 100644 --- a/integration-test/src/test/java/org/apache/iotdb/db/it/IoTDBRestartIT.java +++ b/integration-test/src/test/java/org/apache/iotdb/db/it/IoTDBRestartIT.java @@ -22,9 +22,13 @@ import org.apache.iotdb.db.conf.IoTDBConfig; import org.apache.iotdb.db.conf.IoTDBDescriptor; import org.apache.iotdb.it.env.EnvFactory; import org.apache.iotdb.it.framework.IoTDBTestRunner; +import org.apache.iotdb.it.utils.TsFileGenerator; import org.apache.iotdb.itbase.category.ClusterIT; import org.apache.iotdb.itbase.category.LocalStandaloneIT; +import org.apache.commons.io.FileUtils; +import org.apache.tsfile.enums.TSDataType; +import org.apache.tsfile.write.schema.MeasurementSchema; import org.junit.After; import org.junit.Before; import org.junit.Ignore; @@ -34,10 +38,13 @@ import org.junit.runner.RunWith; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import java.io.File; +import java.nio.file.Files; import java.sql.Connection; import java.sql.ResultSet; import java.sql.SQLException; import java.sql.Statement; +import java.util.Collections; import static org.apache.iotdb.db.utils.constant.TestConstant.TIMESTAMP_STR; import static org.junit.Assert.assertEquals; @@ -379,4 +386,56 @@ public class IoTDBRestartIT { } } } + + @Test + public void testInsertLoadAndRecover() throws Exception { + try (Connection connection = EnvFactory.getEnv().getConnection(); + Statement statement = connection.createStatement()) { + statement.execute("create timeseries root.sg.d1.s1 with datatype=int32"); + statement.execute("insert into root.sg.d1(time,s1) values(2,2)"); + statement.execute("flush"); + } + File tmpDir = new File(Files.createTempDirectory("load").toUri()); + File tsfile = new File(tmpDir, "0-0-0-0.tsfile"); + try { + try (final TsFileGenerator generator = new TsFileGenerator(tsfile)) { + generator.registerTimeseries( + "root.sg.d1", Collections.singletonList(new MeasurementSchema("s1", TSDataType.INT32))); + generator.generateData("root.sg.d1", 1, 2, false); + } + try (Connection connection = EnvFactory.getEnv().getConnection(); + Statement statement = connection.createStatement()) { + statement.execute("insert into root.sg.d1(time,s1) values(1,1)"); + statement.execute(String.format("load \"%s\" ", tsfile.getAbsolutePath())); + try (ResultSet resultSet = statement.executeQuery("select s1 from root.sg.d1")) { + assertNotNull(resultSet); + int cnt = 0; + while (resultSet.next()) { + assertEquals(String.valueOf(cnt + 1), resultSet.getString(1)); + cnt++; + } + assertEquals(2, cnt); + } + } + + // restart dn + TestUtils.stopForciblyAndRestartDataNodes(); + + try (Connection connection = EnvFactory.getEnv().getConnection(); + Statement statement = connection.createStatement()) { + try (ResultSet resultSet = statement.executeQuery("select s1 from root.sg.d1")) { + assertNotNull(resultSet); + int cnt = 0; + while (resultSet.next()) { + assertEquals(String.valueOf(cnt + 1), resultSet.getString(1)); + cnt++; + } + assertEquals(2, cnt); + } + } + + } finally { + FileUtils.deleteDirectory(tmpDir); + } + } } 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 7200d3da018..a04c81f4801 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 @@ -531,9 +531,10 @@ public class DataRegion implements IDataRegionForQuery { } } } - for (List<TsFileResource> value : partitionTmpUnseqTsFiles.values()) { + for (List<TsFileResource> unseqTsFiles : partitionTmpUnseqTsFiles.values()) { + List<TsFileResource> unsealedTsFiles = new ArrayList<>(); // tsFiles without resource file are unsealed - for (TsFileResource resource : value) { + for (TsFileResource resource : unseqTsFiles) { if (resource.resourceFileExists()) { FileMetrics.getInstance() .addTsFile( @@ -542,25 +543,20 @@ public class DataRegion implements IDataRegionForQuery { resource.getTsFile().length(), false, resource.getTsFile().getName()); - } - if (resource.getModFile().exists()) { - FileMetrics.getInstance().increaseModFileNum(1); - FileMetrics.getInstance().increaseModFileSize(resource.getModFile().getSize()); - } - } - while (!value.isEmpty()) { - TsFileResource tsFileResource = value.get(value.size() - 1); - if (tsFileResource.resourceFileExists()) { - break; } else { - value.remove(value.size() - 1); WALRecoverListener recoverListener = - recoverUnsealedTsFile(tsFileResource, dataRegionRecoveryContext, false); + recoverUnsealedTsFile(resource, dataRegionRecoveryContext, false); if (recoverListener != null) { recoverListeners.add(recoverListener); } + unsealedTsFiles.add(resource); + } + if (resource.getModFile().exists()) { + FileMetrics.getInstance().increaseModFileNum(1); + FileMetrics.getInstance().increaseModFileSize(resource.getModFile().getSize()); } } + unseqTsFiles.removeAll(unsealedTsFiles); } // signal wal recover manager to recover this region's files WALRecoverManager.getInstance().getAllDataRegionScannedLatch().countDown(); @@ -910,7 +906,12 @@ public class DataRegion implements IDataRegionForQuery { new SealedTsFileRecoverPerformer(sealedTsFile)) { recoverPerformer.recover(); sealedTsFile.close(); - tsFileResourceManager.registerSealedTsFileResource(sealedTsFile); + if (!TsFileValidator.getInstance().validateTsFile(sealedTsFile)) { + sealedTsFile.remove(); + tsFileManager.remove(sealedTsFile, sealedTsFile.isSeq()); + } else { + tsFileResourceManager.registerSealedTsFileResource(sealedTsFile); + } } catch (Throwable e) { logger.error("Fail to recover sealed TsFile {}, skip it.", sealedTsFile.getTsFilePath(), e); } finally { @@ -1015,7 +1016,9 @@ public class DataRegion implements IDataRegionForQuery { lastFlushTimeMap.getMemSize(partitionId))); } for (TsFileResource tsFileResource : resourceList) { - updateDeviceLastFlushTime(tsFileResource); + if (!tsFileResource.isDeleted()) { + updateDeviceLastFlushTime(tsFileResource); + } } TimePartitionManager.getInstance() .updateAfterFlushing(
