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.
+ }
+ }
}