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


##########
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:
   Not introduced here: the base-file overlay has been unconditional since 
#18923, and this PR leaves the payload-based-with-logs arm as master reads it. 
Gating the base overlay alone would not fix it either, since the merged rows 
would then reach a plan expecting the projection struct as a `VariantVal`; that 
arm also needs the merged output projected (`buildVariantProjector`). Tracked 
on #17744 (https://github.com/apache/hudi/issues/17744#issuecomment-5826048476) 
rather than widening this PR.
   



##########
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:
   Done.
   



##########
hudi-client/hudi-spark-client/src/main/scala/org/apache/hudi/BaseSparkInternalRecordContext.java:
##########
@@ -214,10 +215,30 @@ public InternalRow toBinaryRow(HoodieSchema schema, 
InternalRow internalRow) {
     return unsafeProjection.apply(internalRow);
   }
 
+  /**
+   * Installs the Spark type the rows of an engine schema actually carry in 
this read; null, the default, means the
+   * plain conversion. {@code SparkFileFormatInternalRowReaderContext} 
installs the PushVariantIntoScan overlay here
+   * (see its setSchemaHandler), because every row writer this context builds 
from an engine schema has to be typed
+   * over that shape: a VariantType-typed writer re-encodes a projection 
struct through UnsafeRow.getVariant instead
+   * of copying it across.
+   */
+  public void setRowShape(UnaryOperator<StructType> rowShape) {
+    this.rowShape = rowShape;
+  }
+
+  /**
+   * The Spark type the rows of {@code schema} carry in this read: the plain 
conversion, or the row shape installed
+   * by {@link #setRowShape} when the reader hands its rows over in a 
rewritten shape.
+   */
+  public StructType getRowStructType(HoodieSchema schema) {
+    StructType structType = getCachedSchema(schema);
+    return rowShape == null ? structType : rowShape.apply(structType);
+  }
+
   @Override
   public UnaryOperator<InternalRow> projectRecord(HoodieSchema from, 
HoodieSchema to, Map<String, String> renamedColumns) {
     Function1<InternalRow, UnsafeRow> unsafeRowWriter =
-        HoodieInternalRowUtils.getCachedUnsafeRowWriter(getCachedSchema(from), 
getCachedSchema(to), renamedColumns, Collections.emptyMap());
+        
HoodieInternalRowUtils.getCachedUnsafeRowWriter(getRowStructType(from), 
getRowStructType(to), renamedColumns, Collections.emptyMap());

Review Comment:
   Yes. Partial updates are on by default for MOR `MERGE INTO`, and 
`mergePartialRecords` / `mergeWithEngineRecord` read the base record as 
`VariantType`.
   
   Fixed: `getValue`, `toBinaryRow` and `mergePartialRecords` now read through 
`getRowStructType`. The MOR `variant_get` leg adds a ts-only `MERGE INTO` over 
a base-file record (asserts the partial block) and reads it back on both 
`pushVariantIntoScan` arms.
   



-- 
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