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]