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)
