This is an automated email from the ASF dual-hosted git repository.
danny0405 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/hudi.git
The following commit(s) were added to refs/heads/master by this push:
new 0425c5c458d6 fix(flink): preserve Avro fixed decimal widths in Parquet
writes (#19522)
0425c5c458d6 is described below
commit 0425c5c458d691198dfcdaa818ab2f066a691b2a
Author: Shuo Cheng <[email protected]>
AuthorDate: Thu Aug 6 17:18:13 2026 +0800
fix(flink): preserve Avro fixed decimal widths in Parquet writes (#19522)
* fix(flink): preserve Avro fixed decimal widths in Parquet writes
Honor declared Avro fixed sizes in both the Flink Parquet schema converter
and RowData value writer. Safely sign-extend compact decimals when the declared
fixed width exceeds eight bytes.
* style(flink): clarify decimal byte length resolver name
Rename the helper to describe both fixed-schema resolution and
precision-based fallback behavior. Addresses review comment 3718983684.
---
.../storage/row/parquet/ParquetRowDataWriter.java | 27 +++++++-------
.../row/parquet/ParquetSchemaConverter.java | 12 ++++++-
.../row/parquet/TestParquetRowDataWriter.java | 42 ++++++++++++++++++++++
.../row/parquet/TestParquetSchemaConverter.java | 15 ++++++++
4 files changed, 80 insertions(+), 16 deletions(-)
diff --git
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/storage/row/parquet/ParquetRowDataWriter.java
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/storage/row/parquet/ParquetRowDataWriter.java
index 4a3c13064685..09c7c8d6d625 100644
---
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/storage/row/parquet/ParquetRowDataWriter.java
+++
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/storage/row/parquet/ParquetRowDataWriter.java
@@ -46,7 +46,6 @@ import java.nio.ByteOrder;
import java.sql.Timestamp;
import java.util.Arrays;
-import static
org.apache.flink.formats.parquet.utils.ParquetSchemaConverter.computeMinBytesForDecimalPrecision;
import static
org.apache.flink.formats.parquet.vector.reader.TimestampColumnReader.JULIAN_EPOCH_OFFSET_DAYS;
import static
org.apache.flink.formats.parquet.vector.reader.TimestampColumnReader.MILLIS_IN_DAY;
import static
org.apache.flink.formats.parquet.vector.reader.TimestampColumnReader.NANOS_PER_MILLISECOND;
@@ -116,7 +115,7 @@ public class ParquetRowDataWriter {
return new BinaryWriter();
case DECIMAL:
DecimalType decimalType = (DecimalType) t;
- return createDecimalWriter(decimalType.getPrecision(),
decimalType.getScale());
+ return createDecimalWriter(decimalType.getPrecision(),
decimalType.getScale(), fieldSchema);
case TINYINT:
return new ByteWriter();
case SMALLINT:
@@ -426,24 +425,21 @@ public class ParquetRowDataWriter {
return Binary.fromConstantByteBuffer(buf);
}
- private FieldWriter createDecimalWriter(int precision, int scale) {
+ private FieldWriter createDecimalWriter(int precision, int scale,
HoodieSchema fieldSchema) {
Preconditions.checkArgument(
precision <= DecimalType.MAX_PRECISION,
"Decimal precision %s exceeds max precision %s",
precision,
DecimalType.MAX_PRECISION);
+ int numBytes =
ParquetSchemaConverter.resolveDecimalByteLength(fieldSchema, precision);
/*
* This is optimizer for UnscaledBytesWriter.
*/
class LongUnscaledBytesWriter implements FieldWriter {
- private final int numBytes;
- private final int initShift;
private final byte[] decimalBuffer;
private LongUnscaledBytesWriter() {
- this.numBytes = computeMinBytesForDecimalPrecision(precision);
- this.initShift = 8 * (numBytes - 1);
this.decimalBuffer = new byte[numBytes];
}
@@ -460,12 +456,15 @@ public class ParquetRowDataWriter {
}
private void doWrite(long unscaled) {
- int i = 0;
- int shift = initShift;
- while (i < numBytes) {
- decimalBuffer[i] = (byte) (unscaled >> shift);
- i += 1;
- shift -= 8;
+ // Parquet encodes FIXED_LEN_BYTE_ARRAY decimals as big-endian two's
complement. A compact
+ // Flink decimal provides at most eight value bytes, so pad wider Avro
fixed types with the
+ // sign byte to preserve the value.
+ int firstValueByte = Math.max(0, numBytes - Long.BYTES);
+ Arrays.fill(decimalBuffer, 0, firstValueByte, unscaled < 0 ? (byte) -1
: (byte) 0);
+ // Copy from the least-significant byte backwards to produce the
big-endian representation.
+ for (int i = numBytes - 1; i >= firstValueByte; i--) {
+ decimalBuffer[i] = (byte) unscaled;
+ unscaled >>= Byte.SIZE;
}
recordConsumer.addBinary(Binary.fromReusedByteArray(decimalBuffer, 0,
numBytes));
@@ -473,11 +472,9 @@ public class ParquetRowDataWriter {
}
class UnscaledBytesWriter implements FieldWriter {
- private final int numBytes;
private final byte[] decimalBuffer;
private UnscaledBytesWriter() {
- this.numBytes = computeMinBytesForDecimalPrecision(precision);
this.decimalBuffer = new byte[numBytes];
}
diff --git
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/storage/row/parquet/ParquetSchemaConverter.java
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/storage/row/parquet/ParquetSchemaConverter.java
index 5efe4aeee62c..85b65b60dc07 100644
---
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/storage/row/parquet/ParquetSchemaConverter.java
+++
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/storage/row/parquet/ParquetSchemaConverter.java
@@ -313,7 +313,7 @@ public class ParquetSchemaConverter {
case DECIMAL:
int precision = ((DecimalType) type).getPrecision();
int scale = ((DecimalType) type).getScale();
- int numBytes = computeMinBytesForDecimalPrecision(precision);
+ int numBytes = resolveDecimalByteLength(fieldSchema, precision);
return Types.primitive(
PrimitiveType.PrimitiveTypeName.FIXED_LEN_BYTE_ARRAY,
repetition)
.as(LogicalTypeAnnotation.decimalType(scale, precision))
@@ -447,4 +447,14 @@ public class ParquetSchemaConverter {
}
return numBytes;
}
+
+ static int resolveDecimalByteLength(HoodieSchema fieldSchema, int precision)
{
+ if (fieldSchema instanceof HoodieSchema.Decimal) {
+ HoodieSchema.Decimal decimalSchema = (HoodieSchema.Decimal) fieldSchema;
+ if (decimalSchema.isFixed()) {
+ return decimalSchema.getFixedSize();
+ }
+ }
+ return computeMinBytesForDecimalPrecision(precision);
+ }
}
diff --git
a/hudi-client/hudi-flink-client/src/test/java/org/apache/hudi/io/storage/row/parquet/TestParquetRowDataWriter.java
b/hudi-client/hudi-flink-client/src/test/java/org/apache/hudi/io/storage/row/parquet/TestParquetRowDataWriter.java
index 7d4c685c7a5f..35c45383fd65 100644
---
a/hudi-client/hudi-flink-client/src/test/java/org/apache/hudi/io/storage/row/parquet/TestParquetRowDataWriter.java
+++
b/hudi-client/hudi-flink-client/src/test/java/org/apache/hudi/io/storage/row/parquet/TestParquetRowDataWriter.java
@@ -29,14 +29,18 @@ import org.apache.flink.table.data.GenericRowData;
import org.apache.flink.table.data.StringData;
import org.apache.flink.table.data.TimestampData;
import org.apache.flink.table.types.logical.RowType;
+import org.apache.parquet.io.api.Binary;
import org.apache.parquet.io.api.RecordConsumer;
import org.junit.jupiter.api.Test;
+import org.mockito.ArgumentCaptor;
import java.math.BigDecimal;
import java.time.Instant;
+import java.util.Arrays;
import java.util.LinkedHashMap;
import java.util.Map;
+import static org.junit.jupiter.api.Assertions.assertArrayEquals;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyInt;
import static org.mockito.ArgumentMatchers.anyLong;
@@ -49,6 +53,37 @@ import static org.mockito.Mockito.verify;
class TestParquetRowDataWriter {
+ @Test
+ void testWriteDecimalWithWidthFromHoodieSchema() {
+ HoodieSchema schema = HoodieSchema.parse(
+ "{\"type\":\"record\",\"name\":\"rec\",\"fields\":["
+ +
"{\"name\":\"large_fixed\",\"type\":{\"type\":\"fixed\",\"name\":\"large_fixed_type\","
+ +
"\"size\":10,\"logicalType\":\"decimal\",\"precision\":20,\"scale\":2}},"
+ +
"{\"name\":\"small_fixed\",\"type\":{\"type\":\"fixed\",\"name\":\"small_fixed_type\","
+ +
"\"size\":10,\"logicalType\":\"decimal\",\"precision\":10,\"scale\":2}},"
+ +
"{\"name\":\"bytes_decimal\",\"type\":{\"type\":\"bytes\",\"logicalType\":\"decimal\","
+ + "\"precision\":20,\"scale\":2}}]}");
+ BigDecimal largeValue = new BigDecimal("123456789.12");
+ BigDecimal smallValue = new BigDecimal("-12.34");
+ BigDecimal bytesValue = new BigDecimal("223456789.34");
+ GenericRowData row = GenericRowData.of(
+ DecimalData.fromBigDecimal(largeValue, 20, 2),
+ DecimalData.fromBigDecimal(smallValue, 10, 2),
+ DecimalData.fromBigDecimal(bytesValue, 20, 2));
+ RecordConsumer consumer = mock(RecordConsumer.class);
+
+ new ParquetRowDataWriter(consumer, true, schema).write(row);
+
+ ArgumentCaptor<Binary> binaryCaptor =
ArgumentCaptor.forClass(Binary.class);
+ verify(consumer, times(3)).addBinary(binaryCaptor.capture());
+ assertArrayEquals(signExtend(largeValue.unscaledValue().toByteArray(), 10),
+ binaryCaptor.getAllValues().get(0).getBytes());
+ assertArrayEquals(signExtend(smallValue.unscaledValue().toByteArray(), 10),
+ binaryCaptor.getAllValues().get(1).getBytes());
+ assertArrayEquals(signExtend(bytesValue.unscaledValue().toByteArray(), 9),
+ binaryCaptor.getAllValues().get(2).getBytes());
+ }
+
@Test
void testWritePrimitiveNestedArrayMapDecimalAndTimestampValues() {
RowType rowType = (RowType) DataTypes.ROW(
@@ -173,4 +208,11 @@ class TestParquetRowDataWriter {
verify(consumer, atLeastOnce()).addDouble(3.5d);
verify(consumer, atLeastOnce()).addBinary(any());
}
+
+ private static byte[] signExtend(byte[] bytes, int length) {
+ byte[] result = new byte[length];
+ Arrays.fill(result, 0, length - bytes.length, bytes[0] < 0 ? (byte) -1 :
(byte) 0);
+ System.arraycopy(bytes, 0, result, length - bytes.length, bytes.length);
+ return result;
+ }
}
diff --git
a/hudi-client/hudi-flink-client/src/test/java/org/apache/hudi/io/storage/row/parquet/TestParquetSchemaConverter.java
b/hudi-client/hudi-flink-client/src/test/java/org/apache/hudi/io/storage/row/parquet/TestParquetSchemaConverter.java
index 962af8f6e04f..f33270136d7a 100644
---
a/hudi-client/hudi-flink-client/src/test/java/org/apache/hudi/io/storage/row/parquet/TestParquetSchemaConverter.java
+++
b/hudi-client/hudi-flink-client/src/test/java/org/apache/hudi/io/storage/row/parquet/TestParquetSchemaConverter.java
@@ -241,6 +241,21 @@ public class TestParquetSchemaConverter {
assertEquals(16, featuresType.getTypeLength());
}
+ @Test
+ void testDecimalFixedLenWidthFromHoodieSchema() {
+ HoodieSchema hoodieSchema = HoodieSchema.parse(
+ "{\"type\":\"record\",\"name\":\"rec\",\"fields\":["
+ +
"{\"name\":\"fixed_decimal\",\"type\":{\"type\":\"fixed\",\"name\":\"dec_fixed\","
+ +
"\"size\":10,\"logicalType\":\"decimal\",\"precision\":20,\"scale\":2}},"
+ +
"{\"name\":\"bytes_decimal\",\"type\":{\"type\":\"bytes\",\"logicalType\":\"decimal\","
+ + "\"precision\":20,\"scale\":2}}]}");
+
+ MessageType messageType =
ParquetSchemaConverter.convertToParquetMessageType("converted", hoodieSchema);
+
+ assertEquals(10,
messageType.getType("fixed_decimal").asPrimitiveType().getTypeLength());
+ assertEquals(9,
messageType.getType("bytes_decimal").asPrimitiveType().getTypeLength());
+ }
+
@Test
void testUnannotatedFixedLenByteArrayConvertsToBytes() {
MessageType messageType = new MessageType(