hudi-agent commented on code in PR #20041:
URL: https://github.com/apache/hudi/pull/20041#discussion_r4084808297
##########
hudi-client/hudi-spark-client/src/main/scala/org/apache/hudi/SparkFileFormatInternalRowReaderContext.scala:
##########
@@ -124,20 +127,45 @@ class
SparkFileFormatInternalRowReaderContext(baseFileReader: SparkColumnarFileR
})
}
+ // Whether the query carries a Spark 4.1 PushVariantIntoScan projection at
all.
+ private lazy val hasVariantProjection: Boolean =
+ sparkRequiredSchema.exists(_.fields.exists(f =>
sparkAdapter.containsVariantProjection(f.dataType)))
+
+ private def isPayloadBasedMerge: Boolean = {
+ // getRecordMerger() is a Lombok getter over a field initialized to null
(not Option.empty());
+ // it stays null until HoodieReaderContext.initRecordMerger runs
(HoodieFileGroupReader calls it
+ // from its constructor), so the null guard is required.
+ val merger = getRecordMerger()
+ merger != null && merger.isPresent && merger.get.getMergingStrategy ==
PAYLOAD_BASED_MERGE_STRATEGY_UUID
+ }
+
// True only when there is a Spark 4.1 PushVariantIntoScan projection to
apply AND the table is
// not using a custom (payload-based) merger. Payload-based tables
round-trip records through
// PayloadUpdateProcessor.convertToAvroRecord against a schema that still
types variant fields as
// VariantType, so a row already rewritten into the projected struct shape
would be mis-decoded.
// Single source of truth for both reader paths (parquet native projection +
avro rewrite).
- private def shouldProjectVariants(): Boolean = {
- val hasVariantProjection =
- sparkRequiredSchema.exists(_.fields.exists(f =>
sparkAdapter.containsVariantProjection(f.dataType)))
- // getRecordMerger() is a Lombok getter over a field initialized to null
(not Option.empty());
- // it stays null until HoodieReaderContext.initRecordMerger runs
(HoodieFileGroupReader calls it
- // from its constructor), so the null guard is required.
- val merger = getRecordMerger()
- val isPayloadBased = merger != null && merger.isPresent &&
merger.get.getMergingStrategy == PAYLOAD_BASED_MERGE_STRATEGY_UUID
- hasVariantProjection && !isPayloadBased
+ private def shouldProjectVariants(): Boolean = hasVariantProjection &&
!isPayloadBasedMerge
+
+ // HoodieFileGroupReader installs the schema handler after initRecordMerger
and before it asks for
+ // the output converter, so this is where the record context learns what
shape its rows carry.
+ // Base-file rows are ALWAYS read in the projected shape
(getFileRecordIterator overlays it
+ // unconditionally); log rows only when shouldProjectVariants rewrites them,
so the payload-based
+ // exclusion applies only when there are log files to merge. Two consumers
need the shape: the
+ // output converter, which projects the reader's required schema down to the
requested one
+ // (FileGroupReaderSchemaHandler.getOutputConverter), and the bootstrap
skeleton/data join
+ // (getBootstrapProjection). Without it both build a VariantType-typed row
writer that re-encodes
+ // the projection struct through UnsafeRow.getVariant - byte-identical only
while the struct's
+ // null bitset stays small, and a NegativeArraySizeException once enough
pushed fields are null.
+ override def setSchemaHandler(schemaHandler:
FileGroupReaderSchemaHandler[InternalRow]): Unit = {
+ super.setSchemaHandler(schemaHandler)
+ if (hasVariantProjection && (!isPayloadBasedMerge || !getHasLogFiles)) {
Review Comment:
🤖 For a payload-based table that *does* have log files this skips the row
shape, but `getFileRecordIterator` still overlays the projection onto the
base-file read unconditionally (the `sparkRequiredSchema match` above is not
gated on `shouldProjectVariants`), so base rows are projection structs while
the output converter stays VariantType-typed — the same
`NegativeArraySizeException` this PR fixes elsewhere, plus a base/log shape
mismatch inside the buffer. Would gating the base-file overlay on
`shouldProjectVariants()` too make the payload-based arm consistently
unprojected?
<sub><i>⚠️ AI-generated; verify before applying. React 👍/👎 to flag
quality.</i></sub>
##########
hudi-client/hudi-spark-client/src/main/scala/org/apache/hudi/SparkFileFormatInternalRowReaderContext.scala:
##########
@@ -124,20 +127,45 @@ class
SparkFileFormatInternalRowReaderContext(baseFileReader: SparkColumnarFileR
})
}
+ // Whether the query carries a Spark 4.1 PushVariantIntoScan projection at
all.
+ private lazy val hasVariantProjection: Boolean =
+ sparkRequiredSchema.exists(_.fields.exists(f =>
sparkAdapter.containsVariantProjection(f.dataType)))
+
+ private def isPayloadBasedMerge: Boolean = {
+ // getRecordMerger() is a Lombok getter over a field initialized to null
(not Option.empty());
+ // it stays null until HoodieReaderContext.initRecordMerger runs
(HoodieFileGroupReader calls it
+ // from its constructor), so the null guard is required.
+ val merger = getRecordMerger()
+ merger != null && merger.isPresent && merger.get.getMergingStrategy ==
PAYLOAD_BASED_MERGE_STRATEGY_UUID
+ }
+
// True only when there is a Spark 4.1 PushVariantIntoScan projection to
apply AND the table is
// not using a custom (payload-based) merger. Payload-based tables
round-trip records through
// PayloadUpdateProcessor.convertToAvroRecord against a schema that still
types variant fields as
// VariantType, so a row already rewritten into the projected struct shape
would be mis-decoded.
// Single source of truth for both reader paths (parquet native projection +
avro rewrite).
- private def shouldProjectVariants(): Boolean = {
- val hasVariantProjection =
- sparkRequiredSchema.exists(_.fields.exists(f =>
sparkAdapter.containsVariantProjection(f.dataType)))
- // getRecordMerger() is a Lombok getter over a field initialized to null
(not Option.empty());
- // it stays null until HoodieReaderContext.initRecordMerger runs
(HoodieFileGroupReader calls it
- // from its constructor), so the null guard is required.
- val merger = getRecordMerger()
- val isPayloadBased = merger != null && merger.isPresent &&
merger.get.getMergingStrategy == PAYLOAD_BASED_MERGE_STRATEGY_UUID
- hasVariantProjection && !isPayloadBased
+ private def shouldProjectVariants(): Boolean = hasVariantProjection &&
!isPayloadBasedMerge
+
+ // HoodieFileGroupReader installs the schema handler after initRecordMerger
and before it asks for
+ // the output converter, so this is where the record context learns what
shape its rows carry.
+ // Base-file rows are ALWAYS read in the projected shape
(getFileRecordIterator overlays it
+ // unconditionally); log rows only when shouldProjectVariants rewrites them,
so the payload-based
+ // exclusion applies only when there are log files to merge. Two consumers
need the shape: the
+ // output converter, which projects the reader's required schema down to the
requested one
+ // (FileGroupReaderSchemaHandler.getOutputConverter), and the bootstrap
skeleton/data join
+ // (getBootstrapProjection). Without it both build a VariantType-typed row
writer that re-encodes
+ // the projection struct through UnsafeRow.getVariant - byte-identical only
while the struct's
+ // null bitset stays small, and a NegativeArraySizeException once enough
pushed fields are null.
+ override def setSchemaHandler(schemaHandler:
FileGroupReaderSchemaHandler[InternalRow]): Unit = {
+ super.setSchemaHandler(schemaHandler)
+ if (hasVariantProjection && (!isPayloadBasedMerge || !getHasLogFiles)) {
+ val requiredStruct = sparkRequiredSchema.get
+ recordContext.asInstanceOf[BaseSparkInternalRecordContext].setRowShape(
+ new UnaryOperator[StructType] {
+ override def apply(structType: StructType): StructType =
Review Comment:
🤖 nit: since this is Scala, could this be a plain lambda
(`overlayVariantProjections(_, requiredStruct)`) instead of instantiating an
anonymous `java.util.function.UnaryOperator`? It's already imported as a Java
functional interface in the signature, but the implementation doesn't need the
Java SAM boilerplate.
<sub><i>⚠️ AI-generated; verify before applying. React 👍/👎 to flag
quality.</i></sub>
--
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]