voonhous commented on code in PR #20130:
URL: https://github.com/apache/hudi/pull/20130#discussion_r4135561945


##########
hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/common/TestVectorizedReadWithSchemaEvolution.scala:
##########
@@ -66,5 +66,44 @@ class TestVectorizedReadWithSchemaEvolution extends 
HoodieSparkSqlTestBase {
         }
       }
     }
+
+    test(s"Test INT to DECIMAL schema evolution with precision overflow for 
$tableType table") {
+      if (HoodieSparkUtils.isSpark3) {

Review Comment:
   **major:** `isSpark3` skips this test on every Spark 4 profile, and the 
GitHub Scala-test job that feeds codecov runs only the `spark4.2` matrix entry, 
which is why the patch shows 0% coverage (Azure runs it once on 3.5). The 
reader is shared: `HoodieVectorizedParquetRecordReader` lives in 
hudi-spark-common and Spark40/41/42ParquetReader all call 
`buildVectorizedReader`. The guard was inherited from #12560, before Spark 4 
support. Could we drop it here (and on the test above at line 25)?



##########
hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/client/utils/SparkInternalSchemaConverter.java:
##########
@@ -349,9 +349,13 @@ private static boolean 
convertIntLongType(WritableColumnVector oldV, WritableCol
         } else if (newType instanceof StringType) {
           newV.putByteArray(i, getUTF8Bytes((isInt ? oldV.getInt(i) : 
oldV.getLong(i)) + ""));
         } else if (newType instanceof DecimalType) {
+          DecimalType decimalType = (DecimalType) newType;
           Decimal oldDecimal = Decimal.apply(isInt ? oldV.getInt(i) : 
oldV.getLong(i));
-          oldDecimal.changePrecision(((DecimalType) newType).precision(), 
((DecimalType) newType).scale());
-          newV.putDecimal(i, oldDecimal, ((DecimalType) newType).precision());
+          if (oldDecimal.changePrecision(decimalType.precision(), 
decimalType.scale())) {
+            newV.putDecimal(i, oldDecimal, decimalType.precision());
+          } else {
+            newV.putNull(i);

Review Comment:
   **major:** After this change the vectorized path returns NULL on overflow 
whatever `spark.sql.ansi.enabled` is, but the row path casts through Spark 
`Cast` (`SparkSchemaTransformUtils.recursivelyCastExpressions`: 
`Cast(Cast(expr, String), dec)` with the session eval mode) and throws under 
ANSI. Spark 4 defaults ANSI on, so a COW vectorized read returns NULL where a 
MOR or `enableVectorizedReader=false` read of the same file throws. Could we 
throw here when `SQLConf.get().ansiEnabled()` is true (it resolves on executors 
via the task context) and add an ANSI-on case expecting the exception?



##########
hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/common/TestVectorizedReadWithSchemaEvolution.scala:
##########
@@ -66,5 +66,44 @@ class TestVectorizedReadWithSchemaEvolution extends 
HoodieSparkSqlTestBase {
         }
       }
     }
+
+    test(s"Test INT to DECIMAL schema evolution with precision overflow for 
$tableType table") {
+      if (HoodieSparkUtils.isSpark3) {
+        withSQLConf(
+          "hoodie.schema.on.read.enable" -> "true",
+          "spark.sql.ansi.enabled" -> "false",
+          "spark.sql.parquet.enableVectorizedReader" -> "true",
+          "hoodie.parquet.small.file.limit" -> "0"
+        ) {
+          withTempDir { tmp =>
+            val tableName = generateTableName
+            val tablePath = s"${tmp.getCanonicalPath}/$tableName"
+
+            spark.sql(
+              s"""
+                 |create table $tableName (
+                 |  id int,
+                 |  price int,
+                 |  ts long
+                 |) using hudi
+                 | location '$tablePath'
+                 | tblproperties (
+                 |  type = '$tableType',
+                 |  primaryKey = 'id',
+                 |  orderingFields = 'ts'
+                 | )
+                 |""".stripMargin)
+
+            spark.sql(s"insert into $tableName values (1, 12345, 1000)")

Review Comment:
   **minor:** (not blocking) A single overflowing INT row leaves two cheap 
cases uncovered: a value that fits in the same batch, e.g. `(2, 12, 1000)` 
expecting `12.00`, pins `putNull` and `putDecimal` interleaving in one vector, 
and a `bigint` column with `9223372036854775807 -> decimal(20, 2)` exercises 
the LONG source and the byte-array-backed vector (precision > 18). Could we add 
both here?



##########
hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/common/TestVectorizedReadWithSchemaEvolution.scala:
##########
@@ -66,5 +66,44 @@ class TestVectorizedReadWithSchemaEvolution extends 
HoodieSparkSqlTestBase {
         }
       }
     }
+
+    test(s"Test INT to DECIMAL schema evolution with precision overflow for 
$tableType table") {

Review Comment:
   **minor:** (not blocking) The MOR leg never reaches the changed code: 
`HoodieFileGroupReaderBasedFileFormat` builds the MOR base-file reader with 
`enableVectorizedRead = false`, so MOR slices take the row path, where the 
Spark `Cast` already returns NULL with ANSI off. This leg passes with or 
without the fix. Was the MOR `123.45` in the description observed under a 
different config? If not, could we keep the MOR leg as an explicit row-path 
cross-check and name it so, or drop it?



##########
hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/client/utils/SparkInternalSchemaConverter.java:
##########
@@ -349,9 +349,13 @@ private static boolean 
convertIntLongType(WritableColumnVector oldV, WritableCol
         } else if (newType instanceof StringType) {
           newV.putByteArray(i, getUTF8Bytes((isInt ? oldV.getInt(i) : 
oldV.getLong(i)) + ""));
         } else if (newType instanceof DecimalType) {
+          DecimalType decimalType = (DecimalType) newType;
           Decimal oldDecimal = Decimal.apply(isInt ? oldV.getInt(i) : 
oldV.getLong(i));

Review Comment:
   **nit:** (feel free to ignore) The overflow reaches this line only because 
`SchemaChangeUtils.isTypeUpdateAllowInternal` accepts INT or LONG to any 
`DECIMAL(p, s)` without checking `p - s >= 10` (INT) or `>= 19` (LONG); 
`decimal(4, 2)` can never hold an int. Would a follow-up that rejects the lossy 
DDL at the source be worth filing?



##########
hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/client/utils/SparkInternalSchemaConverter.java:
##########
@@ -349,9 +349,13 @@ private static boolean 
convertIntLongType(WritableColumnVector oldV, WritableCol
         } else if (newType instanceof StringType) {
           newV.putByteArray(i, getUTF8Bytes((isInt ? oldV.getInt(i) : 
oldV.getLong(i)) + ""));
         } else if (newType instanceof DecimalType) {
+          DecimalType decimalType = (DecimalType) newType;

Review Comment:
   **minor:** (not blocking) `TestSparkInternalSchemaConverter.scala` in 
hudi-spark already exists with one test and never calls 
`convertColumnVectorType`. A direct test there, filling an 
`OnHeapColumnVector(IntegerType)` with `12345`, converting into 
`OnHeapColumnVector(DecimalType(4, 2))` and asserting `isNullAt(0)`, would pin 
int, long, float, double and string overflow in one place with no Spark SQL and 
no version guard. Would it be worth adding that alongside the SQL test?



##########
hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/client/utils/SparkInternalSchemaConverter.java:
##########
@@ -349,9 +349,13 @@ private static boolean 
convertIntLongType(WritableColumnVector oldV, WritableCol
         } else if (newType instanceof StringType) {
           newV.putByteArray(i, getUTF8Bytes((isInt ? oldV.getInt(i) : 
oldV.getLong(i)) + ""));
         } else if (newType instanceof DecimalType) {
+          DecimalType decimalType = (DecimalType) newType;
           Decimal oldDecimal = Decimal.apply(isInt ? oldV.getInt(i) : 
oldV.getLong(i));
-          oldDecimal.changePrecision(((DecimalType) newType).precision(), 
((DecimalType) newType).scale());
-          newV.putDecimal(i, oldDecimal, ((DecimalType) newType).precision());
+          if (oldDecimal.changePrecision(decimalType.precision(), 
decimalType.scale())) {
+            newV.putDecimal(i, oldDecimal, decimalType.precision());

Review Comment:
   **minor:** (not blocking) The description says "Closes #20108", but that 
issue also asks for the types this converter still returns `false` for 
(Boolean, Byte, Short, Binary, Timestamp) to be routed to the row reader, which 
the description itself defers. Merging as-is would close the issue with that 
half open. Could we change it to "Part of #20108", or add the fallback here?



##########
hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/client/utils/SparkInternalSchemaConverter.java:
##########
@@ -349,9 +349,13 @@ private static boolean 
convertIntLongType(WritableColumnVector oldV, WritableCol
         } else if (newType instanceof StringType) {
           newV.putByteArray(i, getUTF8Bytes((isInt ? oldV.getInt(i) : 
oldV.getLong(i)) + ""));
         } else if (newType instanceof DecimalType) {
+          DecimalType decimalType = (DecimalType) newType;
           Decimal oldDecimal = Decimal.apply(isInt ? oldV.getInt(i) : 
oldV.getLong(i));
-          oldDecimal.changePrecision(((DecimalType) newType).precision(), 
((DecimalType) newType).scale());
-          newV.putDecimal(i, oldDecimal, ((DecimalType) newType).precision());
+          if (oldDecimal.changePrecision(decimalType.precision(), 
decimalType.scale())) {

Review Comment:
   **major:** Adding the mechanism behind this: a failed `changePrecision` 
leaves the `Decimal` untouched, and `putDecimal` then writes its unscaled long 
into the narrower vector, so double `12345.6` read as `decimal(4, 2)` comes 
back as `1234.56`, a wrong value rather than a rescaled one. 
`SchemaChangeUtils` allows FLOAT, DOUBLE and STRING to any decimal, so all 
three are reachable via DDL; DECIMAL to DECIMAL is widening-only, so the check 
at line 435 cannot fail. Could we pull the check into a `putDecimalOrNull(newV, 
i, decimal, decimalType)` helper used at the three sites?



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to