This is an automated email from the ASF dual-hosted git repository.

voonhous 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 76be6c65b2c1 fix(variant): row writers over the projected shape 
(#20041)
76be6c65b2c1 is described below

commit 76be6c65b2c1fd0bc74379dc1f30a2181bd2b8b9
Author: voonhous <[email protected]>
AuthorDate: Mon Sep 28 17:10:53 2026 +0800

    fix(variant): row writers over the projected shape (#20041)
    
    * fix(variant): row writers over the projected shape
    
    Spark 4.1+ PushVariantIntoScan rewrites a variant reached by
    variant_get, cast or is null into a projection struct for every
    HadoopFsRelation over a ParquetFileFormat, which is every Hudi read
    relation. SparkFileFormatInternalRowReaderContext reads base-file rows
    in that shape and rewrites log rows into it, but two row writers were
    still built from the engine HoodieSchema, whose variant field converts
    to VariantType: the output converter that projects the reader's
    required schema down to the requested one (merge columns on MOR,
    _hoodie_commit_time on incremental) and the bootstrap skeleton/data
    join. A VariantType-typed writer copies the struct through
    UnsafeRow.getVariant, which reads the struct's null bitset as the
    variant value length. That is byte-identical while the low bitset word
    stays below the struct size, so small projections passed, and throws
    NegativeArraySizeException once enough pushed fields are null: eight
    extracted paths of which a row holds one fail on any MOR read with a
    log and on any bootstrap read.
    
    - BaseSparkInternalRecordContext: setRowShape / getRowStructType; the
      row writers of projectRecord and getBootstrapProjection are typed
      over that shape.
    - SparkFileFormatInternalRowReaderContext.setSchemaHandler installs
      the PushVariantIntoScan overlay once the merger is known; base rows
      always carry it, log rows only when shouldProjectVariants rewrites
      them, so the payload-based exclusion applies only with log files.
    
    Tests in TestVariantShreddingMixedLayouts, both pushVariantIntoScan
    arms, each asserting whether the scan carries the projection struct:
    time travel at the base instant, incremental V2 over the timeline,
    incremental V1 selected through
    hoodie.datasource.read.incr.table.version and bounded at the base
    instant's requested time (the zero-row V2 read of that bound proves the
    option took effect), and a METADATA_ONLY
    bootstrap over a Spark-written variant parquet on COW and MOR. Both
    carry the eight-path query that failed before the fix. CDC is untouched
    by the rule (its relation schema is four strings) and the legacy
    streaming RDD path is not a HadoopFsRelation; both stay pinned by their
    existing tests. hoodie.datasource.read.use.new.parquet.file.format no
    longer exists.
    
    The bootstrap legs run on the default record type: a METADATA_ONLY
    bootstrap on the SPARK record type fails in the skeleton-file write
    (HUDI-5807), independently of variants.
    
    * fix(variant): shape-aware accessors for partial merges
    
    A partial MERGE INTO on MOR rebuilds the row field by field through
    SparkRecordMergingUtils.mergePartialRecords and
    mergeWithEngineRecord (via getValue). Both read the base-file record
    as VariantType, so under a PushVariantIntoScan projection they put a
    VariantVal where the now shape-typed output converter expects the
    projection struct.
    
    - getValue, toBinaryRow and mergePartialRecords read through
      getRowStructType; the shaped types are cached per engine schema.
    - setRowShape takes a Scala lambda instead of an anonymous
      UnaryOperator.
    - Extend the MOR variant_get leg with a ts-only MERGE INTO that
      writes a partial log block and reads id 4 back on both arms.
---
 .../apache/hudi/merge/SparkRecordMergingUtils.java |  12 +-
 .../hudi/BaseSparkInternalRecordContext.java       |  44 +++-
 .../hudi/BaseSparkInternalRowReaderContext.java    |   7 +-
 .../SparkFileFormatInternalRowReaderContext.scala  |  47 ++++-
 .../schema/TestVariantShreddingMixedLayouts.scala  | 226 ++++++++++++++++++++-
 .../dml/schema/VariantShreddingTestSupport.scala   |  19 +-
 6 files changed, 332 insertions(+), 23 deletions(-)

diff --git 
a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/merge/SparkRecordMergingUtils.java
 
b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/merge/SparkRecordMergingUtils.java
index 29408c131dfa..ebb91f48d717 100644
--- 
a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/merge/SparkRecordMergingUtils.java
+++ 
b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/merge/SparkRecordMergingUtils.java
@@ -19,6 +19,7 @@
 
 package org.apache.hudi.merge;
 
+import org.apache.hudi.BaseSparkInternalRecordContext;
 import org.apache.hudi.HoodieSchemaConversionUtils;
 import org.apache.hudi.common.engine.RecordContext;
 import org.apache.hudi.common.model.HoodieRecordMerger;
@@ -32,6 +33,7 @@ import org.apache.hudi.common.util.collection.Pair;
 import org.apache.spark.sql.HoodieInternalRowUtils;
 import org.apache.spark.sql.catalyst.InternalRow;
 import org.apache.spark.sql.catalyst.expressions.GenericInternalRow;
+import org.apache.spark.sql.types.DataType;
 import org.apache.spark.sql.types.StructField;
 import org.apache.spark.sql.types.StructType;
 
@@ -119,16 +121,22 @@ public class SparkRecordMergingUtils {
       Map<Integer, StructField> mergedIdToFieldMapping = 
mergedSchemaPair.getLeft();
       Map<String, Integer> oldNameToIdMapping = 
getCachedFieldNameToIdMapping(oldSchema);
       Map<String, Integer> newPartialNameToIdMapping = 
getCachedFieldNameToIdMapping(newSchema);
+      // The values are read in the shape the rows carry: under a 
PushVariantIntoScan projection a variant is the
+      // projection struct, which a VariantType read would decode through 
UnsafeRow.getVariant.
+      StructType mergedRowStruct = recordContext instanceof 
BaseSparkInternalRecordContext
+          ? ((BaseSparkInternalRecordContext) 
recordContext).getRowStructType(mergedSchemaPair.getRight().getRight())
+          : mergedSchemaPair.getRight().getLeft();
       List<Object> values = new ArrayList<>(mergedIdToFieldMapping.size());
       for (int fieldId = 0; fieldId < mergedIdToFieldMapping.size(); 
fieldId++) {
         StructField structField = mergedIdToFieldMapping.get(fieldId);
+        DataType dataType = mergedRowStruct.fields()[fieldId].dataType();
         Integer ordInPartialUpdate = 
newPartialNameToIdMapping.get(structField.name());
         if (ordInPartialUpdate != null) {
           // The field exists in the newer record; picks the value from newer 
record
-          values.add(newPartialRow.get(ordInPartialUpdate, 
structField.dataType()));
+          values.add(newPartialRow.get(ordInPartialUpdate, dataType));
         } else {
           // The field does not exist in the newer record; picks the value 
from older record
-          values.add(oldRow.get(oldNameToIdMapping.get(structField.name()), 
structField.dataType()));
+          values.add(oldRow.get(oldNameToIdMapping.get(structField.name()), 
dataType));
         }
       }
       InternalRow mergedRow = new GenericInternalRow(values.toArray());
diff --git 
a/hudi-client/hudi-spark-client/src/main/scala/org/apache/hudi/BaseSparkInternalRecordContext.java
 
b/hudi-client/hudi-spark-client/src/main/scala/org/apache/hudi/BaseSparkInternalRecordContext.java
index 1ed7f804503c..7b236d3419ed 100644
--- 
a/hudi-client/hudi-spark-client/src/main/scala/org/apache/hudi/BaseSparkInternalRecordContext.java
+++ 
b/hudi-client/hudi-spark-client/src/main/scala/org/apache/hudi/BaseSparkInternalRecordContext.java
@@ -46,6 +46,7 @@ import java.nio.ByteBuffer;
 import java.util.Collections;
 import java.util.List;
 import java.util.Map;
+import java.util.concurrent.ConcurrentHashMap;
 import java.util.function.UnaryOperator;
 
 import scala.Function1;
@@ -55,6 +56,9 @@ import static 
org.apache.spark.sql.HoodieInternalRowUtils.getCachedSchema;
 public abstract class BaseSparkInternalRecordContext extends 
RecordContext<InternalRow> {
 
   private OrderingValueEngineTypeConverter orderingValueConverter;
+  private UnaryOperator<StructType> rowShape;
+  // The row shape applied per engine schema, so the per-record accessors do 
not rebuild the overlay.
+  private final Map<HoodieSchema, StructType> rowStructTypes = new 
ConcurrentHashMap<>();
 
   protected BaseSparkInternalRecordContext(HoodieTableConfig tableConfig) {
     super(tableConfig, new DefaultJavaTypeConverter());
@@ -73,7 +77,10 @@ public abstract class BaseSparkInternalRecordContext extends 
RecordContext<Inter
   }
 
   private static Object getFieldValueFromInternalRowInternal(InternalRow row, 
HoodieSchema recordSchema, String fieldName, boolean convertToJavaType) {
-    StructType structType = getCachedSchema(recordSchema);
+    return getFieldValueFromInternalRowInternal(row, 
getCachedSchema(recordSchema), fieldName, convertToJavaType);
+  }
+
+  private static Object getFieldValueFromInternalRowInternal(InternalRow row, 
StructType structType, String fieldName, boolean convertToJavaType) {
     scala.Option<HoodieUnsafeRowUtils.NestedFieldPath> cachedNestedFieldPath =
         HoodieInternalRowUtils.getCachedPosList(structType, fieldName);
     if (cachedNestedFieldPath.isDefined()) {
@@ -106,7 +113,7 @@ public abstract class BaseSparkInternalRecordContext 
extends RecordContext<Inter
 
   @Override
   public Object getValue(InternalRow row, HoodieSchema schema, String 
fieldName) {
-    return getFieldValueFromInternalRow(row, schema, fieldName);
+    return getFieldValueFromInternalRowInternal(row, getRowStructType(schema), 
fieldName, false);
   }
 
   @Override
@@ -210,14 +217,43 @@ public abstract class BaseSparkInternalRecordContext 
extends RecordContext<Inter
     if (internalRow instanceof UnsafeRow) {
       return internalRow;
     }
-    final UnsafeProjection unsafeProjection = 
HoodieInternalRowUtils.getCachedUnsafeProjection(schema);
+    final UnsafeProjection unsafeProjection;
+    if (rowShape == null) {
+      unsafeProjection = 
HoodieInternalRowUtils.getCachedUnsafeProjection(schema);
+    } else {
+      StructType rowStructType = getRowStructType(schema);
+      unsafeProjection = 
HoodieInternalRowUtils.getCachedUnsafeProjection(rowStructType, rowStructType);
+    }
     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 and field accessor 
this context builds from an engine
+   * schema has to be typed over that shape: a VariantType-typed one reads a 
projection struct through
+   * UnsafeRow.getVariant instead of copying it across. That covers {@link 
#projectRecord}, {@link #getValue} (the
+   * partial-update merges read the older record through it) and {@link 
#toBinaryRow}; SparkRecordMergingUtils
+   * reads the merged record's fields through {@link #getRowStructType} too.
+   */
+  public void setRowShape(UnaryOperator<StructType> rowShape) {
+    this.rowShape = rowShape;
+    rowStructTypes.clear();
+  }
+
+  /**
+   * 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 : 
rowStructTypes.computeIfAbsent(schema, s -> 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());
     return row -> (InternalRow) unsafeRowWriter.apply(row);
   }
 
diff --git 
a/hudi-client/hudi-spark-client/src/main/scala/org/apache/hudi/BaseSparkInternalRowReaderContext.java
 
b/hudi-client/hudi-spark-client/src/main/scala/org/apache/hudi/BaseSparkInternalRowReaderContext.java
index b0e2e938b06c..8da532bbcd17 100644
--- 
a/hudi-client/hudi-spark-client/src/main/scala/org/apache/hudi/BaseSparkInternalRowReaderContext.java
+++ 
b/hudi-client/hudi-spark-client/src/main/scala/org/apache/hudi/BaseSparkInternalRowReaderContext.java
@@ -45,7 +45,6 @@ import java.util.stream.Collectors;
 import scala.Function1;
 
 import static 
org.apache.hudi.common.config.HoodieReaderConfig.RECORD_MERGE_IMPL_CLASSES_WRITE_CONFIG_KEY;
-import static org.apache.spark.sql.HoodieInternalRowUtils.getCachedSchema;
 
 /**
  * An abstract class implementing {@link HoodieReaderContext} to handle {@link 
InternalRow}s.
@@ -83,6 +82,8 @@ public abstract class BaseSparkInternalRowReaderContext 
extends HoodieReaderCont
   /**
    * Constructs a transformation that will take a row and convert it to a new 
row with the given schema and adds in the values for the partition columns if 
they are missing in the returned row.
    * It is assumed that the `to` schema will contain the partition fields.
+   * The data-file rows arrive in the shape {@link #getFileRecordIterator} 
read them in, so the writer is typed over
+   * the record context's row shape rather than over the plain engine-schema 
conversion.
    * @param from the original schema
    * @param to the schema the row will be converted to
    * @param partitionFieldAndValues the partition fields and their values, if 
any are required by the reader
@@ -91,8 +92,10 @@ public abstract class BaseSparkInternalRowReaderContext 
extends HoodieReaderCont
   protected UnaryOperator<InternalRow> getBootstrapProjection(HoodieSchema 
from, HoodieSchema to, List<Pair<String, Object>> partitionFieldAndValues) {
     Map<Integer, Object> partitionValuesByIndex = 
partitionFieldAndValues.stream()
         .collect(Collectors.toMap(pair -> 
to.getField(pair.getKey()).orElseThrow(() -> new 
IllegalArgumentException("Missing field: " + pair.getKey())).pos(), 
Pair::getRight));
+    BaseSparkInternalRecordContext sparkRecordContext = 
(BaseSparkInternalRecordContext) recordContext;
     Function1<InternalRow, UnsafeRow> unsafeRowWriter =
-        HoodieInternalRowUtils.getCachedUnsafeRowWriter(getCachedSchema(from), 
getCachedSchema(to), Collections.emptyMap(), partitionValuesByIndex);
+        
HoodieInternalRowUtils.getCachedUnsafeRowWriter(sparkRecordContext.getRowStructType(from),
+            sparkRecordContext.getRowStructType(to), Collections.emptyMap(), 
partitionValuesByIndex);
     return row -> (InternalRow) unsafeRowWriter.apply(row);
   }
 
diff --git 
a/hudi-client/hudi-spark-client/src/main/scala/org/apache/hudi/SparkFileFormatInternalRowReaderContext.scala
 
b/hudi-client/hudi-spark-client/src/main/scala/org/apache/hudi/SparkFileFormatInternalRowReaderContext.scala
index 71118d235b13..ce35119ab6df 100644
--- 
a/hudi-client/hudi-spark-client/src/main/scala/org/apache/hudi/SparkFileFormatInternalRowReaderContext.scala
+++ 
b/hudi-client/hudi-spark-client/src/main/scala/org/apache/hudi/SparkFileFormatInternalRowReaderContext.scala
@@ -28,6 +28,7 @@ import 
org.apache.hudi.common.model.HoodieRecordMerger.PAYLOAD_BASED_MERGE_STRAT
 import org.apache.hudi.common.schema.{HoodieSchema, HoodieSchemaUtils}
 import org.apache.hudi.common.table.HoodieTableConfig
 import org.apache.hudi.common.table.log.InstantRange
+import org.apache.hudi.common.table.read.FileGroupReaderSchemaHandler
 import 
org.apache.hudi.common.table.read.buffer.PositionBasedFileGroupRecordBuffer.ROW_INDEX_TEMPORARY_COLUMN_NAME
 import org.apache.hudi.common.util.{HoodieVectorUtils, Option => HOption}
 import org.apache.hudi.common.util.ValidationUtils.checkState
@@ -69,7 +70,9 @@ import scala.collection.JavaConverters._
  *                            collapses the projected struct back to a plain 
VARIANT, so the projected Spark
  *                            schema cannot be recovered from a HoodieSchema 
round-trip (#18739 sub-task 4).
  *                            Kept Spark-side so the engine-neutral schema 
model stays free of Spark 4.1
- *                            variant concepts.
+ *                            variant concepts. The same overlay also types 
the record context's row writers
+ *                            - the required-to-requested output converter and 
the bootstrap skeleton/data
+ *                            join - through 
BaseSparkInternalRecordContext.setRowShape, see setSchemaHandler.
  * @param instantRangeOpt optional requested-time range applied to base and 
log records before merging
  */
 class SparkFileFormatInternalRowReaderContext(baseFileReader: 
SparkColumnarFileReader,
@@ -124,20 +127,44 @@ 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. Every row 
writer and field accessor
+  // the record context builds from an engine schema needs the shape: the 
output converter, which
+  // projects the reader's required schema down to the requested one
+  // (FileGroupReaderSchemaHandler.getOutputConverter), the bootstrap 
skeleton/data join
+  // (getBootstrapProjection), and the partial-update merges that rebuild a 
row field by field
+  // (SparkRecordMergingUtils.mergePartialRecords, mergeWithEngineRecord). 
Without it they read 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(
+        overlayVariantProjections(_, requiredStruct))
+    }
   }
 
   // Aligns avro log-block records with the PushVariantIntoScan-projected 
variant shape before
diff --git 
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/dml/schema/TestVariantShreddingMixedLayouts.scala
 
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/dml/schema/TestVariantShreddingMixedLayouts.scala
index bbe0deddf02b..a9921d8a8d0c 100644
--- 
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/dml/schema/TestVariantShreddingMixedLayouts.scala
+++ 
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/dml/schema/TestVariantShreddingMixedLayouts.scala
@@ -19,13 +19,16 @@
 
 package org.apache.spark.sql.hudi.dml.schema
 
-import org.apache.hudi.{HoodieSchemaConversionUtils, HoodieSparkUtils, 
HoodieTableSchema, SparkAdapterSupport}
+import org.apache.hudi.{DataSourceReadOptions, DataSourceWriteOptions, 
HoodieSchemaConversionUtils, HoodieSparkUtils, HoodieTableSchema, 
SparkAdapterSupport}
+import org.apache.hudi.common.config.HoodieMetadataConfig
 import org.apache.hudi.common.fs.FSUtils
 import org.apache.hudi.common.model.HoodieFileFormat
 import org.apache.hudi.common.model.HoodieRecord.HoodieRecordType
 import org.apache.hudi.common.schema.{HoodieSchema, HoodieSchemaField, 
HoodieSchemaType}
 import org.apache.hudi.common.table.log.block.HoodieLogBlock.HoodieLogBlockType
+import org.apache.hudi.config.{HoodieBootstrapConfig, HoodieWriteConfig}
 import org.apache.hudi.core.io.storage.VariantShreddingInferenceFileWriter
+import org.apache.hudi.keygen.NonpartitionedKeyGenerator
 import org.apache.hudi.testutils.DataSourceTestUtils
 
 import org.apache.hadoop.fs.{Path => HadoopPath}
@@ -664,6 +667,227 @@ class TestVariantShreddingMixedLayouts extends 
HoodieSparkSqlTestBase with Varia
             s"select count(*) from $tableName where try_variant_get(v, '$$.a', 
'bigint') > 100")(Seq(2))
           checkAnswer(
             s"select id from $tableName where variant_get(v, '$$.b', 'string') 
= 'b7'")(Seq(7))
+
+          // A MERGE INTO that assigns ts alone writes a partial log block
+          // (hoodie.spark.sql.merge.into.partial.updates defaults to true on 
MOR), so the read
+          // rebuilds id 4 field by field 
(SparkRecordMergingUtils.mergePartialRecords) with v taken
+          // from the base-file record, which carries the projection struct on 
the pushed arm. Eight
+          // pushed paths of which the row holds one make seven of its fields 
null, as in the read
+          // paths legs below.
+          spark.sql(s"merge into $tableName t using (select 4 as id, 1003L as 
ts) s on t.id = s.id " +
+            "when matched then update set ts = s.ts")
+          assert(hasPartialLogBlock(tablePath), s"[$leg] expected the MERGE 
INTO to write a partial log block")
+          checkAnswer(s"select id, variant_get(v, '$$.a', 'bigint'), ts from 
$tableName where id = 4")(
+            Seq(4, 4L, 1003L))
+          val widePaths = ('a' to 'h').map(c => s"try_variant_get(v, '$$.$c', 
'bigint')").mkString(", ")
+          checkAnswer(s"select $widePaths from $tableName where id = 4")(
+            Seq(4L, null, null, null, null, null, null, null))
+        }
+      }
+    }
+  }
+
+  test("Time travel and incremental V1/V2 reads carry the variant projection 
on both pushVariantIntoScan arms") {
+    assume(HoodieSparkUtils.gteqSpark4_1, SPARK_4_1_GATE)
+
+    // PushVariantIntoScan matches every LogicalRelation whose 
HadoopFsRelation carries a
+    // ParquetFileFormat, and each Hudi read relation - snapshot, time travel, 
incremental V1 and
+    // V2 - is exactly that over HoodieFileGroupReaderBasedFileFormat, so all 
of them get the
+    // projection struct pushed into the scan. The legs above pin the snapshot 
relation only.
+    // Read-mode test; SPARK pinned (see the mixed-files test above).
+    Seq("true", "false").foreach { pushIntoScan =>
+      withSQLConf("spark.sql.variant.pushVariantIntoScan" -> pushIntoScan) {
+        withVariantTable(s"read paths pushVariantIntoScan=$pushIntoScan", 
"mor",
+          props = Seq("hoodie.compact.inline = 'false'"), recordTypes = 
Seq(HoodieRecordType.SPARK)) {
+          (tableName, tablePath, leg) =>
+          withWriteLayout(Forced("a bigint")) {
+            spark.sql(s"insert into $tableName ${variantSourceSql(Seq((0 until 
10, ObjA)))}")
+          }
+          val baseInstant = latestCompletedInstant(tablePath)
+          // One unshredded log over the shredded base: ids 0-4 stay in the 
base file's typed slot,
+          // ids 5-8 come out of the log's residual and id 9's variant is 
nulled out by the log.
+          withWriteLayout(Unshredded) {
+            spark.sql(s"update $tableName set " +
+              s"""v = case when id = 9 then null else 
parse_json(concat('{"a":', 100 + id, '}')) end, """ +
+              "ts = 1001 where id >= 5")
+          }
+
+          // Every leg below expects the same rows on both arms, so the plan 
assertions are what
+          // tell them apart.
+          val pushed = pushIntoScan.toBoolean
+          val verdict = if (pushed) "should have" else "must not have"
+          val projectedA = s"id, variant_get(v, '$$.a', 'bigint')"
+          val latestRows: Seq[Seq[Any]] =
+            (0 until 9).map(id => Seq(id, if (id < 5) id.toLong else 100L + 
id)) :+ Seq(9, null)
+          val preLogRows: Seq[Seq[Any]] = (0 until 10).map(id => Seq(id, 
id.toLong))
+
+          // Time travel builds its own relation, at an instant before the log 
existed.
+          val asOfSql = s"select $projectedA from $tableName timestamp as of 
'$baseInstant' order by id"
+          checkAnswer(asOfSql)(preLogRows: _*)
+          checkAnswer(s"select count(*) from $tableName timestamp as of 
'$baseInstant' where v is null")(Seq(0))
+          assert(variantProjectionPushedIntoScan(asOfSql) == pushed,
+            s"[$leg] PushVariantIntoScan $verdict rewritten v into a 
projection struct (time travel)")
+
+          // The latest snapshot is the reference both incremental relations 
have to reproduce.
+          checkAnswer(s"select $projectedA from $tableName order by 
id")(latestRows: _*)
+          checkAnswer(s"select id from $tableName where v is null")(Seq(9))
+          // Eight pushed paths of which the row holds one, so seven fields of 
the projection struct
+          // are null. Before the record context typed its row writers over 
the projected shape, the
+          // output converter (the reader's required schema down to the 
requested one) rebuilt that
+          // struct through a VariantType-typed writer, which read the null 
bitset as a variant
+          // length and threw NegativeArraySizeException.
+          val widePaths = ('a' to 'h').map(c => s"try_variant_get(v, '$$.$c', 
'bigint')").mkString(", ")
+          checkAnswer(s"select $widePaths from $tableName where id = 2")(
+            Seq(2L, null, null, null, null, null, null, null))
+
+          // Incremental V2 is the completion-time relation used on table 
version 8+, and over the
+          // whole timeline it returns the latest state of every key. V1 is 
the requested-time
+          // relation, which DefaultSource selects on this same table once the 
read option asks
+          // for a table version below 8: it is bounded at the base instant's 
REQUESTED time and
+          // must come back with the pre-log rows. Both relations feed the 
same file format, the
+          // V1 one through required filters on _hoodie_commit_time that the 
format adds back as
+          // filter-only read columns beside the projected variant.
+          Seq(
+            ("v2", Map.empty[String, String], latestRows, 1),
+            ("v1", 
Map(DataSourceReadOptions.INCREMENTAL_READ_TABLE_VERSION.key -> "6",
+              DataSourceReadOptions.END_COMMIT.key -> baseInstant), 
preLogRows, 0)
+          ).foreach { case (version, versionOpts, expectedRows, nullCount) =>
+            val view = s"${tableName}_inc_$version"
+            spark.read.format("hudi")
+              .option(DataSourceReadOptions.QUERY_TYPE.key, 
DataSourceReadOptions.QUERY_TYPE_INCREMENTAL_OPT_VAL)
+              .option(DataSourceReadOptions.START_COMMIT.key, "000")
+              .options(versionOpts)
+              .load(tablePath)
+              .createOrReplaceTempView(view)
+            val incSql = s"select $projectedA from $view order by id"
+            checkAnswer(incSql)(expectedRows: _*)
+            checkAnswer(s"select count(*) from $view where v is 
null")(Seq(nullCount))
+            assert(variantProjectionPushedIntoScan(incSql) == pushed,
+              s"[$leg] PushVariantIntoScan $verdict rewritten v into a 
projection struct " +
+                s"(incremental $version)")
+            spark.catalog.dropTempView(view)
+          }
+
+          // What proves the read option really selected V1 rather than 
falling through to V2:
+          // the same END_COMMIT bound read by V2 is a COMPLETION time, and 
the base commit
+          // completed after its own requested time, so that read selects 
nothing.
+          val v2AtRequestedTime = spark.read.format("hudi")
+            .option(DataSourceReadOptions.QUERY_TYPE.key, 
DataSourceReadOptions.QUERY_TYPE_INCREMENTAL_OPT_VAL)
+            .option(DataSourceReadOptions.START_COMMIT.key, "000")
+            .option(DataSourceReadOptions.END_COMMIT.key, baseInstant)
+            .load(tablePath)
+          assert(v2AtRequestedTime.count() == 0,
+            s"[$leg] V2 reads END_COMMIT as a completion time, which the base 
commit's requested time precedes")
+        }
+      }
+    }
+  }
+
+  test("Bootstrapped tables read the variant projection through the skeleton 
and data file join") {
+    assume(HoodieSparkUtils.gteqSpark4_1, SPARK_4_1_GATE)
+
+    // A METADATA_ONLY bootstrap leaves the data in the source parquet file 
and writes a skeleton
+    // base file holding the meta columns only, so HoodieFileGroupReader joins 
the two
+    // (mergeBootstrapReaders) whenever a meta column is required beside the 
data columns: always
+    // on MOR, where the record key drives the log merge, and on COW whenever 
the query asks for
+    // one. The bootstrap relation is a HadoopFsRelation over the same file 
format, so it receives
+    // the PushVariantIntoScan projection struct too, and the variant then has 
to survive that
+    // skeleton/data join rather than a plain single-file read.
+    Seq("cow", "mor").foreach { tableType =>
+      Seq("true", "false").foreach { pushIntoScan =>
+        withSQLConf("spark.sql.variant.pushVariantIntoScan" -> pushIntoScan) {
+          withTempDir { tmp =>
+            val tableName = generateTableName
+            val leg = s"bootstrap $tableType 
pushVariantIntoScan=$pushIntoScan, $tableName"
+            val srcPath = s"${tmp.getCanonicalPath}/source"
+            val tablePath = s"${tmp.getCanonicalPath}/hudi"
+
+            // The source is written by Spark's own parquet writer: a plain 
unshredded variant
+            // column that no Hudi write path ever touched.
+            spark.sql("""select cast(id as int) as id, 
parse_json(concat('{"a":', id, '}')) as v, """ +
+              "1000L as ts from range(0, 10, 1, 1)")
+              .write.parquet(srcPath)
+
+            // The leg runs on the default record type, like 
TestDataSourceForBootstrap's own
+            // metadata-only legs: a METADATA_ONLY bootstrap on the SPARK 
record type fails in the
+            // skeleton-file write (HUDI-5807 - HoodieRowParquetWriteSupport 
cannot resolve the
+            // meta-only schema), which is not a variant matter, so the record 
type is not swept
+            // here. On the default write version the MOR upsert below writes 
a native parquet log
+            // file either way (pinned below), so the log side is projected 
natively like a base
+            // file; the avro rewrite path is owned by the avro-block legs 
elsewhere in this suite.
+            val writeOpts = Map(
+              DataSourceWriteOptions.TABLE_TYPE.key ->
+                (if (tableType == "mor") 
DataSourceWriteOptions.MOR_TABLE_TYPE_OPT_VAL
+                 else DataSourceWriteOptions.COW_TABLE_TYPE_OPT_VAL),
+              HoodieWriteConfig.TBL_NAME.key -> tableName,
+              DataSourceWriteOptions.RECORDKEY_FIELD.key -> "id",
+              DataSourceWriteOptions.ORDERING_FIELDS.key -> "ts",
+              DataSourceWriteOptions.KEYGENERATOR_CLASS_NAME.key -> 
classOf[NonpartitionedKeyGenerator].getName,
+              // The writer's default column-stats index is rejected by the 
bootstrap commit
+              // ("col stats is not supported with bootstrap operation"), as 
in TestDataSourceForBootstrap.
+              HoodieMetadataConfig.ENABLE_METADATA_INDEX_COLUMN_STATS.key -> 
"false")
+
+            // METADATA_ONLY is the default bootstrap mode selector, so the 
data stays in srcPath.
+            spark.emptyDataFrame.write.format("hudi")
+              .options(writeOpts)
+              .option(DataSourceWriteOptions.OPERATION.key, 
DataSourceWriteOptions.BOOTSTRAP_OPERATION_OPT_VAL)
+              .option(HoodieBootstrapConfig.BASE_PATH.key, srcPath)
+              .mode(SaveMode.Overwrite)
+              .save(tablePath)
+
+            if (tableType == "mor") {
+              // A log over the bootstrapped file group, so the merged read 
joins skeleton, source
+              // file and log. COW stays read-only: an upsert there would 
rewrite the file group
+              // into a regular base file and the bootstrap read path would be 
gone.
+              spark.sql("""select cast(id as int) as id, case when id = 9 then 
null """ +
+                """else parse_json(concat('{"a":', 100 + id, '}')) end as v, 
""" +
+                "1001L as ts from range(5, 10, 1, 1)")
+                .write.format("hudi")
+                .options(writeOpts)
+                .option(DataSourceWriteOptions.OPERATION.key, 
DataSourceWriteOptions.UPSERT_OPERATION_OPT_VAL)
+                .mode(SaveMode.Append)
+                .save(tablePath)
+              
assert(listDataParquetFiles(tablePath).exists(_.endsWith(".log.parquet")),
+                s"[$leg] expected the upsert to write a native parquet log 
file")
+            }
+
+            // The data is still served out of the source file: every base 
file the table itself
+            // wrote is a skeleton, meta columns only. Log files are left out 
- the MOR upsert's
+            // log block carries the whole record, variant included.
+            val skeletonFiles = listDataParquetFiles(tablePath)
+              .filter(f => FSUtils.isBaseFile(new HadoopPath(f).getName))
+            assert(skeletonFiles.nonEmpty, s"[$leg] expected at least one 
skeleton base file")
+            skeletonFiles.foreach(f => 
assert(!readParquetSchema(f).containsField("v"),
+              s"[$leg] skeleton file must not carry the data column: $f"))
+
+            val view = s"${tableName}_view"
+            
spark.read.format("hudi").load(tablePath).createOrReplaceTempView(view)
+            val pushed = pushIntoScan.toBoolean
+            val verdict = if (pushed) "should have" else "must not have"
+            val expected: Seq[Seq[Any]] = if (tableType == "mor") {
+              (0 until 9).map(id => Seq(id, if (id < 5) id.toLong else 100L + 
id)) :+ Seq(9, null)
+            } else {
+              (0 until 10).map(id => Seq(id, id.toLong))
+            }
+            val projectedSql = s"select id, variant_get(v, '$$.a', 'bigint') 
from $view order by id"
+            checkAnswer(projectedSql)(expected: _*)
+            // A meta column beside the variant forces the skeleton/data join 
on COW as well; on
+            // MOR the record key is already required for the log merge.
+            checkAnswer(s"select _hoodie_record_key, variant_get(v, '$$.a', 
'bigint') from $view " +
+              "where id in (2, 7) order by id")(
+              Seq("2", 2L), Seq("7", if (tableType == "mor") 107L else 7L))
+            checkAnswer(s"select count(*) from $view where v is null")(Seq(if 
(tableType == "mor") 1 else 0))
+            // Eight pushed paths of which the row holds one, so seven fields 
of the projection
+            // struct are null. Before the record context typed its row 
writers over the projected
+            // shape, the bootstrap skeleton/data join rebuilt that struct 
through a VariantType-typed
+            // writer, which read the null bitset as a variant length and 
threw NegativeArraySizeException.
+            val widePaths = ('a' to 'h').map(c => s"try_variant_get(v, 
'$$.$c', 'bigint')").mkString(", ")
+            checkAnswer(s"select _hoodie_record_key, $widePaths from $view 
where id = 2")(
+              Seq("2", 2L, null, null, null, null, null, null, null))
+            assert(variantProjectionPushedIntoScan(projectedSql) == pushed,
+              s"[$leg] PushVariantIntoScan $verdict rewritten v into a 
projection struct")
+            spark.catalog.dropTempView(view)
+          }
         }
       }
     }
diff --git 
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/dml/schema/VariantShreddingTestSupport.scala
 
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/dml/schema/VariantShreddingTestSupport.scala
index 3ed37a776080..8c0d5a980170 100644
--- 
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/dml/schema/VariantShreddingTestSupport.scala
+++ 
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/dml/schema/VariantShreddingTestSupport.scala
@@ -26,6 +26,7 @@ import 
org.apache.hudi.common.model.HoodieRecord.HoodieRecordType
 import org.apache.hudi.common.model.WriteOperationType
 import org.apache.hudi.common.table.TableSchemaResolver
 import org.apache.hudi.common.table.log.HoodieLogFormat
+import org.apache.hudi.common.table.log.block.{HoodieDataBlock, HoodieLogBlock}
 import org.apache.hudi.common.table.log.block.HoodieLogBlock.HoodieLogBlockType
 import org.apache.hudi.common.testutils.HoodieTestUtils
 import org.apache.hudi.storage.StoragePath
@@ -743,7 +744,17 @@ trait VariantShreddingTestSupport { self: 
HoodieSparkSqlTestBase =>
    * a log format assert on this rather than on file names: native logs carry 
a .log.parquet suffix,
    * but an inline log file is named the same whether its data blocks are avro 
or parquet.
    */
-  protected def listLogBlockTypes(tablePath: String): Seq[HoodieLogBlockType] 
= {
+  protected def listLogBlockTypes(tablePath: String): Seq[HoodieLogBlockType] =
+    mapLogBlocks(tablePath)(_.getBlockType)
+
+  /** Whether any data block in the table's log files carries a partial-update 
schema (IS_PARTIAL). */
+  protected def hasPartialLogBlock(tablePath: String): Boolean =
+    mapLogBlocks(tablePath) {
+      case dataBlock: HoodieDataBlock => dataBlock.containsPartialUpdates()
+      case _ => false
+    }.contains(true)
+
+  private def mapLogBlocks[T](tablePath: String)(f: HoodieLogBlock => T): 
Seq[T] = {
     val (metaClient, fsView) = getMetaClientAndFileSystemView(tablePath)
     val schema = new TableSchemaResolver(metaClient).getTableSchema
     val logFiles = fsView.getAllFileSlices("").iterator().asScala
@@ -752,11 +763,11 @@ trait VariantShreddingTestSupport { self: 
HoodieSparkSqlTestBase =>
     logFiles.flatMap { path =>
       val reader = HoodieLogFormat.newReader(metaClient, new 
HoodieLogFile(path), schema)
       try {
-        val types = mutable.ArrayBuffer[HoodieLogBlockType]()
+        val results = mutable.ArrayBuffer[T]()
         while (reader.hasNext) {
-          types += reader.next().getBlockType
+          results += f(reader.next())
         }
-        types.toSeq
+        results.toSeq
       } finally {
         reader.close()
       }

Reply via email to