This is an automated email from the ASF dual-hosted git repository.
davidzollo pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/seatunnel.git
The following commit(s) were added to refs/heads/dev by this push:
new b94a65ac3e [Fix][Connector-V2] Fix nested ROW type conversion crash in
Paimon connector (#10952)
b94a65ac3e is described below
commit b94a65ac3edb1f315a76356e19ce0716c192bf15
Author: cosmosni <[email protected]>
AuthorDate: Tue May 26 20:37:41 2026 +0800
[Fix][Connector-V2] Fix nested ROW type conversion crash in Paimon
connector (#10952)
---
.../seatunnel/paimon/utils/RowConverter.java | 25 +++--
.../seatunnel/paimon/utils/RowConverterTest.java | 108 +++++++++++++++++++++
2 files changed, 127 insertions(+), 6 deletions(-)
diff --git
a/seatunnel-connectors-v2/connector-paimon/src/main/java/org/apache/seatunnel/connectors/seatunnel/paimon/utils/RowConverter.java
b/seatunnel-connectors-v2/connector-paimon/src/main/java/org/apache/seatunnel/connectors/seatunnel/paimon/utils/RowConverter.java
index 7766aef546..c47e5e140b 100644
---
a/seatunnel-connectors-v2/connector-paimon/src/main/java/org/apache/seatunnel/connectors/seatunnel/paimon/utils/RowConverter.java
+++
b/seatunnel-connectors-v2/connector-paimon/src/main/java/org/apache/seatunnel/connectors/seatunnel/paimon/utils/RowConverter.java
@@ -58,6 +58,7 @@ import java.math.RoundingMode;
import java.time.LocalDate;
import java.time.LocalDateTime;
import java.time.LocalTime;
+import java.util.Collections;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
@@ -338,7 +339,9 @@ public class RowConverter {
SeaTunnelDataType<?> rowType =
seaTunnelRowType.getFieldType(i);
InternalRow row =
rowData.getRow(i, ((SeaTunnelRowType)
rowType).getTotalFields());
- objects[i] = convert(row, (SeaTunnelRowType) rowType,
tableSchema);
+ TableSchema nestedSchema =
+
extractNestedSchema(tableSchema.fields().get(i).type());
+ objects[i] = convert(row, (SeaTunnelRowType) rowType,
nestedSchema);
break;
default:
throw CommonError.unsupportedDataType(
@@ -480,13 +483,11 @@ public class RowConverter {
case ROW:
SeaTunnelDataType<?> rowType =
seaTunnelRowType.getFieldType(i);
Object row = fieldValue;
+ TableSchema nestedSchema =
extractNestedSchema(sinkTotalFields.get(i).type());
InternalRow paimonRow =
- reconvert(
- (SeaTunnelRow) row,
- (SeaTunnelRowType) rowType,
- sinkTableSchema);
+ reconvert((SeaTunnelRow) row, (SeaTunnelRowType)
rowType, nestedSchema);
RowType paimonRowType =
- RowTypeConverter.reconvert((SeaTunnelRowType)
rowType, sinkTableSchema);
+ RowTypeConverter.reconvert((SeaTunnelRowType)
rowType, nestedSchema);
binaryWriter.writeRow(i, paimonRow, new
InternalRowSerializer(paimonRowType));
break;
default:
@@ -499,6 +500,18 @@ public class RowConverter {
return binaryRow;
}
+ private static TableSchema extractNestedSchema(DataType nestedType) {
+ RowType nestedRowType = (RowType) nestedType;
+ return new TableSchema(
+ 0,
+ nestedRowType.getFields(),
+ nestedRowType.getFieldCount(),
+ Collections.emptyList(),
+ Collections.emptyList(),
+ Collections.emptyMap(),
+ null);
+ }
+
private static void checkCanWriteWithSchema(
int i, SeaTunnelRowType seaTunnelRowType, List<DataField> fields,
Object fieldValue) {
String sourceFieldName = seaTunnelRowType.getFieldName(i);
diff --git
a/seatunnel-connectors-v2/connector-paimon/src/test/java/org/apache/seatunnel/connectors/seatunnel/paimon/utils/RowConverterTest.java
b/seatunnel-connectors-v2/connector-paimon/src/test/java/org/apache/seatunnel/connectors/seatunnel/paimon/utils/RowConverterTest.java
index 6a8a445e78..f62884c005 100644
---
a/seatunnel-connectors-v2/connector-paimon/src/test/java/org/apache/seatunnel/connectors/seatunnel/paimon/utils/RowConverterTest.java
+++
b/seatunnel-connectors-v2/connector-paimon/src/test/java/org/apache/seatunnel/connectors/seatunnel/paimon/utils/RowConverterTest.java
@@ -43,6 +43,7 @@ import org.apache.paimon.data.InternalRow;
import org.apache.paimon.data.Timestamp;
import org.apache.paimon.data.serializer.InternalArraySerializer;
import org.apache.paimon.data.serializer.InternalMapSerializer;
+import org.apache.paimon.data.serializer.InternalRowSerializer;
import org.apache.paimon.schema.TableSchema;
import org.apache.paimon.types.DataField;
import org.apache.paimon.types.DataType;
@@ -367,4 +368,111 @@ public class RowConverterTest {
}
});
}
+
+ @Test
+ public void nestedRowPaimonToSeaTunnel() {
+ // Top-level: c_int INT, c_string STRING, c_nested ROW<a INT, b STRING>
+ RowType nestedRowType =
+ RowType.of(
+ new DataType[] {DataTypes.INT(), DataTypes.STRING()},
+ new String[] {"a", "b"});
+
+ RowType topLevelRowType =
+ RowType.of(
+ new DataType[] {DataTypes.INT(), DataTypes.STRING(),
nestedRowType},
+ new String[] {"c_int", "c_string", "c_nested"});
+
+ TableSchema tableSchema =
+ new TableSchema(
+ 0,
+ TableSchema.newFields(topLevelRowType),
+ topLevelRowType.getFieldCount(),
+ Collections.emptyList(),
+ Collections.emptyList(),
+ Collections.emptyMap(),
+ "");
+
+ // Build nested Paimon row
+ BinaryRow nestedBinaryRow = new BinaryRow(2);
+ BinaryRowWriter nestedWriter = new BinaryRowWriter(nestedBinaryRow);
+ nestedWriter.writeInt(0, 42);
+ nestedWriter.writeString(1, BinaryString.fromString("world"));
+
+ // Build top-level Paimon row
+ BinaryRow binaryRow = new BinaryRow(3);
+ BinaryRowWriter writer = new BinaryRowWriter(binaryRow);
+ writer.writeInt(0, 1);
+ writer.writeString(1, BinaryString.fromString("hello"));
+ writer.writeRow(2, nestedBinaryRow, new
InternalRowSerializer(nestedRowType));
+
+ // SeaTunnel types
+ SeaTunnelRowType seaTunnelNestedType =
+ new SeaTunnelRowType(
+ new String[] {"a", "b"},
+ new SeaTunnelDataType[] {BasicType.INT_TYPE,
BasicType.STRING_TYPE});
+ SeaTunnelRowType seaTunnelTopLevelType =
+ new SeaTunnelRowType(
+ new String[] {"c_int", "c_string", "c_nested"},
+ new SeaTunnelDataType[] {
+ BasicType.INT_TYPE, BasicType.STRING_TYPE,
seaTunnelNestedType
+ });
+
+ // Convert Paimon → SeaTunnel
+ SeaTunnelRow result = RowConverter.convert(binaryRow,
seaTunnelTopLevelType, tableSchema);
+
+ Assertions.assertEquals(1, result.getField(0));
+ Assertions.assertEquals("hello", result.getField(1));
+ SeaTunnelRow nestedResult = (SeaTunnelRow) result.getField(2);
+ Assertions.assertEquals(42, nestedResult.getField(0));
+ Assertions.assertEquals("world", nestedResult.getField(1));
+ }
+
+ @Test
+ public void nestedRowSeaTunnelToPaimon() {
+ // Top-level: c_int INT, c_string STRING, c_nested ROW<a INT, b STRING>
+ RowType nestedRowType =
+ RowType.of(
+ new DataType[] {DataTypes.INT(), DataTypes.STRING()},
+ new String[] {"a", "b"});
+
+ RowType topLevelRowType =
+ RowType.of(
+ new DataType[] {DataTypes.INT(), DataTypes.STRING(),
nestedRowType},
+ new String[] {"c_int", "c_string", "c_nested"});
+
+ TableSchema tableSchema =
+ new TableSchema(
+ 0,
+ TableSchema.newFields(topLevelRowType),
+ topLevelRowType.getFieldCount(),
+ Collections.emptyList(),
+ Collections.emptyList(),
+ Collections.emptyMap(),
+ "");
+
+ // SeaTunnel types
+ SeaTunnelRowType seaTunnelNestedType =
+ new SeaTunnelRowType(
+ new String[] {"a", "b"},
+ new SeaTunnelDataType[] {BasicType.INT_TYPE,
BasicType.STRING_TYPE});
+ SeaTunnelRowType seaTunnelTopLevelType =
+ new SeaTunnelRowType(
+ new String[] {"c_int", "c_string", "c_nested"},
+ new SeaTunnelDataType[] {
+ BasicType.INT_TYPE, BasicType.STRING_TYPE,
seaTunnelNestedType
+ });
+
+ SeaTunnelRow nestedRow = new SeaTunnelRow(new Object[] {42, "world"});
+ SeaTunnelRow topLevelRow = new SeaTunnelRow(new Object[] {1, "hello",
nestedRow});
+
+ // Convert SeaTunnel → Paimon (would throw before fix due to field
count mismatch)
+ InternalRow result =
+ RowConverter.reconvert(topLevelRow, seaTunnelTopLevelType,
tableSchema);
+
+ Assertions.assertEquals(1, result.getInt(0));
+ Assertions.assertEquals("hello", result.getString(1).toString());
+ InternalRow nestedResult = result.getRow(2, 2);
+ Assertions.assertEquals(42, nestedResult.getInt(0));
+ Assertions.assertEquals("world", nestedResult.getString(1).toString());
+ }
}