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(

Reply via email to