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 429587f629a Load: preserve pending tablets when parser fails (#18487)
429587f629a is described below

commit 429587f629a597e63d47a026258a308b2dbee6ac
Author: Caideyipi <[email protected]>
AuthorDate: Mon Aug 24 12:14:17 2026 +0800

    Load: preserve pending tablets when parser fails (#18487)
---
 ...eeStatementDataTypeConvertExecutionVisitor.java | 63 ++++++++++++++++++++--
 ...atementDataTypeConvertExecutionVisitorTest.java | 59 ++++++++++++++++++++
 2 files changed, 119 insertions(+), 3 deletions(-)

diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/converter/LoadTreeStatementDataTypeConvertExecutionVisitor.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/converter/LoadTreeStatementDataTypeConvertExecutionVisitor.java
index a8f4c28e7c8..3f2a4fcef87 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/converter/LoadTreeStatementDataTypeConvertExecutionVisitor.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/converter/LoadTreeStatementDataTypeConvertExecutionVisitor.java
@@ -42,6 +42,7 @@ import java.io.File;
 import java.util.ArrayList;
 import java.util.List;
 import java.util.Optional;
+import java.util.function.Function;
 import java.util.stream.Collectors;
 
 import static 
org.apache.iotdb.db.pipe.resource.memory.PipeMemoryWeightUtil.calculateTabletSizeInBytes;
@@ -57,6 +58,7 @@ public class LoadTreeStatementDataTypeConvertExecutionVisitor
           .getLoadTsFileTabletConversionBatchMemorySizeInBytes();
 
   private final StatementExecutor statementExecutor;
+  private final Function<File, LoadTreeTsFileTabletIterator> 
tabletIteratorFactory;
 
   @FunctionalInterface
   public interface StatementExecutor {
@@ -65,7 +67,14 @@ public class LoadTreeStatementDataTypeConvertExecutionVisitor
 
   public LoadTreeStatementDataTypeConvertExecutionVisitor(
       final StatementExecutor statementExecutor) {
+    this(statementExecutor, file -> new LoadTreeTsFileTabletIterator(file, 
true));
+  }
+
+  LoadTreeStatementDataTypeConvertExecutionVisitor(
+      final StatementExecutor statementExecutor,
+      final Function<File, LoadTreeTsFileTabletIterator> 
tabletIteratorFactory) {
     this.statementExecutor = statementExecutor;
+    this.tabletIteratorFactory = tabletIteratorFactory;
   }
 
   @Override
@@ -87,11 +96,24 @@ public class 
LoadTreeStatementDataTypeConvertExecutionVisitor
     boolean shouldReleaseContext = !isManagedTask;
 
     try {
+      if (conversionContext.deferredStatus != null) {
+        final TSStatus result =
+            flushPendingTablets(conversionContext, 
loadTsFileStatement.isConvertOnTypeMismatch());
+        if (!handleTSStatus(result, loadTsFileStatement)) {
+          shouldReleaseContext = !isManagedTask || 
!isTemporaryUnavailable(result);
+          return Optional.of(result);
+        }
+        final TSStatus deferredStatus = conversionContext.deferredStatus;
+        conversionContext.deferredStatus = null;
+        shouldReleaseContext = true;
+        return Optional.of(deferredStatus);
+      }
+
       final List<File> files = loadTsFileStatement.getTsFiles();
       while (conversionContext.fileIndex < files.size()) {
         if (conversionContext.tabletIterator == null) {
           conversionContext.tabletIterator =
-              new 
LoadTreeTsFileTabletIterator(files.get(conversionContext.fileIndex), true);
+              
tabletIteratorFactory.apply(files.get(conversionContext.fileIndex));
         }
 
         if (conversionContext.deferredTabletRawReq != null) {
@@ -164,11 +186,29 @@ public class 
LoadTreeStatementDataTypeConvertExecutionVisitor
               
.STORAGE_LOG_FAILED_TO_CONVERT_DATA_TYPE_FOR_LOADTSFILESTATEMENT_5D132E57,
           loadTsFileStatement,
           e);
+      final boolean retryable = isRetryableConversionException(e);
       final TSStatus status =
           loadTsFileStatement.accept(
               LoadTsFileDataTypeConverter.TREE_STATEMENT_EXCEPTION_VISITOR, e);
-      shouldReleaseContext =
-          !isManagedTask || 
!LoadTsFileDataTypeConverter.isMemoryPressureException(e);
+
+      // A parser can fail after producing tablets that are still waiting for 
the next batch
+      // boundary. Submit those tablets before reporting the parser error so 
the error does not
+      // discard successfully converted data.
+      if (!retryable && !conversionContext.tabletRawReqs.isEmpty()) {
+        final TSStatus flushStatus =
+            flushPendingTablets(conversionContext, 
loadTsFileStatement.isConvertOnTypeMismatch());
+        if (!handleTSStatus(flushStatus, loadTsFileStatement)) {
+          if (isManagedTask && isTemporaryUnavailable(flushStatus)) {
+            conversionContext.deferredStatus = status;
+            shouldReleaseContext = false;
+          } else {
+            shouldReleaseContext = true;
+          }
+          return Optional.of(flushStatus);
+        }
+      }
+
+      shouldReleaseContext = !isManagedTask || !retryable;
       return Optional.of(status);
     } finally {
       if (shouldReleaseContext) {
@@ -201,6 +241,21 @@ public class 
LoadTreeStatementDataTypeConvertExecutionVisitor
                 == 
TSStatusCode.PIPE_RECEIVER_TEMPORARY_UNAVAILABLE_EXCEPTION.getStatusCode());
   }
 
+  private static boolean isRetryableConversionException(final Throwable 
throwable) {
+    if (LoadTsFileDataTypeConverter.isMemoryPressureException(throwable)) {
+      return true;
+    }
+
+    Throwable current = throwable;
+    while (current != null) {
+      if (current instanceof InterruptedException) {
+        return true;
+      }
+      current = current.getCause();
+    }
+    return false;
+  }
+
   private static void deleteSourceFiles(final LoadTsFileStatement statement) {
     statement
         .getTsFiles()
@@ -228,6 +283,7 @@ public class 
LoadTreeStatementDataTypeConvertExecutionVisitor
     private Pair<Tablet, Boolean> deferredTabletWithIsAligned;
     private PipeTransferTabletRawReq deferredTabletRawReq;
     private long deferredTabletRawReqSize;
+    private TSStatus deferredStatus;
 
     private void addTablet(final PipeTransferTabletRawReq request, final long 
size) {
       tabletRawReqs.add(request);
@@ -259,6 +315,7 @@ public class 
LoadTreeStatementDataTypeConvertExecutionVisitor
       deferredTabletWithIsAligned = null;
       deferredTabletRawReq = null;
       deferredTabletRawReqSize = 0;
+      deferredStatus = null;
       block.close();
     }
   }
diff --git 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/load/converter/LoadTreeStatementDataTypeConvertExecutionVisitorTest.java
 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/load/converter/LoadTreeStatementDataTypeConvertExecutionVisitorTest.java
index c82b9bf7a2f..b374880d6a8 100644
--- 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/load/converter/LoadTreeStatementDataTypeConvertExecutionVisitorTest.java
+++ 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/load/converter/LoadTreeStatementDataTypeConvertExecutionVisitorTest.java
@@ -40,6 +40,7 @@ import org.apache.tsfile.file.metadata.IDeviceID;
 import org.apache.tsfile.file.metadata.enums.TSEncoding;
 import org.apache.tsfile.read.TsFileSequenceReader;
 import org.apache.tsfile.read.common.Path;
+import org.apache.tsfile.utils.Pair;
 import org.apache.tsfile.write.TsFileWriter;
 import org.apache.tsfile.write.record.Tablet;
 import org.apache.tsfile.write.schema.IMeasurementSchema;
@@ -113,6 +114,32 @@ public class 
LoadTreeStatementDataTypeConvertExecutionVisitorTest {
     Assert.assertEquals(loadedPointCountBeforeCorruption, 
loadedPointCountAfterFallback);
   }
 
+  @Test
+  public void testFlushesPendingTabletsWhenIteratorFails() throws Exception {
+    tsFile = File.createTempFile("load-tree-pending-tablet", ".tsfile");
+    final List<IMeasurementSchema> schemaList =
+        Arrays.asList(new MeasurementSchema("s0", TSDataType.INT64, 
TSEncoding.PLAIN));
+    final Tablet tablet = new Tablet(DEVICE_0, schemaList, 1);
+    tablet.addTimestamp(0, 1);
+    tablet.addValue("s0", 0, 1L);
+
+    final Map<String, Integer> pointCountByDevice = new HashMap<>();
+    final LoadTreeStatementDataTypeConvertExecutionVisitor visitor =
+        new LoadTreeStatementDataTypeConvertExecutionVisitor(
+            statement -> {
+              collectLoadedPoints(statement, pointCountByDevice);
+              return new TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode());
+            },
+            file -> new ThrowingTabletIterator(file, tablet));
+
+    final Optional<TSStatus> status =
+        
visitor.visitLoadFile(LoadTsFileStatement.createUnchecked(tsFile.getAbsolutePath()),
 null);
+
+    Assert.assertTrue(status.isPresent());
+    Assert.assertEquals(TSStatusCode.LOAD_FILE_ERROR.getStatusCode(), 
status.get().getCode());
+    Assert.assertEquals(1, pointCountByDevice.getOrDefault(DEVICE_0, 
0).intValue());
+  }
+
   @Test
   public void testLoadScanParserUsesQueryMemoryPoolInsteadOfPipeMemory() 
throws Exception {
     tsFile = new File("load-tree-parser-query-memory.tsfile");
@@ -421,4 +448,36 @@ public class 
LoadTreeStatementDataTypeConvertExecutionVisitorTest {
       
CommonDescriptor.getInstance().getConfig().setPipeMemoryManagementEnabled(enabled);
     }
   }
+
+  private static class ThrowingTabletIterator extends 
LoadTreeTsFileTabletIterator {
+    private final Pair<Tablet, Boolean> tabletWithIsAligned;
+    private boolean tabletAvailable = true;
+
+    private ThrowingTabletIterator(final File file, final Tablet tablet) {
+      super(file, true);
+      tabletWithIsAligned = new Pair<>(tablet, false);
+    }
+
+    @Override
+    public boolean hasNext() {
+      if (tabletAvailable) {
+        return true;
+      }
+      throw new IllegalStateException("synthetic parser failure");
+    }
+
+    @Override
+    public Pair<Tablet, Boolean> next() {
+      if (!tabletAvailable) {
+        throw new IllegalStateException("synthetic parser failure");
+      }
+      tabletAvailable = false;
+      return tabletWithIsAligned;
+    }
+
+    @Override
+    public void close() {
+      // No parser resources are allocated by this test iterator.
+    }
+  }
 }

Reply via email to