This is an automated email from the ASF dual-hosted git repository.
wombatu-kun 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 39aa022c8f26 refactor(spark): share Spark 4.x partition-values and
mapping classes in hudi-spark4-common (#19148)
39aa022c8f26 is described below
commit 39aa022c8f2686c0aa0e9745da1ed203ca35e581
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.
---
.../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 ++---
...Spark42HoodiePartitionCDCFileGroupMapping.scala | 9 +--
.../Spark42HoodiePartitionFileSliceMapping.scala | 11 +--
.../apache/hudi/Spark42HoodiePartitionValues.scala | 86 +---------------------
.../client/model/Spark42HoodieInternalRow.scala | 19 ++---
18 files changed, 92 insertions(+), 456 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.2.x/src/main/scala/org/apache/hudi/Spark42HoodiePartitionCDCFileGroupMapping.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.2.x/src/main/scala/org/apache/hudi/Spark42HoodiePartitionCDCFileGroupMapping.scala
copy to
hudi-spark-datasource/hudi-spark4-common/src/main/scala/org/apache/hudi/Spark4HoodiePartitionCDCFileGroupMapping.scala
index 50c2c584bd30..428ad6a14112 100644
---
a/hudi-spark-datasource/hudi-spark4.2.x/src/main/scala/org/apache/hudi/Spark42HoodiePartitionCDCFileGroupMapping.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 Spark42HoodiePartitionCDCFileGroupMapping(partitionValues: InternalRow,
- fileSplits:
List[HoodieCDCFileSplit])
- extends Spark42HoodiePartitionValues(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.2.x/src/main/scala/org/apache/hudi/client/model/Spark42HoodieInternalRow.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.2.x/src/main/scala/org/apache/hudi/client/model/Spark42HoodieInternalRow.scala
copy to
hudi-spark-datasource/hudi-spark4-common/src/main/scala/org/apache/hudi/client/model/Spark4HoodieInternalRow.scala
index 115c2e7151f1..dd3b7cd0e3bc 100644
---
a/hudi-spark-datasource/hudi-spark4.2.x/src/main/scala/org/apache/hudi/client/model/Spark42HoodieInternalRow.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 Spark42HoodieInternalRow(
+/**
+ * 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 Spark42HoodieInternalRow(
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 Spark42HoodieInternalRow(
+ 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)
}
}
diff --git
a/hudi-spark-datasource/hudi-spark4.2.x/src/main/scala/org/apache/hudi/Spark42HoodiePartitionCDCFileGroupMapping.scala
b/hudi-spark-datasource/hudi-spark4.2.x/src/main/scala/org/apache/hudi/Spark42HoodiePartitionCDCFileGroupMapping.scala
index 50c2c584bd30..d2874fabc356 100644
---
a/hudi-spark-datasource/hudi-spark4.2.x/src/main/scala/org/apache/hudi/Spark42HoodiePartitionCDCFileGroupMapping.scala
+++
b/hudi-spark-datasource/hudi-spark4.2.x/src/main/scala/org/apache/hudi/Spark42HoodiePartitionCDCFileGroupMapping.scala
@@ -24,11 +24,6 @@ import org.apache.hudi.common.table.cdc.HoodieCDCFileSplit
import org.apache.spark.sql.catalyst.InternalRow
class Spark42HoodiePartitionCDCFileGroupMapping(partitionValues: InternalRow,
- fileSplits:
List[HoodieCDCFileSplit])
+ protected val fileSplits:
List[HoodieCDCFileSplit])
extends Spark42HoodiePartitionValues(partitionValues)
- with HoodiePartitionCDCFileGroupMapping {
-
- override def getFileSplits(): List[HoodieCDCFileSplit] = {
- fileSplits
- }
-}
+ with Spark4HoodiePartitionCDCFileGroupMapping
diff --git
a/hudi-spark-datasource/hudi-spark4.2.x/src/main/scala/org/apache/hudi/Spark42HoodiePartitionFileSliceMapping.scala
b/hudi-spark-datasource/hudi-spark4.2.x/src/main/scala/org/apache/hudi/Spark42HoodiePartitionFileSliceMapping.scala
index f4b3838c7947..b9396ca5f8da 100644
---
a/hudi-spark-datasource/hudi-spark4.2.x/src/main/scala/org/apache/hudi/Spark42HoodiePartitionFileSliceMapping.scala
+++
b/hudi-spark-datasource/hudi-spark4.2.x/src/main/scala/org/apache/hudi/Spark42HoodiePartitionFileSliceMapping.scala
@@ -24,13 +24,6 @@ import org.apache.hudi.common.model.FileSlice
import org.apache.spark.sql.catalyst.InternalRow
class Spark42HoodiePartitionFileSliceMapping(values: InternalRow,
- slices: Map[String, FileSlice])
+ protected val slices: Map[String,
FileSlice])
extends Spark42HoodiePartitionValues(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.2.x/src/main/scala/org/apache/hudi/Spark42HoodiePartitionValues.scala
b/hudi-spark-datasource/hudi-spark4.2.x/src/main/scala/org/apache/hudi/Spark42HoodiePartitionValues.scala
index 831475d01364..aa53fa4a5b58 100644
---
a/hudi-spark-datasource/hudi-spark4.2.x/src/main/scala/org/apache/hudi/Spark42HoodiePartitionValues.scala
+++
b/hudi-spark-datasource/hudi-spark4.2.x/src/main/scala/org/apache/hudi/Spark42HoodiePartitionValues.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 Spark42HoodiePartitionValues(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 Spark42HoodiePartitionValues(override val values: InternalRow)
+ extends Spark4HoodiePartitionValues(values) {
override def copy(): InternalRow = {
Spark42HoodiePartitionValues(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 Spark42HoodiePartitionValues(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.2.x/src/main/scala/org/apache/hudi/client/model/Spark42HoodieInternalRow.scala
b/hudi-spark-datasource/hudi-spark4.2.x/src/main/scala/org/apache/hudi/client/model/Spark42HoodieInternalRow.scala
index 115c2e7151f1..74a5355404b4 100644
---
a/hudi-spark-datasource/hudi-spark4.2.x/src/main/scala/org/apache/hudi/client/model/Spark42HoodieInternalRow.scala
+++
b/hudi-spark-datasource/hudi-spark4.2.x/src/main/scala/org/apache/hudi/client/model/Spark42HoodieInternalRow.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 Spark42HoodieInternalRow(
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 Spark42HoodieInternalRow(
sourceRow.getGeometry(rebaseOrdinal(ordinal))
}
- override def copy(): InternalRow = {
- val copyMetaFields = metaFields.map(f => if (f != null) f.copy() else null)
- new Spark42HoodieInternalRow(
- copyMetaFields,
- if (sourceRow == null) null else sourceRow.copy(),
- sourceContainsMetaFields)
+ override protected def newInternalRow(metaFields: Array[UTF8String],
+ sourceRow: InternalRow,
+ sourceContainsMetaFields: Boolean):
Spark4HoodieInternalRow = {
+ new Spark42HoodieInternalRow(metaFields, sourceRow,
sourceContainsMetaFields)
}
}