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 31e75348112 [Pipe] Fix tablet type conversion for failed columns
(#18595)
31e75348112 is described below
commit 31e7534811217440172d3b2edd27170ef8328c19
Author: Caideyipi <[email protected]>
AuthorDate: Tue Sep 8 14:55:47 2026 +0800
[Pipe] Fix tablet type conversion for failed columns (#18595)
---
.../PipeConvertedInsertTabletStatement.java | 25 ++++
.../LoadConvertedInsertTabletStatement.java | 4 +
.../PipeConvertedInsertTabletStatementTest.java | 133 +++++++++++++++++++++
3 files changed, 162 insertions(+)
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/receiver/transform/statement/PipeConvertedInsertTabletStatement.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/receiver/transform/statement/PipeConvertedInsertTabletStatement.java
index f7e6f0be1e6..a4f0b8b29ee 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/receiver/transform/statement/PipeConvertedInsertTabletStatement.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/receiver/transform/statement/PipeConvertedInsertTabletStatement.java
@@ -99,6 +99,9 @@ public class PipeConvertedInsertTabletStatement extends
InsertTabletStatement {
@Override
protected boolean checkAndCastDataType(int columnIndex, TSDataType dataType)
{
+ if (!isValidColumnForTypeConversion(columnIndex, dataType)) {
+ return false;
+ }
if (LOGGER.isInfoEnabled()) {
PipeLogger.log(
LOGGER::info,
@@ -112,6 +115,28 @@ public class PipeConvertedInsertTabletStatement extends
InsertTabletStatement {
return true;
}
+ /**
+ * Returns whether a tablet column has all state required for type
conversion.
+ *
+ * <p>Partial insert processing can leave a measurement slot without a data
type or value column
+ * (and deserialized statements can contain arrays of different lengths). In
that case the column
+ * must be treated as failed instead of being indexed by the conversion path.
+ */
+ protected boolean isValidColumnForTypeConversion(
+ final int columnIndex, final TSDataType dataType) {
+ return dataType != null
+ && dataTypes != null
+ && columns != null
+ && measurements != null
+ && columnIndex >= 0
+ && columnIndex < measurements.length
+ && columnIndex < dataTypes.length
+ && columnIndex < columns.length
+ && measurements[columnIndex] != null
+ && dataTypes[columnIndex] != null
+ && columns[columnIndex] != null;
+ }
+
protected boolean originalCheckAndCastDataType(int columnIndex, TSDataType
dataType) {
return super.checkAndCastDataType(columnIndex, dataType);
}
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/converter/LoadConvertedInsertTabletStatement.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/converter/LoadConvertedInsertTabletStatement.java
index 591756e9f11..3edf9e150b8 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/converter/LoadConvertedInsertTabletStatement.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/converter/LoadConvertedInsertTabletStatement.java
@@ -49,6 +49,10 @@ public class LoadConvertedInsertTabletStatement extends
PipeConvertedInsertTable
return originalCheckAndCastDataType(columnIndex, dataType);
}
+ if (!isValidColumnForTypeConversion(columnIndex, dataType)) {
+ return false;
+ }
+
LOGGER.info(
StorageEngineMessages.STORAGE_LOG_LOAD_INSERTING_TABLET_TO_CASTING_TYPE_FROM_TO_AE808A8B,
devicePath,
diff --git
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/receiver/transform/statement/PipeConvertedInsertTabletStatementTest.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/receiver/transform/statement/PipeConvertedInsertTabletStatementTest.java
new file mode 100644
index 00000000000..a6bd537bee4
--- /dev/null
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/receiver/transform/statement/PipeConvertedInsertTabletStatementTest.java
@@ -0,0 +1,133 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.iotdb.db.pipe.receiver.transform.statement;
+
+import org.apache.iotdb.commons.path.PartialPath;
+import org.apache.iotdb.db.conf.IoTDBDescriptor;
+import
org.apache.iotdb.db.queryengine.plan.statement.crud.InsertTabletStatement;
+import
org.apache.iotdb.db.storageengine.load.converter.LoadConvertedInsertTabletStatement;
+
+import org.apache.tsfile.enums.TSDataType;
+import org.apache.tsfile.write.record.Tablet;
+import org.apache.tsfile.write.schema.MeasurementSchema;
+import org.junit.After;
+import org.junit.Assert;
+import org.junit.Before;
+import org.junit.Test;
+
+public class PipeConvertedInsertTabletStatementTest {
+
+ private boolean enablePartialInsert;
+
+ @Before
+ public void setUp() {
+ enablePartialInsert =
IoTDBDescriptor.getInstance().getConfig().isEnablePartialInsert();
+ IoTDBDescriptor.getInstance().getConfig().setEnablePartialInsert(true);
+ }
+
+ @After
+ public void tearDown() {
+
IoTDBDescriptor.getInstance().getConfig().setEnablePartialInsert(enablePartialInsert);
+ }
+
+ @Test
+ public void testTypeConversionSkipsMissingColumnWithoutException() throws
Exception {
+ final InsertTabletStatement source =
createPartiallyFailedStatementWithMissingColumn();
+ final PipeConvertedInsertTabletStatement converted =
+ new PipeConvertedInsertTabletStatement(source);
+
+ converted.selfCheckDataTypes(0);
+ converted.selfCheckDataTypes(1);
+
+ final Tablet tablet = converted.convertToTablet();
+ Assert.assertEquals(1, tablet.getSchemas().size());
+ Assert.assertEquals("s1", tablet.getSchemas().get(0).getMeasurementName());
+ Assert.assertEquals(TSDataType.DOUBLE,
tablet.getSchemas().get(0).getType());
+ Assert.assertArrayEquals(new double[] {1.0, 2.0}, (double[])
tablet.getValues()[0], 0.0);
+ Assert.assertNull(converted.getMeasurements()[1]);
+ Assert.assertNull(converted.getDataTypes()[1]);
+ }
+
+ @Test
+ public void testTypeConversionSkipsNullMeasurementWithoutException() throws
Exception {
+ final InsertTabletStatement source = new InsertTabletStatement();
+ source.setDevicePath(new PartialPath("root.sg.d1"));
+ source.setMeasurements(new String[] {"s1", "s2"});
+ source.setDataTypes(new TSDataType[] {TSDataType.INT32, null});
+ source.setMeasurementSchemas(
+ new MeasurementSchema[] {
+ new MeasurementSchema("s1", TSDataType.DOUBLE),
+ new MeasurementSchema("s2", TSDataType.DOUBLE)
+ });
+ source.setTimes(new long[] {1L, 2L});
+ source.setColumns(new Object[] {new int[] {1, 2}, null});
+ source.setRowCount(2);
+ source.selfCheckDataTypes(1);
+
+ Assert.assertNull(source.getMeasurements()[1]);
+
+ final PipeConvertedInsertTabletStatement converted =
+ new PipeConvertedInsertTabletStatement(source);
+ converted.selfCheckDataTypes(0);
+ converted.selfCheckDataTypes(1);
+
+ final Tablet tablet = converted.convertToTablet();
+
+ Assert.assertEquals(1, tablet.getSchemas().size());
+ Assert.assertEquals("s1", tablet.getSchemas().get(0).getMeasurementName());
+ Assert.assertEquals(TSDataType.DOUBLE,
tablet.getSchemas().get(0).getType());
+ Assert.assertArrayEquals(new double[] {1.0, 2.0}, (double[])
tablet.getValues()[0], 0.0);
+ Assert.assertNull(converted.getMeasurements()[1]);
+ Assert.assertNull(converted.getDataTypes()[1]);
+ }
+
+ @Test
+ public void testLoadTypeConversionSkipsMissingColumnWithoutException()
throws Exception {
+ final LoadConvertedInsertTabletStatement converted =
+ new LoadConvertedInsertTabletStatement(
+ createPartiallyFailedStatementWithMissingColumn(), true);
+
+ converted.selfCheckDataTypes(0);
+ converted.selfCheckDataTypes(1);
+
+ final Tablet tablet = converted.convertToTablet();
+ Assert.assertEquals(1, tablet.getSchemas().size());
+ Assert.assertEquals(TSDataType.DOUBLE,
tablet.getSchemas().get(0).getType());
+ Assert.assertArrayEquals(new double[] {1.0, 2.0}, (double[])
tablet.getValues()[0], 0.0);
+ }
+
+ private static InsertTabletStatement
createPartiallyFailedStatementWithMissingColumn()
+ throws Exception {
+ final InsertTabletStatement source = new InsertTabletStatement();
+ source.setDevicePath(new PartialPath("root.sg.d1"));
+ source.setMeasurements(new String[] {"s1", "s2"});
+ source.setDataTypes(new TSDataType[] {TSDataType.INT32, TSDataType.INT64});
+ source.setMeasurementSchemas(
+ new MeasurementSchema[] {
+ new MeasurementSchema("s1", TSDataType.DOUBLE),
+ new MeasurementSchema("s2", TSDataType.DOUBLE)
+ });
+ source.setTimes(new long[] {1L, 2L});
+ source.setColumns(new Object[] {new int[] {1, 2}});
+ source.setRowCount(2);
+ source.selfCheckDataTypes(1);
+ return source;
+ }
+}