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

voonhous pushed a commit to branch variant-read-paths
in repository https://gitbox.apache.org/repos/asf/hudi.git

commit 96ae47f44ff6ec8cb88dbb2f69d96726e067afab
Author: voon <[email protected]>
AuthorDate: Wed Sep 23 23:46:56 2026 +0800

    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.
---
 .../hudi/BaseSparkInternalRecordContext.java       |  23 ++-
 .../hudi/BaseSparkInternalRowReaderContext.java    |   7 +-
 .../SparkFileFormatInternalRowReaderContext.scala  |  50 +++--
 .../schema/TestVariantShreddingMixedLayouts.scala  | 211 ++++++++++++++++++++-
 4 files changed, 276 insertions(+), 15 deletions(-)

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..5d289c35a940 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
@@ -55,6 +55,7 @@ import static 
org.apache.spark.sql.HoodieInternalRowUtils.getCachedSchema;
 public abstract class BaseSparkInternalRecordContext extends 
RecordContext<InternalRow> {
 
   private OrderingValueEngineTypeConverter orderingValueConverter;
+  private UnaryOperator<StructType> rowShape;
 
   protected BaseSparkInternalRecordContext(HoodieTableConfig tableConfig) {
     super(tableConfig, new DefaultJavaTypeConverter());
@@ -214,10 +215,30 @@ public abstract class BaseSparkInternalRecordContext 
extends RecordContext<Inter
     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());
     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..b7af6c820382 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
@@ -47,7 +48,7 @@ import org.apache.spark.sql.sources.Filter
 import org.apache.spark.sql.types.{ArrayType, ByteType, DoubleType, FloatType, 
LongType, MetadataBuilder, StructField, StructType}
 import org.apache.spark.sql.vectorized.{ColumnVector, ColumnarBatch}
 
-import java.util.function.{Function => JFunction}
+import java.util.function.{Function => JFunction, UnaryOperator}
 
 import scala.collection.JavaConverters._
 
@@ -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,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 =
+            overlayVariantProjections(structType, 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 345e82d44242..6332c1b69325 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}
@@ -667,6 +670,212 @@ class TestVariantShreddingMixedLayouts extends 
HoodieSparkSqlTestBase with Varia
     }
   }
 
+  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)
+          }
+        }
+      }
+    }
+  }
+
   test("Schema-on-read reads of shredded variant files fail fast") {
     assume(HoodieSparkUtils.gteqSpark4_1, SPARK_4_1_GATE)
 

Reply via email to