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 406c66789e [Fix][Connector-V2][File] Fix Parquet INT96 writes for
uppercase fields (#10943)
406c66789e is described below
commit 406c66789ec297e1bdde1b6f61875d61096f2355
Author: Asish Kumar <[email protected]>
AuthorDate: Mon Aug 24 08:31:22 2026 +0530
[Fix][Connector-V2][File] Fix Parquet INT96 writes for uppercase fields
(#10943)
Co-authored-by: davidzollo <[email protected]>
Co-authored-by: davidzollo <[email protected]>
---
.../file/sink/writer/ParquetWriteStrategy.java | 25 ++++++++++++----------
.../file/writer/ParquetWriteStrategyTest.java | 20 +++++++++++++----
2 files changed, 30 insertions(+), 15 deletions(-)
diff --git
a/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/sink/writer/ParquetWriteStrategy.java
b/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/sink/writer/ParquetWriteStrategy.java
index e72436338f..53bd50692e 100644
---
a/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/sink/writer/ParquetWriteStrategy.java
+++
b/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/sink/writer/ParquetWriteStrategy.java
@@ -101,7 +101,7 @@ public class ParquetWriteStrategy extends
AbstractWriteStrategy<ParquetWriter<Ge
Configuration configuration = getConfiguration(hadoopConf);
writePathsAsInt96 =
fileSinkConfig.getParquetAvroWriteFixedAsInt96().stream()
- .map(this::normalizeFieldName)
+ .map(ParquetWriteStrategy::normalizeFieldName)
.collect(Collectors.toSet());
if (fileSinkConfig.getParquetWriteTimestampAsInt96()) {
List<String> timestampFields = new ArrayList<>();
@@ -127,10 +127,11 @@ public class ParquetWriteStrategy extends
AbstractWriteStrategy<ParquetWriter<Ge
GenericRecordBuilder recordBuilder = new GenericRecordBuilder(schema);
for (Integer integer : sinkColumnsIndexInRow) {
String fieldName = seaTunnelRowType.getFieldName(integer);
+ String parquetFieldName = normalizeFieldName(fieldName);
Object field = getFieldSafe(seaTunnelRow, integer);
recordBuilder.set(
- fieldName.toLowerCase(),
- resolveObject(fieldName, field,
seaTunnelRowType.getFieldType(integer)));
+ parquetFieldName,
+ resolveObject(parquetFieldName, field,
seaTunnelRowType.getFieldType(integer)));
}
GenericData.Record record = recordBuilder.build();
try {
@@ -306,9 +307,11 @@ public class ParquetWriteStrategy extends
AbstractWriteStrategy<ParquetWriter<Ge
(SeaTunnelRowType) seaTunnelDataType,
sinkColumnsIndex);
GenericRecordBuilder recordBuilder = new
GenericRecordBuilder(recordSchema);
for (int i = 0; i < fieldNames.length; i++) {
+ String parquetFieldName =
normalizeFieldName(fieldNames[i]);
recordBuilder.set(
- fieldNames[i].toLowerCase(),
- resolveObject(fieldNames[i],
seaTunnelRow.getField(i), fieldTypes[i]));
+ parquetFieldName,
+ resolveObject(
+ parquetFieldName,
seaTunnelRow.getField(i), fieldTypes[i]));
}
return recordBuilder.build();
default:
@@ -321,6 +324,10 @@ public class ParquetWriteStrategy extends
AbstractWriteStrategy<ParquetWriter<Ge
}
}
+ private static String normalizeFieldName(String fieldName) {
+ return fieldName.toLowerCase(Locale.ROOT);
+ }
+
public Type seaTunnelDataType2ParquetDataType(
String fieldName, SeaTunnelDataType<?> seaTunnelDataType) {
switch (seaTunnelDataType.getSqlType()) {
@@ -447,7 +454,7 @@ public class ParquetWriteStrategy extends
AbstractWriteStrategy<ParquetWriter<Ge
if (fileSinkConfig.getParquetWriteTimestampAsInt96()) {
writePathsAsInt96 =
fileSinkConfig.getParquetAvroWriteFixedAsInt96().stream()
- .map(this::normalizeFieldName)
+ .map(ParquetWriteStrategy::normalizeFieldName)
.collect(Collectors.toSet());
for (int i = 0; i < seaTunnelRowType.getTotalFields(); i++) {
if
(SqlType.TIMESTAMP.equals(seaTunnelRowType.getFieldType(i).getSqlType())) {
@@ -468,15 +475,11 @@ public class ParquetWriteStrategy extends
AbstractWriteStrategy<ParquetWriter<Ge
index -> {
Type type =
seaTunnelDataType2ParquetDataType(
- fieldNames[index].toLowerCase(),
fieldTypes[index]);
+ normalizeFieldName(fieldNames[index]),
fieldTypes[index]);
types.add(type);
});
MessageType seaTunnelRow =
Types.buildMessage().addFields(types.toArray(new
Type[0])).named("SeaTunnelRecord");
return schemaConverter.convert(seaTunnelRow);
}
-
- private String normalizeFieldName(String fieldName) {
- return fieldName.toLowerCase(Locale.ROOT);
- }
}
diff --git
a/seatunnel-connectors-v2/connector-file/connector-file-base/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/writer/ParquetWriteStrategyTest.java
b/seatunnel-connectors-v2/connector-file/connector-file-base/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/writer/ParquetWriteStrategyTest.java
index 7c1a71eef2..3599a06392 100644
---
a/seatunnel-connectors-v2/connector-file/connector-file-base/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/writer/ParquetWriteStrategyTest.java
+++
b/seatunnel-connectors-v2/connector-file/connector-file-base/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/writer/ParquetWriteStrategyTest.java
@@ -111,15 +111,16 @@ public class ParquetWriteStrategyTest {
writeConfig.put("path", "file:///tmp/seatunnel/parquet/int96");
writeConfig.put("file_format_type", FileFormat.PARQUET.name());
writeConfig.put("parquet_avro_write_timestamp_as_int96", "true");
- writeConfig.put("parquet_avro_write_fixed_as_int96",
Arrays.asList("f3_bytes"));
+ writeConfig.put("parquet_avro_write_fixed_as_int96",
Arrays.asList("F3_Bytes"));
SeaTunnelRowType writeRowType =
new SeaTunnelRowType(
- new String[] {"f1_text", "f2_timestamp", "f3_bytes"},
+ new String[] {"f1_text", "f2_timestamp", "F3_Bytes",
"createTime"},
new SeaTunnelDataType[] {
BasicType.STRING_TYPE,
LocalTimeType.LOCAL_DATE_TIME_TYPE,
- PrimitiveByteArrayType.INSTANCE
+ PrimitiveByteArrayType.INSTANCE,
+ LocalTimeType.LOCAL_DATE_TIME_TYPE
});
FileSinkConfig writeSinkConfig =
new FileSinkConfig(ReadonlyConfig.fromMap(writeConfig),
writeRowType);
@@ -131,7 +132,10 @@ public class ParquetWriteStrategyTest {
writeStrategy.init(hadoopConf, "test1", "test1", 0);
writeStrategy.beginTransaction(1L);
writeStrategy.write(
- new SeaTunnelRow(new Object[] {"test", LocalDateTime.now(),
new byte[12]}));
+ new SeaTunnelRow(
+ new Object[] {
+ "test", LocalDateTime.now(), new byte[12],
LocalDateTime.now()
+ }));
writeStrategy.finishAndCloseFile();
writeStrategy.close();
@@ -161,6 +165,10 @@ public class ParquetWriteStrategyTest {
Assertions.assertEquals(
PrimitiveType.PrimitiveTypeName.INT96,
f3Type.asPrimitiveType().getPrimitiveTypeName());
+ Type createTimeType = metadata.getSchema().getType("createtime");
+ Assertions.assertEquals(
+ PrimitiveType.PrimitiveTypeName.INT96,
+ createTimeType.asPrimitiveType().getPrimitiveTypeName());
}
SeaTunnelRowType readRowType =
readStrategy.getSeaTunnelRowTypeInfo(readFilePath);
@@ -172,6 +180,9 @@ public class ParquetWriteStrategyTest {
Assertions.assertEquals(
LocalTimeType.LOCAL_DATE_TIME_TYPE.getSqlType(),
readRowType.getFieldType(2).getSqlType());
+ Assertions.assertEquals(
+ LocalTimeType.LOCAL_DATE_TIME_TYPE.getSqlType(),
+ readRowType.getFieldType(3).getSqlType());
List<SeaTunnelRow> readRows = new ArrayList<>();
Collector<SeaTunnelRow> readCollector =
new Collector<SeaTunnelRow>() {
@@ -180,6 +191,7 @@ public class ParquetWriteStrategyTest {
Assertions.assertTrue(record.getField(0) instanceof
String);
Assertions.assertTrue(record.getField(1) instanceof
LocalDateTime);
Assertions.assertTrue(record.getField(2) instanceof
LocalDateTime);
+ Assertions.assertTrue(record.getField(3) instanceof
LocalDateTime);
readRows.add(record);
}