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()
}