voonhous commented on code in PR #20041:
URL: https://github.com/apache/hudi/pull/20041#discussion_r4100675376
##########
hudi-client/hudi-spark-client/src/main/scala/org/apache/hudi/SparkFileFormatInternalRowReaderContext.scala:
##########
@@ -124,20 +127,45 @@ class
SparkFileFormatInternalRowReaderContext(baseFileReader: SparkColumnarFileR
})
}
+ // Whether the query carries a Spark 4.1 PushVariantIntoScan projection at
all.
+ private lazy val hasVariantProjection: Boolean =
+ sparkRequiredSchema.exists(_.fields.exists(f =>
sparkAdapter.containsVariantProjection(f.dataType)))
+
+ private def isPayloadBasedMerge: Boolean = {
+ // getRecordMerger() is a Lombok getter over a field initialized to null
(not Option.empty());
+ // it stays null until HoodieReaderContext.initRecordMerger runs
(HoodieFileGroupReader calls it
+ // from its constructor), so the null guard is required.
+ val merger = getRecordMerger()
+ merger != null && merger.isPresent && merger.get.getMergingStrategy ==
PAYLOAD_BASED_MERGE_STRATEGY_UUID
+ }
+
// True only when there is a Spark 4.1 PushVariantIntoScan projection to
apply AND the table is
// not using a custom (payload-based) merger. Payload-based tables
round-trip records through
// PayloadUpdateProcessor.convertToAvroRecord against a schema that still
types variant fields as
// VariantType, so a row already rewritten into the projected struct shape
would be mis-decoded.
// Single source of truth for both reader paths (parquet native projection +
avro rewrite).
- private def shouldProjectVariants(): Boolean = {
- val hasVariantProjection =
- sparkRequiredSchema.exists(_.fields.exists(f =>
sparkAdapter.containsVariantProjection(f.dataType)))
- // getRecordMerger() is a Lombok getter over a field initialized to null
(not Option.empty());
- // it stays null until HoodieReaderContext.initRecordMerger runs
(HoodieFileGroupReader calls it
- // from its constructor), so the null guard is required.
- val merger = getRecordMerger()
- val isPayloadBased = merger != null && merger.isPresent &&
merger.get.getMergingStrategy == PAYLOAD_BASED_MERGE_STRATEGY_UUID
- hasVariantProjection && !isPayloadBased
+ private def shouldProjectVariants(): Boolean = hasVariantProjection &&
!isPayloadBasedMerge
+
+ // HoodieFileGroupReader installs the schema handler after initRecordMerger
and before it asks for
+ // the output converter, so this is where the record context learns what
shape its rows carry.
+ // Base-file rows are ALWAYS read in the projected shape
(getFileRecordIterator overlays it
+ // unconditionally); log rows only when shouldProjectVariants rewrites them,
so the payload-based
+ // exclusion applies only when there are log files to merge. Two consumers
need the shape: the
+ // output converter, which projects the reader's required schema down to the
requested one
+ // (FileGroupReaderSchemaHandler.getOutputConverter), and the bootstrap
skeleton/data join
+ // (getBootstrapProjection). Without it both build a VariantType-typed row
writer that re-encodes
+ // the projection struct through UnsafeRow.getVariant - byte-identical only
while the struct's
+ // null bitset stays small, and a NegativeArraySizeException once enough
pushed fields are null.
+ override def setSchemaHandler(schemaHandler:
FileGroupReaderSchemaHandler[InternalRow]): Unit = {
+ super.setSchemaHandler(schemaHandler)
+ if (hasVariantProjection && (!isPayloadBasedMerge || !getHasLogFiles)) {
Review Comment:
Not introduced here: the base-file overlay has been unconditional since
#18923, and this PR leaves the payload-based-with-logs arm as master reads it.
Gating the base overlay alone would not fix it either, since the merged rows
would then reach a plan expecting the projection struct as a `VariantVal`; that
arm also needs the merged output projected (`buildVariantProjector`). Tracked
on #17744 (https://github.com/apache/hudi/issues/17744#issuecomment-5826048476)
rather than widening this PR.
##########
hudi-client/hudi-spark-client/src/main/scala/org/apache/hudi/SparkFileFormatInternalRowReaderContext.scala:
##########
@@ -124,20 +127,45 @@ class
SparkFileFormatInternalRowReaderContext(baseFileReader: SparkColumnarFileR
})
}
+ // Whether the query carries a Spark 4.1 PushVariantIntoScan projection at
all.
+ private lazy val hasVariantProjection: Boolean =
+ sparkRequiredSchema.exists(_.fields.exists(f =>
sparkAdapter.containsVariantProjection(f.dataType)))
+
+ private def isPayloadBasedMerge: Boolean = {
+ // getRecordMerger() is a Lombok getter over a field initialized to null
(not Option.empty());
+ // it stays null until HoodieReaderContext.initRecordMerger runs
(HoodieFileGroupReader calls it
+ // from its constructor), so the null guard is required.
+ val merger = getRecordMerger()
+ merger != null && merger.isPresent && merger.get.getMergingStrategy ==
PAYLOAD_BASED_MERGE_STRATEGY_UUID
+ }
+
// True only when there is a Spark 4.1 PushVariantIntoScan projection to
apply AND the table is
// not using a custom (payload-based) merger. Payload-based tables
round-trip records through
// PayloadUpdateProcessor.convertToAvroRecord against a schema that still
types variant fields as
// VariantType, so a row already rewritten into the projected struct shape
would be mis-decoded.
// Single source of truth for both reader paths (parquet native projection +
avro rewrite).
- private def shouldProjectVariants(): Boolean = {
- val hasVariantProjection =
- sparkRequiredSchema.exists(_.fields.exists(f =>
sparkAdapter.containsVariantProjection(f.dataType)))
- // getRecordMerger() is a Lombok getter over a field initialized to null
(not Option.empty());
- // it stays null until HoodieReaderContext.initRecordMerger runs
(HoodieFileGroupReader calls it
- // from its constructor), so the null guard is required.
- val merger = getRecordMerger()
- val isPayloadBased = merger != null && merger.isPresent &&
merger.get.getMergingStrategy == PAYLOAD_BASED_MERGE_STRATEGY_UUID
- hasVariantProjection && !isPayloadBased
+ private def shouldProjectVariants(): Boolean = hasVariantProjection &&
!isPayloadBasedMerge
+
+ // HoodieFileGroupReader installs the schema handler after initRecordMerger
and before it asks for
+ // the output converter, so this is where the record context learns what
shape its rows carry.
+ // Base-file rows are ALWAYS read in the projected shape
(getFileRecordIterator overlays it
+ // unconditionally); log rows only when shouldProjectVariants rewrites them,
so the payload-based
+ // exclusion applies only when there are log files to merge. Two consumers
need the shape: the
+ // output converter, which projects the reader's required schema down to the
requested one
+ // (FileGroupReaderSchemaHandler.getOutputConverter), and the bootstrap
skeleton/data join
+ // (getBootstrapProjection). Without it both build a VariantType-typed row
writer that re-encodes
+ // the projection struct through UnsafeRow.getVariant - byte-identical only
while the struct's
+ // null bitset stays small, and a NegativeArraySizeException once enough
pushed fields are null.
+ override def setSchemaHandler(schemaHandler:
FileGroupReaderSchemaHandler[InternalRow]): Unit = {
+ super.setSchemaHandler(schemaHandler)
+ if (hasVariantProjection && (!isPayloadBasedMerge || !getHasLogFiles)) {
+ val requiredStruct = sparkRequiredSchema.get
+ recordContext.asInstanceOf[BaseSparkInternalRecordContext].setRowShape(
+ new UnaryOperator[StructType] {
+ override def apply(structType: StructType): StructType =
Review Comment:
Done.
##########
hudi-client/hudi-spark-client/src/main/scala/org/apache/hudi/BaseSparkInternalRecordContext.java:
##########
@@ -214,10 +215,30 @@ public InternalRow toBinaryRow(HoodieSchema schema,
InternalRow internalRow) {
return unsafeProjection.apply(internalRow);
}
+ /**
+ * Installs the Spark type the rows of an engine schema actually carry in
this read; null, the default, means the
+ * plain conversion. {@code SparkFileFormatInternalRowReaderContext}
installs the PushVariantIntoScan overlay here
+ * (see its setSchemaHandler), because every row writer this context builds
from an engine schema has to be typed
+ * over that shape: a VariantType-typed writer re-encodes a projection
struct through UnsafeRow.getVariant instead
+ * of copying it across.
+ */
+ public void setRowShape(UnaryOperator<StructType> rowShape) {
+ this.rowShape = rowShape;
+ }
+
+ /**
+ * The Spark type the rows of {@code schema} carry in this read: the plain
conversion, or the row shape installed
+ * by {@link #setRowShape} when the reader hands its rows over in a
rewritten shape.
+ */
+ public StructType getRowStructType(HoodieSchema schema) {
+ StructType structType = getCachedSchema(schema);
+ return rowShape == null ? structType : rowShape.apply(structType);
+ }
+
@Override
public UnaryOperator<InternalRow> projectRecord(HoodieSchema from,
HoodieSchema to, Map<String, String> renamedColumns) {
Function1<InternalRow, UnsafeRow> unsafeRowWriter =
- HoodieInternalRowUtils.getCachedUnsafeRowWriter(getCachedSchema(from),
getCachedSchema(to), renamedColumns, Collections.emptyMap());
+
HoodieInternalRowUtils.getCachedUnsafeRowWriter(getRowStructType(from),
getRowStructType(to), renamedColumns, Collections.emptyMap());
Review Comment:
Yes. Partial updates are on by default for MOR `MERGE INTO`, and
`mergePartialRecords` / `mergeWithEngineRecord` read the base record as
`VariantType`.
Fixed: `getValue`, `toBinaryRow` and `mergePartialRecords` now read through
`getRowStructType`. The MOR `variant_get` leg adds a ts-only `MERGE INTO` over
a base-file record (asserts the partial block) and reads it back on both
`pushVariantIntoScan` arms.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]