This is an automated email from the ASF dual-hosted git repository. voonhous pushed a commit to branch release-1.2.1 in repository https://gitbox.apache.org/repos/asf/hudi.git
commit 8327a728d8784d73901cdb7ef94ced9f2b90cf3c Author: Y Ethan Guo <[email protected]> AuthorDate: Thu Jul 2 22:48:50 2026 -0700 refactor(spark): share Spark 4.x partition-values and mapping classes in hudi-spark4-common (#19148) * refactor(spark): extract shared base classes for HoodiePartitionValues Move the InternalRow delegation methods shared by all Spark versions from Spark3HoodiePartitionValues and the per-version Spark 4.x copies into BaseHoodiePartitionValues in hudi-spark-common, and the getVariant delegation shared by all Spark 4.x versions into an abstract Spark4HoodiePartitionValues in hudi-spark4-common. The Spark 4.0/4.1/4.2 classes now only carry copy() plus, on 4.1/4.2, the geography/geometry getters introduced by Spark 4.1. * refactor(spark): share Spark 4.x mapping and internal-row impls in hudi-spark4-common Extract the bodies duplicated across the Spark 4.0/4.1/4.2 copies of HoodiePartitionFileSliceMapping, HoodiePartitionCDCFileGroupMapping, and HoodieInternalRow into shared traits/abstract classes in hudi-spark4-common, mirroring how hudi-spark3-common already shares them for the 3.x family. Version classes keep only genuine deltas: the Spark 4.1/4.2 geography and geometry getters, and the concrete type returned by copy() via a newInternalRow factory hook. (cherry picked from commit 39aa022c8f2686c0aa0e9745da1ed203ca35e581) --- .../apache/hudi/BaseHoodiePartitionValues.scala} | 11 +-- .../apache/hudi/Spark3HoodiePartitionValues.scala | 81 +------------------- ...Spark4HoodiePartitionCDCFileGroupMapping.scala} | 11 +-- .../Spark4HoodiePartitionFileSliceMapping.scala} | 13 +++- .../apache/hudi/Spark4HoodiePartitionValues.scala} | 21 +++--- .../client/model/Spark4HoodieInternalRow.scala} | 29 ++++---- ...Spark40HoodiePartitionCDCFileGroupMapping.scala | 11 +-- .../Spark40HoodiePartitionFileSliceMapping.scala | 11 +-- .../apache/hudi/Spark40HoodiePartitionValues.scala | 85 +-------------------- .../client/model/Spark40HoodieInternalRow.scala | 25 +++---- ...Spark41HoodiePartitionCDCFileGroupMapping.scala | 9 +-- .../Spark41HoodiePartitionFileSliceMapping.scala | 11 +-- .../apache/hudi/Spark41HoodiePartitionValues.scala | 86 +--------------------- .../client/model/Spark41HoodieInternalRow.scala | 19 ++--- 14 files changed, 79 insertions(+), 344 deletions(-) diff --git a/hudi-spark-datasource/hudi-spark3-common/src/main/scala/org/apache/hudi/Spark3HoodiePartitionValues.scala b/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/BaseHoodiePartitionValues.scala similarity index 91% copy from hudi-spark-datasource/hudi-spark3-common/src/main/scala/org/apache/hudi/Spark3HoodiePartitionValues.scala copy to hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/BaseHoodiePartitionValues.scala index bd50e3ebfca2..a6fd7e238f69 100644 --- a/hudi-spark-datasource/hudi-spark3-common/src/main/scala/org/apache/hudi/Spark3HoodiePartitionValues.scala +++ b/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/BaseHoodiePartitionValues.scala @@ -24,7 +24,12 @@ import org.apache.spark.sql.catalyst.util.{ArrayData, MapData} import org.apache.spark.sql.types.{DataType, Decimal} import org.apache.spark.unsafe.types.{CalendarInterval, UTF8String} -case class Spark3HoodiePartitionValues(values: InternalRow) extends HoodiePartitionValues { +/** + * Abstract base class for HoodiePartitionValues implementations. + * Contains all common delegation logic to the underlying InternalRow. + */ +abstract class BaseHoodiePartitionValues(val values: InternalRow) extends HoodiePartitionValues { + override def numFields: Int = { values.numFields } @@ -37,10 +42,6 @@ case class Spark3HoodiePartitionValues(values: InternalRow) extends HoodiePartit values.update(i, value) } - override def copy(): InternalRow = { - Spark3HoodiePartitionValues(values.copy()) - } - override def isNullAt(ordinal: Int): Boolean = { values.isNullAt(ordinal) } diff --git a/hudi-spark-datasource/hudi-spark3-common/src/main/scala/org/apache/hudi/Spark3HoodiePartitionValues.scala b/hudi-spark-datasource/hudi-spark3-common/src/main/scala/org/apache/hudi/Spark3HoodiePartitionValues.scala index bd50e3ebfca2..acdc175be285 100644 --- a/hudi-spark-datasource/hudi-spark3-common/src/main/scala/org/apache/hudi/Spark3HoodiePartitionValues.scala +++ b/hudi-spark-datasource/hudi-spark3-common/src/main/scala/org/apache/hudi/Spark3HoodiePartitionValues.scala @@ -20,88 +20,11 @@ package org.apache.hudi import org.apache.spark.sql.catalyst.InternalRow -import org.apache.spark.sql.catalyst.util.{ArrayData, MapData} -import org.apache.spark.sql.types.{DataType, Decimal} -import org.apache.spark.unsafe.types.{CalendarInterval, UTF8String} -case class Spark3HoodiePartitionValues(values: InternalRow) extends HoodiePartitionValues { - override def numFields: Int = { - values.numFields - } - - override def setNullAt(i: Int): Unit = { - values.setNullAt(i) - } - - override def update(i: Int, value: Any): Unit = { - values.update(i, value) - } +case class Spark3HoodiePartitionValues(override val values: InternalRow) + extends BaseHoodiePartitionValues(values) { override def copy(): InternalRow = { Spark3HoodiePartitionValues(values.copy()) } - - override def isNullAt(ordinal: Int): Boolean = { - values.isNullAt(ordinal) - } - - override def getBoolean(ordinal: Int): Boolean = { - values.getBoolean(ordinal) - } - - override def getByte(ordinal: Int): Byte = { - values.getByte(ordinal) - } - - override def getShort(ordinal: Int): Short = { - values.getShort(ordinal) - } - - override def getInt(ordinal: Int): Int = { - values.getInt(ordinal) - } - - override def getLong(ordinal: Int): Long = { - values.getLong(ordinal) - } - - override def getFloat(ordinal: Int): Float = { - values.getFloat(ordinal) - } - - override def getDouble(ordinal: Int): Double = { - values.getDouble(ordinal) - } - - override def getDecimal(ordinal: Int, precision: Int, scale: Int): Decimal = { - values.getDecimal(ordinal, precision, scale) - } - - override def getUTF8String(ordinal: Int): UTF8String = { - values.getUTF8String(ordinal) - } - - override def getBinary(ordinal: Int): Array[Byte] = { - values.getBinary(ordinal) - } - - override def getInterval(ordinal: Int): CalendarInterval = { - values.getInterval(ordinal) - } - - override def getStruct(ordinal: Int, numFields: Int): InternalRow = { - values.getStruct(ordinal, numFields) - } - - override def getArray(ordinal: Int): ArrayData = { - values.getArray(ordinal) - } - - override def getMap(ordinal: Int): MapData = { - values.getMap(ordinal) - } - - override def get(ordinal: Int, dataType: DataType): AnyRef = { - values.get(ordinal, dataType) - } } diff --git a/hudi-spark-datasource/hudi-spark4.1.x/src/main/scala/org/apache/hudi/Spark41HoodiePartitionCDCFileGroupMapping.scala b/hudi-spark-datasource/hudi-spark4-common/src/main/scala/org/apache/hudi/Spark4HoodiePartitionCDCFileGroupMapping.scala similarity index 75% copy from hudi-spark-datasource/hudi-spark4.1.x/src/main/scala/org/apache/hudi/Spark41HoodiePartitionCDCFileGroupMapping.scala copy to hudi-spark-datasource/hudi-spark4-common/src/main/scala/org/apache/hudi/Spark4HoodiePartitionCDCFileGroupMapping.scala index 005c17eb4128..428ad6a14112 100644 --- a/hudi-spark-datasource/hudi-spark4.1.x/src/main/scala/org/apache/hudi/Spark41HoodiePartitionCDCFileGroupMapping.scala +++ b/hudi-spark-datasource/hudi-spark4-common/src/main/scala/org/apache/hudi/Spark4HoodiePartitionCDCFileGroupMapping.scala @@ -21,12 +21,13 @@ package org.apache.hudi import org.apache.hudi.common.table.cdc.HoodieCDCFileSplit -import org.apache.spark.sql.catalyst.InternalRow +/** + * Implementation of [[HoodiePartitionCDCFileGroupMapping]] shared by all Spark 4.x + * versions, mixed into the version-specific partition values classes. + */ +trait Spark4HoodiePartitionCDCFileGroupMapping extends HoodiePartitionCDCFileGroupMapping { -class Spark41HoodiePartitionCDCFileGroupMapping(partitionValues: InternalRow, - fileSplits: List[HoodieCDCFileSplit]) - extends Spark41HoodiePartitionValues(partitionValues) - with HoodiePartitionCDCFileGroupMapping { + protected def fileSplits: List[HoodieCDCFileSplit] override def getFileSplits(): List[HoodieCDCFileSplit] = { fileSplits diff --git a/hudi-spark-datasource/hudi-spark4.1.x/src/main/scala/org/apache/hudi/Spark41HoodiePartitionFileSliceMapping.scala b/hudi-spark-datasource/hudi-spark4-common/src/main/scala/org/apache/hudi/Spark4HoodiePartitionFileSliceMapping.scala similarity index 77% copy from hudi-spark-datasource/hudi-spark4.1.x/src/main/scala/org/apache/hudi/Spark41HoodiePartitionFileSliceMapping.scala copy to hudi-spark-datasource/hudi-spark4-common/src/main/scala/org/apache/hudi/Spark4HoodiePartitionFileSliceMapping.scala index 07e199e4a9b1..3046a0ddbf2e 100644 --- a/hudi-spark-datasource/hudi-spark4.1.x/src/main/scala/org/apache/hudi/Spark41HoodiePartitionFileSliceMapping.scala +++ b/hudi-spark-datasource/hudi-spark4-common/src/main/scala/org/apache/hudi/Spark4HoodiePartitionFileSliceMapping.scala @@ -23,10 +23,15 @@ import org.apache.hudi.common.model.FileSlice import org.apache.spark.sql.catalyst.InternalRow -class Spark41HoodiePartitionFileSliceMapping(values: InternalRow, - slices: Map[String, FileSlice]) - extends Spark41HoodiePartitionValues(values) - with HoodiePartitionFileSliceMapping { +/** + * Implementation of [[HoodiePartitionFileSliceMapping]] shared by all Spark 4.x + * versions, mixed into the version-specific partition values classes. + */ +trait Spark4HoodiePartitionFileSliceMapping extends HoodiePartitionFileSliceMapping { + + def values: InternalRow + + protected def slices: Map[String, FileSlice] override def getSlice(fileId: String): Option[FileSlice] = { slices.get(fileId) diff --git a/hudi-spark-datasource/hudi-spark4.1.x/src/main/scala/org/apache/hudi/Spark41HoodiePartitionFileSliceMapping.scala b/hudi-spark-datasource/hudi-spark4-common/src/main/scala/org/apache/hudi/Spark4HoodiePartitionValues.scala similarity index 62% copy from hudi-spark-datasource/hudi-spark4.1.x/src/main/scala/org/apache/hudi/Spark41HoodiePartitionFileSliceMapping.scala copy to hudi-spark-datasource/hudi-spark4-common/src/main/scala/org/apache/hudi/Spark4HoodiePartitionValues.scala index 07e199e4a9b1..a729c6a0574a 100644 --- a/hudi-spark-datasource/hudi-spark4.1.x/src/main/scala/org/apache/hudi/Spark41HoodiePartitionFileSliceMapping.scala +++ b/hudi-spark-datasource/hudi-spark4-common/src/main/scala/org/apache/hudi/Spark4HoodiePartitionValues.scala @@ -19,18 +19,19 @@ package org.apache.hudi -import org.apache.hudi.common.model.FileSlice - import org.apache.spark.sql.catalyst.InternalRow +import org.apache.spark.unsafe.types.VariantVal -class Spark41HoodiePartitionFileSliceMapping(values: InternalRow, - slices: Map[String, FileSlice]) - extends Spark41HoodiePartitionValues(values) - with HoodiePartitionFileSliceMapping { +/** + * Base class for Spark 4.x HoodiePartitionValues implementations. + * Adds the delegation logic available in all Spark 4.x versions on top of + * [[BaseHoodiePartitionValues]]. Version-specific subclasses only implement + * `copy()` and the getters introduced by a newer Spark version. + */ +abstract class Spark4HoodiePartitionValues(values: InternalRow) + extends BaseHoodiePartitionValues(values) { - override def getSlice(fileId: String): Option[FileSlice] = { - slices.get(fileId) + override def getVariant(ordinal: Int): VariantVal = { + values.getVariant(ordinal) } - - override def getPartitionValues: InternalRow = values } diff --git a/hudi-spark-datasource/hudi-spark4.1.x/src/main/scala/org/apache/hudi/client/model/Spark41HoodieInternalRow.scala b/hudi-spark-datasource/hudi-spark4-common/src/main/scala/org/apache/hudi/client/model/Spark4HoodieInternalRow.scala similarity index 67% copy from hudi-spark-datasource/hudi-spark4.1.x/src/main/scala/org/apache/hudi/client/model/Spark41HoodieInternalRow.scala copy to hudi-spark-datasource/hudi-spark4-common/src/main/scala/org/apache/hudi/client/model/Spark4HoodieInternalRow.scala index a56368eee1e4..dd3b7cd0e3bc 100644 --- a/hudi-spark-datasource/hudi-spark4.1.x/src/main/scala/org/apache/hudi/client/model/Spark41HoodieInternalRow.scala +++ b/hudi-spark-datasource/hudi-spark4-common/src/main/scala/org/apache/hudi/client/model/Spark4HoodieInternalRow.scala @@ -19,9 +19,15 @@ package org.apache.hudi.client.model import org.apache.spark.sql.catalyst.InternalRow -import org.apache.spark.unsafe.types.{GeographyVal, GeometryVal, UTF8String, VariantVal} +import org.apache.spark.unsafe.types.{UTF8String, VariantVal} -class Spark41HoodieInternalRow( +/** + * Base class for Spark 4.x HoodieInternalRow implementations. + * Adds the delegation logic available in all Spark 4.x versions on top of + * [[HoodieInternalRow]]. Version-specific subclasses only implement + * [[newInternalRow]] and the getters introduced by a newer Spark version. + */ +abstract class Spark4HoodieInternalRow( metaFields: Array[UTF8String], sourceRow: InternalRow, sourceContainsMetaFields: Boolean) @@ -32,21 +38,18 @@ class Spark41HoodieInternalRow( sourceRow.getVariant(rebaseOrdinal(ordinal)) } - override def getGeography(ordinal: Int): GeographyVal = { - ruleOutMetaFieldsAccess(ordinal, classOf[GeographyVal]) - sourceRow.getGeography(rebaseOrdinal(ordinal)) - } - - override def getGeometry(ordinal: Int): GeometryVal = { - ruleOutMetaFieldsAccess(ordinal, classOf[GeometryVal]) - sourceRow.getGeometry(rebaseOrdinal(ordinal)) - } - override def copy(): InternalRow = { val copyMetaFields = metaFields.map(f => if (f != null) f.copy() else null) - new Spark41HoodieInternalRow( + newInternalRow( copyMetaFields, if (sourceRow == null) null else sourceRow.copy(), sourceContainsMetaFields) } + + /** + * Creates a new instance of the version-specific row type, used by [[copy]]. + */ + protected def newInternalRow(metaFields: Array[UTF8String], + sourceRow: InternalRow, + sourceContainsMetaFields: Boolean): Spark4HoodieInternalRow } diff --git a/hudi-spark-datasource/hudi-spark4.0.x/src/main/scala/org/apache/hudi/Spark40HoodiePartitionCDCFileGroupMapping.scala b/hudi-spark-datasource/hudi-spark4.0.x/src/main/scala/org/apache/hudi/Spark40HoodiePartitionCDCFileGroupMapping.scala index 58eaac632463..28fcfaa080b2 100644 --- a/hudi-spark-datasource/hudi-spark4.0.x/src/main/scala/org/apache/hudi/Spark40HoodiePartitionCDCFileGroupMapping.scala +++ b/hudi-spark-datasource/hudi-spark4.0.x/src/main/scala/org/apache/hudi/Spark40HoodiePartitionCDCFileGroupMapping.scala @@ -24,11 +24,6 @@ import org.apache.hudi.common.table.cdc.HoodieCDCFileSplit import org.apache.spark.sql.catalyst.InternalRow class Spark40HoodiePartitionCDCFileGroupMapping(partitionValues: InternalRow, - fileSplits: List[HoodieCDCFileSplit]) - extends Spark40HoodiePartitionValues(partitionValues) - with HoodiePartitionCDCFileGroupMapping { - - override def getFileSplits(): List[HoodieCDCFileSplit] = { - fileSplits - } -} + protected val fileSplits: List[HoodieCDCFileSplit]) + extends Spark40HoodiePartitionValues(partitionValues) + with Spark4HoodiePartitionCDCFileGroupMapping diff --git a/hudi-spark-datasource/hudi-spark4.0.x/src/main/scala/org/apache/hudi/Spark40HoodiePartitionFileSliceMapping.scala b/hudi-spark-datasource/hudi-spark4.0.x/src/main/scala/org/apache/hudi/Spark40HoodiePartitionFileSliceMapping.scala index 0f769f5bd7ea..7de6dbf71f95 100644 --- a/hudi-spark-datasource/hudi-spark4.0.x/src/main/scala/org/apache/hudi/Spark40HoodiePartitionFileSliceMapping.scala +++ b/hudi-spark-datasource/hudi-spark4.0.x/src/main/scala/org/apache/hudi/Spark40HoodiePartitionFileSliceMapping.scala @@ -24,13 +24,6 @@ import org.apache.hudi.common.model.FileSlice import org.apache.spark.sql.catalyst.InternalRow class Spark40HoodiePartitionFileSliceMapping(values: InternalRow, - slices: Map[String, FileSlice]) + protected val slices: Map[String, FileSlice]) extends Spark40HoodiePartitionValues(values) - with HoodiePartitionFileSliceMapping { - - override def getSlice(fileId: String): Option[FileSlice] = { - slices.get(fileId) - } - - override def getPartitionValues: InternalRow = values -} + with Spark4HoodiePartitionFileSliceMapping diff --git a/hudi-spark-datasource/hudi-spark4.0.x/src/main/scala/org/apache/hudi/Spark40HoodiePartitionValues.scala b/hudi-spark-datasource/hudi-spark4.0.x/src/main/scala/org/apache/hudi/Spark40HoodiePartitionValues.scala index db6e3f10341a..360a6aa8996e 100644 --- a/hudi-spark-datasource/hudi-spark4.0.x/src/main/scala/org/apache/hudi/Spark40HoodiePartitionValues.scala +++ b/hudi-spark-datasource/hudi-spark4.0.x/src/main/scala/org/apache/hudi/Spark40HoodiePartitionValues.scala @@ -20,92 +20,11 @@ package org.apache.hudi import org.apache.spark.sql.catalyst.InternalRow -import org.apache.spark.sql.catalyst.util.{ArrayData, MapData} -import org.apache.spark.sql.types.{DataType, Decimal} -import org.apache.spark.unsafe.types.{CalendarInterval, UTF8String, VariantVal} -case class Spark40HoodiePartitionValues(values: InternalRow) extends HoodiePartitionValues { - override def numFields: Int = { - values.numFields - } - - override def setNullAt(i: Int): Unit = { - values.setNullAt(i) - } - - override def update(i: Int, value: Any): Unit = { - values.update(i, value) - } +case class Spark40HoodiePartitionValues(override val values: InternalRow) + extends Spark4HoodiePartitionValues(values) { override def copy(): InternalRow = { Spark40HoodiePartitionValues(values.copy()) } - - override def isNullAt(ordinal: Int): Boolean = { - values.isNullAt(ordinal) - } - - override def getBoolean(ordinal: Int): Boolean = { - values.getBoolean(ordinal) - } - - override def getByte(ordinal: Int): Byte = { - values.getByte(ordinal) - } - - override def getShort(ordinal: Int): Short = { - values.getShort(ordinal) - } - - override def getInt(ordinal: Int): Int = { - values.getInt(ordinal) - } - - override def getLong(ordinal: Int): Long = { - values.getLong(ordinal) - } - - override def getFloat(ordinal: Int): Float = { - values.getFloat(ordinal) - } - - override def getDouble(ordinal: Int): Double = { - values.getDouble(ordinal) - } - - override def getDecimal(ordinal: Int, precision: Int, scale: Int): Decimal = { - values.getDecimal(ordinal, precision, scale) - } - - override def getUTF8String(ordinal: Int): UTF8String = { - values.getUTF8String(ordinal) - } - - override def getBinary(ordinal: Int): Array[Byte] = { - values.getBinary(ordinal) - } - - override def getInterval(ordinal: Int): CalendarInterval = { - values.getInterval(ordinal) - } - - override def getVariant(ordinal: Int): VariantVal = { - values.getVariant(ordinal) - } - - override def getStruct(ordinal: Int, numFields: Int): InternalRow = { - values.getStruct(ordinal, numFields) - } - - override def getArray(ordinal: Int): ArrayData = { - values.getArray(ordinal) - } - - override def getMap(ordinal: Int): MapData = { - values.getMap(ordinal) - } - - override def get(ordinal: Int, dataType: DataType): AnyRef = { - values.get(ordinal, dataType) - } } diff --git a/hudi-spark-datasource/hudi-spark4.0.x/src/main/scala/org/apache/hudi/client/model/Spark40HoodieInternalRow.scala b/hudi-spark-datasource/hudi-spark4.0.x/src/main/scala/org/apache/hudi/client/model/Spark40HoodieInternalRow.scala index 1da628d03a51..6b0ce70edd45 100644 --- a/hudi-spark-datasource/hudi-spark4.0.x/src/main/scala/org/apache/hudi/client/model/Spark40HoodieInternalRow.scala +++ b/hudi-spark-datasource/hudi-spark4.0.x/src/main/scala/org/apache/hudi/client/model/Spark40HoodieInternalRow.scala @@ -19,24 +19,17 @@ package org.apache.hudi.client.model import org.apache.spark.sql.catalyst.InternalRow -import org.apache.spark.unsafe.types.{UTF8String, VariantVal} +import org.apache.spark.unsafe.types.UTF8String class Spark40HoodieInternalRow( - metaFields: Array[UTF8String], - sourceRow: InternalRow, - sourceContainsMetaFields: Boolean) - extends HoodieInternalRow(metaFields, sourceRow, sourceContainsMetaFields) { + metaFields: Array[UTF8String], + sourceRow: InternalRow, + sourceContainsMetaFields: Boolean) + extends Spark4HoodieInternalRow(metaFields, sourceRow, sourceContainsMetaFields) { - override def getVariant(ordinal: Int): VariantVal = { - ruleOutMetaFieldsAccess(ordinal, classOf[VariantVal]) - sourceRow.getVariant(rebaseOrdinal(ordinal)) - } - - override def copy(): InternalRow = { - val copyMetaFields = metaFields.map(f => if (f != null) f.copy() else null) - new Spark40HoodieInternalRow( - copyMetaFields, - if (sourceRow == null) null else sourceRow.copy(), - sourceContainsMetaFields) + override protected def newInternalRow(metaFields: Array[UTF8String], + sourceRow: InternalRow, + sourceContainsMetaFields: Boolean): Spark4HoodieInternalRow = { + new Spark40HoodieInternalRow(metaFields, sourceRow, sourceContainsMetaFields) } } diff --git a/hudi-spark-datasource/hudi-spark4.1.x/src/main/scala/org/apache/hudi/Spark41HoodiePartitionCDCFileGroupMapping.scala b/hudi-spark-datasource/hudi-spark4.1.x/src/main/scala/org/apache/hudi/Spark41HoodiePartitionCDCFileGroupMapping.scala index 005c17eb4128..96f53d886ba5 100644 --- a/hudi-spark-datasource/hudi-spark4.1.x/src/main/scala/org/apache/hudi/Spark41HoodiePartitionCDCFileGroupMapping.scala +++ b/hudi-spark-datasource/hudi-spark4.1.x/src/main/scala/org/apache/hudi/Spark41HoodiePartitionCDCFileGroupMapping.scala @@ -24,11 +24,6 @@ import org.apache.hudi.common.table.cdc.HoodieCDCFileSplit import org.apache.spark.sql.catalyst.InternalRow class Spark41HoodiePartitionCDCFileGroupMapping(partitionValues: InternalRow, - fileSplits: List[HoodieCDCFileSplit]) + protected val fileSplits: List[HoodieCDCFileSplit]) extends Spark41HoodiePartitionValues(partitionValues) - with HoodiePartitionCDCFileGroupMapping { - - override def getFileSplits(): List[HoodieCDCFileSplit] = { - fileSplits - } -} + with Spark4HoodiePartitionCDCFileGroupMapping diff --git a/hudi-spark-datasource/hudi-spark4.1.x/src/main/scala/org/apache/hudi/Spark41HoodiePartitionFileSliceMapping.scala b/hudi-spark-datasource/hudi-spark4.1.x/src/main/scala/org/apache/hudi/Spark41HoodiePartitionFileSliceMapping.scala index 07e199e4a9b1..fbd66d8fc8a7 100644 --- a/hudi-spark-datasource/hudi-spark4.1.x/src/main/scala/org/apache/hudi/Spark41HoodiePartitionFileSliceMapping.scala +++ b/hudi-spark-datasource/hudi-spark4.1.x/src/main/scala/org/apache/hudi/Spark41HoodiePartitionFileSliceMapping.scala @@ -24,13 +24,6 @@ import org.apache.hudi.common.model.FileSlice import org.apache.spark.sql.catalyst.InternalRow class Spark41HoodiePartitionFileSliceMapping(values: InternalRow, - slices: Map[String, FileSlice]) + protected val slices: Map[String, FileSlice]) extends Spark41HoodiePartitionValues(values) - with HoodiePartitionFileSliceMapping { - - override def getSlice(fileId: String): Option[FileSlice] = { - slices.get(fileId) - } - - override def getPartitionValues: InternalRow = values -} + with Spark4HoodiePartitionFileSliceMapping diff --git a/hudi-spark-datasource/hudi-spark4.1.x/src/main/scala/org/apache/hudi/Spark41HoodiePartitionValues.scala b/hudi-spark-datasource/hudi-spark4.1.x/src/main/scala/org/apache/hudi/Spark41HoodiePartitionValues.scala index 7d2c71717d57..3964352f3cd8 100644 --- a/hudi-spark-datasource/hudi-spark4.1.x/src/main/scala/org/apache/hudi/Spark41HoodiePartitionValues.scala +++ b/hudi-spark-datasource/hudi-spark4.1.x/src/main/scala/org/apache/hudi/Spark41HoodiePartitionValues.scala @@ -20,71 +20,15 @@ package org.apache.hudi import org.apache.spark.sql.catalyst.InternalRow -import org.apache.spark.sql.catalyst.util.{ArrayData, MapData} -import org.apache.spark.sql.types.{DataType, Decimal} -import org.apache.spark.unsafe.types.{CalendarInterval, GeographyVal, GeometryVal, UTF8String, VariantVal} +import org.apache.spark.unsafe.types.{GeographyVal, GeometryVal} -case class Spark41HoodiePartitionValues(values: InternalRow) extends HoodiePartitionValues { - override def numFields: Int = { - values.numFields - } - - override def setNullAt(i: Int): Unit = { - values.setNullAt(i) - } - - override def update(i: Int, value: Any): Unit = { - values.update(i, value) - } +case class Spark41HoodiePartitionValues(override val values: InternalRow) + extends Spark4HoodiePartitionValues(values) { override def copy(): InternalRow = { Spark41HoodiePartitionValues(values.copy()) } - override def isNullAt(ordinal: Int): Boolean = { - values.isNullAt(ordinal) - } - - override def getBoolean(ordinal: Int): Boolean = { - values.getBoolean(ordinal) - } - - override def getByte(ordinal: Int): Byte = { - values.getByte(ordinal) - } - - override def getShort(ordinal: Int): Short = { - values.getShort(ordinal) - } - - override def getInt(ordinal: Int): Int = { - values.getInt(ordinal) - } - - override def getLong(ordinal: Int): Long = { - values.getLong(ordinal) - } - - override def getFloat(ordinal: Int): Float = { - values.getFloat(ordinal) - } - - override def getDouble(ordinal: Int): Double = { - values.getDouble(ordinal) - } - - override def getDecimal(ordinal: Int, precision: Int, scale: Int): Decimal = { - values.getDecimal(ordinal, precision, scale) - } - - override def getUTF8String(ordinal: Int): UTF8String = { - values.getUTF8String(ordinal) - } - - override def getBinary(ordinal: Int): Array[Byte] = { - values.getBinary(ordinal) - } - override def getGeography(ordinal: Int): GeographyVal = { values.getGeography(ordinal) } @@ -92,28 +36,4 @@ case class Spark41HoodiePartitionValues(values: InternalRow) extends HoodieParti override def getGeometry(ordinal: Int): GeometryVal = { values.getGeometry(ordinal) } - - override def getInterval(ordinal: Int): CalendarInterval = { - values.getInterval(ordinal) - } - - override def getVariant(ordinal: Int): VariantVal = { - values.getVariant(ordinal) - } - - override def getStruct(ordinal: Int, numFields: Int): InternalRow = { - values.getStruct(ordinal, numFields) - } - - override def getArray(ordinal: Int): ArrayData = { - values.getArray(ordinal) - } - - override def getMap(ordinal: Int): MapData = { - values.getMap(ordinal) - } - - override def get(ordinal: Int, dataType: DataType): AnyRef = { - values.get(ordinal, dataType) - } } diff --git a/hudi-spark-datasource/hudi-spark4.1.x/src/main/scala/org/apache/hudi/client/model/Spark41HoodieInternalRow.scala b/hudi-spark-datasource/hudi-spark4.1.x/src/main/scala/org/apache/hudi/client/model/Spark41HoodieInternalRow.scala index a56368eee1e4..308424396a75 100644 --- a/hudi-spark-datasource/hudi-spark4.1.x/src/main/scala/org/apache/hudi/client/model/Spark41HoodieInternalRow.scala +++ b/hudi-spark-datasource/hudi-spark4.1.x/src/main/scala/org/apache/hudi/client/model/Spark41HoodieInternalRow.scala @@ -19,18 +19,13 @@ package org.apache.hudi.client.model import org.apache.spark.sql.catalyst.InternalRow -import org.apache.spark.unsafe.types.{GeographyVal, GeometryVal, UTF8String, VariantVal} +import org.apache.spark.unsafe.types.{GeographyVal, GeometryVal, UTF8String} class Spark41HoodieInternalRow( metaFields: Array[UTF8String], sourceRow: InternalRow, sourceContainsMetaFields: Boolean) - extends HoodieInternalRow(metaFields, sourceRow, sourceContainsMetaFields) { - - override def getVariant(ordinal: Int): VariantVal = { - ruleOutMetaFieldsAccess(ordinal, classOf[VariantVal]) - sourceRow.getVariant(rebaseOrdinal(ordinal)) - } + extends Spark4HoodieInternalRow(metaFields, sourceRow, sourceContainsMetaFields) { override def getGeography(ordinal: Int): GeographyVal = { ruleOutMetaFieldsAccess(ordinal, classOf[GeographyVal]) @@ -42,11 +37,9 @@ class Spark41HoodieInternalRow( sourceRow.getGeometry(rebaseOrdinal(ordinal)) } - override def copy(): InternalRow = { - val copyMetaFields = metaFields.map(f => if (f != null) f.copy() else null) - new Spark41HoodieInternalRow( - copyMetaFields, - if (sourceRow == null) null else sourceRow.copy(), - sourceContainsMetaFields) + override protected def newInternalRow(metaFields: Array[UTF8String], + sourceRow: InternalRow, + sourceContainsMetaFields: Boolean): Spark4HoodieInternalRow = { + new Spark41HoodieInternalRow(metaFields, sourceRow, sourceContainsMetaFields) } }
