This is an automated email from the ASF dual-hosted git repository.
JingsongLi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/paimon.git
The following commit(s) were added to refs/heads/master by this push:
new d637d6507a [spark] data-evolution merge-into supports updating
Map/Array Blobs (#8796)
d637d6507a is described below
commit d637d6507a818061b7be2f072114805695b661f6
Author: Faiz <[email protected]>
AuthorDate: Wed Jul 22 17:46:00 2026 +0800
[spark] data-evolution merge-into supports updating Map/Array Blobs (#8796)
---
.../MergeIntoPaimonDataEvolutionTable.scala | 21 +--
.../MergeIntoPaimonDataEvolutionTable.scala | 22 +--
.../spark/write/DataEvolutionTableDataWrite.scala | 33 +++-
.../paimon/spark/sql/BlobUpdateTestBase.scala | 206 ++++++++++++++++++---
.../spark/sql/DataEvolutionDeletionTestBase.scala | 20 +-
5 files changed, 221 insertions(+), 81 deletions(-)
diff --git
a/paimon-spark/paimon-spark-4.0/src/main/scala/org/apache/paimon/spark/commands/MergeIntoPaimonDataEvolutionTable.scala
b/paimon-spark/paimon-spark-4.0/src/main/scala/org/apache/paimon/spark/commands/MergeIntoPaimonDataEvolutionTable.scala
index 4e4b21f4e2..5690a65d5f 100644
---
a/paimon-spark/paimon-spark-4.0/src/main/scala/org/apache/paimon/spark/commands/MergeIntoPaimonDataEvolutionTable.scala
+++
b/paimon-spark/paimon-spark-4.0/src/main/scala/org/apache/paimon/spark/commands/MergeIntoPaimonDataEvolutionTable.scala
@@ -39,7 +39,7 @@ import org.apache.paimon.table.sink.{CommitMessage,
CommitMessageImpl}
import org.apache.paimon.table.source.DataSplit
import org.apache.paimon.table.source.snapshot.SnapshotReader
import org.apache.paimon.table.source.snapshot.TimeTravelUtil
-import org.apache.paimon.types.{BlobType, DataTypeRoot, RowType}
+import org.apache.paimon.types.{BlobType, RowType}
import org.apache.paimon.types.VectorType.isVectorStoreFile
import org.apache.spark.internal.Logging
@@ -430,14 +430,6 @@ case class MergeIntoPaimonDataEvolutionTable(
val rawBlobFieldNames = rawBlobFields
.map(_.name())
.toSet
- val rawNestedBlobFieldNames = rawBlobFields
- .filter(
- field => {
- val root = field.`type`().getTypeRoot
- root == DataTypeRoot.ARRAY || root == DataTypeRoot.MAP
- })
- .map(_.name())
- .toSet
def isRawBlobUpdateColumn(attr: AttributeReference): Boolean = {
rawBlobFieldNames.exists(rawBlobFieldName => resolver(rawBlobFieldName,
attr.name))
@@ -459,17 +451,6 @@ case class MergeIntoPaimonDataEvolutionTable(
}.toSet
}
- val modifiedRawNestedBlobColumnNames = matchedActions
- .collect { case action: UpdateAction => modifiedRawBlobNames(action) }
- .flatten
- .filter(name => rawNestedBlobFieldNames.exists(blobName =>
resolver(blobName, name)))
- .toSet
- if (modifiedRawNestedBlobColumnNames.nonEmpty) {
- throw new UnsupportedOperationException(
- "Should not append/update raw-data ARRAY<BLOB> or MAP<X, BLOB> column
through MERGE INTO: " +
- modifiedRawNestedBlobColumnNames.toSeq.sorted.mkString(", "))
- }
-
val rawBlobMarkerNames =
rawBlobMarkerNamesAvoiding(
rawBlobUpdateColumns.size,
diff --git
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/MergeIntoPaimonDataEvolutionTable.scala
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/MergeIntoPaimonDataEvolutionTable.scala
index 4e4b21f4e2..f287d31959 100644
---
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/MergeIntoPaimonDataEvolutionTable.scala
+++
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/MergeIntoPaimonDataEvolutionTable.scala
@@ -39,7 +39,7 @@ import org.apache.paimon.table.sink.{CommitMessage,
CommitMessageImpl}
import org.apache.paimon.table.source.DataSplit
import org.apache.paimon.table.source.snapshot.SnapshotReader
import org.apache.paimon.table.source.snapshot.TimeTravelUtil
-import org.apache.paimon.types.{BlobType, DataTypeRoot, RowType}
+import org.apache.paimon.types.{BlobType, RowType}
import org.apache.paimon.types.VectorType.isVectorStoreFile
import org.apache.spark.internal.Logging
@@ -430,15 +430,6 @@ case class MergeIntoPaimonDataEvolutionTable(
val rawBlobFieldNames = rawBlobFields
.map(_.name())
.toSet
- val rawNestedBlobFieldNames = rawBlobFields
- .filter(
- field => {
- val root = field.`type`().getTypeRoot
- root == DataTypeRoot.ARRAY || root == DataTypeRoot.MAP
- })
- .map(_.name())
- .toSet
-
def isRawBlobUpdateColumn(attr: AttributeReference): Boolean = {
rawBlobFieldNames.exists(rawBlobFieldName => resolver(rawBlobFieldName,
attr.name))
}
@@ -459,17 +450,6 @@ case class MergeIntoPaimonDataEvolutionTable(
}.toSet
}
- val modifiedRawNestedBlobColumnNames = matchedActions
- .collect { case action: UpdateAction => modifiedRawBlobNames(action) }
- .flatten
- .filter(name => rawNestedBlobFieldNames.exists(blobName =>
resolver(blobName, name)))
- .toSet
- if (modifiedRawNestedBlobColumnNames.nonEmpty) {
- throw new UnsupportedOperationException(
- "Should not append/update raw-data ARRAY<BLOB> or MAP<X, BLOB> column
through MERGE INTO: " +
- modifiedRawNestedBlobColumnNames.toSeq.sorted.mkString(", "))
- }
-
val rawBlobMarkerNames =
rawBlobMarkerNamesAvoiding(
rawBlobUpdateColumns.size,
diff --git
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/write/DataEvolutionTableDataWrite.scala
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/write/DataEvolutionTableDataWrite.scala
index 19b66785c7..1ba681c813 100644
---
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/write/DataEvolutionTableDataWrite.scala
+++
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/write/DataEvolutionTableDataWrite.scala
@@ -60,7 +60,6 @@ case class DataEvolutionTableDataWrite(
private val toPaimonRow = {
SparkRowUtils.toPaimonRow(writeType, -1, uriReaderFactory)
}
- private lazy val rowSerializer = InternalSerializers.create(writeType)
private val rawBlobFallbackFields =
rawBlobPlaceholderMarkerIndexes.toSeq.sortBy(_._1).toArray
private val rawBlobFallbackFieldIndexes = rawBlobFallbackFields.map(_._1)
private val rawBlobFallbackMappings = {
@@ -72,6 +71,7 @@ case class DataEvolutionTableDataWrite(
mappings
}
private val rawBlobFallbackMarkerIndexes = rawBlobFallbackFields.map(_._2)
+ private val rawBlobFallbackFieldIndexSet = rawBlobFallbackFieldIndexes.toSet
private val rawBlobFallbackPlaceholders: Array[AnyRef] =
rawBlobFallbackFields.map {
case (fieldIndex, _) =>
if (writeType.getTypeAt(fieldIndex).getTypeRoot == DataTypeRoot.ARRAY) {
@@ -84,6 +84,25 @@ case class DataEvolutionTableDataWrite(
}
private val rawBlobFallbackRow = new
GenericRow(rawBlobFallbackMarkerIndexes.length)
private val rawBlobFallbackMappingRow = new
FallbackMappingRow(rawBlobFallbackMappings)
+ private val fillerFieldCopiers =
writeType.getFields.asScala.zipWithIndex.collect {
+ case (field, index) if !rawBlobFallbackFieldIndexSet.contains(index) =>
+ (
+ index,
+ InternalRow.createFieldGetter(field.`type`(), index),
+ InternalSerializers.create[AnyRef](field.`type`()))
+ }.toArray
+
+ private def createFillerRow(source: InternalRow): GenericRow = {
+ val filler = new GenericRow(writeType.getFieldCount)
+ filler.setRowKind(source.getRowKind)
+ // only copy non-blob fields
+ fillerFieldCopiers.foreach {
+ case (index, getter, serializer) =>
+ val value = getter.getFieldOrNull(source)
+ filler.setField(index, if (value == null) null else
serializer.copy(value))
+ }
+ filler
+ }
def write(row: Row): Unit = {
val firstRowId = row.getLong(firstRowIdIndex)
@@ -224,14 +243,10 @@ case class DataEvolutionTableDataWrite(
private def fillGapUntil(rowId: Long, fillerSourceRow: InternalRow =
null): Unit = {
if (fillerRow == null && fillerSourceRow != null) {
- // Copy the first record this file writer met to minimize the
influences on
- // file stats, but keep raw blob fields as NULLs so filler rows do not
trigger
- // blob fallback.
- val copied = rowSerializer
- .copyRowData(fillerSourceRow, new
GenericRow(writeType.getFieldCount))
- .asInstanceOf[GenericRow]
- rawBlobFallbackFieldIndexes.foreach(copied.setField(_, null))
- fillerRow = copied
+ // Cache the first record eagerly because finish may need it to fill a
trailing gap.
+ // Copy it to minimize the influences on file stats, but keep raw blob
fields as
+ // NULLs so filler rows do not trigger blob fallback.
+ fillerRow = createFillerRow(fillerSourceRow)
}
assert(
diff --git
a/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/BlobUpdateTestBase.scala
b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/BlobUpdateTestBase.scala
index cc784cdb12..f9d6597c94 100644
---
a/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/BlobUpdateTestBase.scala
+++
b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/BlobUpdateTestBase.scala
@@ -18,6 +18,7 @@
package org.apache.paimon.spark.sql
+import org.apache.paimon.data.BlobDescriptor
import org.apache.paimon.spark.PaimonSparkTestBase
import org.apache.spark.SparkConf
@@ -59,45 +60,93 @@ class BlobUpdateTestBase extends PaimonSparkTestBase {
}
}
- test("Blob: merge-into rejects updating raw-data array blob column") {
+ test("Blob: merge-into updates raw-data array blob column") {
withTable("s", "t") {
sql("CREATE TABLE t (id INT, name STRING, pictures ARRAY<BINARY>)
TBLPROPERTIES " +
"('row-tracking.enabled'='true', 'data-evolution.enabled'='true',
'blob-field'='pictures')")
- sql("INSERT INTO t VALUES (1, 'name1', array(X'48656C6C6F'))")
-
- sql("CREATE TABLE s (id INT, pictures ARRAY<BINARY>)")
- sql("INSERT INTO s VALUES (1, array(X'4E4557'))")
-
- val e = intercept[UnsupportedOperationException] {
- sql("""
- |MERGE INTO t
- |USING s
- |ON t.id = s.id
- |WHEN MATCHED THEN UPDATE SET t.pictures = s.pictures
- |""".stripMargin)
- }
- assert(e.getMessage.contains("raw-data ARRAY<BLOB>"))
+ sql(
+ "INSERT INTO t VALUES " +
+ "(1, 'name1', array(X'48656C6C6F')), " +
+ "(2, 'name2', array(X'414243')), " +
+ "(3, 'name3', CAST(NULL AS ARRAY<BINARY>))")
+
+ sql("CREATE TABLE s (id INT, pictures ARRAY<BINARY>, picture BINARY)")
+ sql(
+ "INSERT INTO s VALUES " +
+ "(1, array(X'4E4557', CAST(NULL AS BINARY)), X'21'), " +
+ "(3, CAST(array() AS ARRAY<BINARY>), CAST(NULL AS BINARY))")
+
+ sql("""
+ |MERGE INTO t
+ |USING s
+ |ON t.id = s.id
+ |WHEN MATCHED AND s.id = 1 THEN UPDATE SET
+ | t.pictures = concat(t.pictures, s.pictures, array(s.picture))
+ |WHEN MATCHED THEN UPDATE SET t.pictures = s.pictures
+ |""".stripMargin)
+
+ val updated = sql(
+ "SELECT size(pictures), pictures[0], pictures[1], pictures[2] IS NULL,
pictures[3] " +
+ "FROM t WHERE id = 1").collect()(0)
+ assert(updated.getInt(0) == 4)
+ assert(
+ java.util.Arrays.equals(updated.getAs[Array[Byte]](1), Array[Byte](72,
101, 108, 108, 111)))
+ assert(java.util.Arrays.equals(updated.getAs[Array[Byte]](2),
Array[Byte](78, 69, 87)))
+ assert(updated.getBoolean(3))
+ assert(java.util.Arrays.equals(updated.getAs[Array[Byte]](4),
Array[Byte](33)))
+
+ val unchanged = sql("SELECT size(pictures), pictures[0] FROM t WHERE id
= 2").collect()(0)
+ assert(unchanged.getInt(0) == 1)
+ assert(java.util.Arrays.equals(unchanged.getAs[Array[Byte]](1),
Array[Byte](65, 66, 67)))
+
+ checkAnswer(
+ sql("SELECT pictures IS NOT NULL, size(pictures) FROM t WHERE id = 3"),
+ Seq(Row(true, 0)))
}
}
- test("Blob: merge-into rejects updating raw-data map blob column") {
+ test("Blob: merge-into updates raw-data map blob column") {
withTable("s", "t") {
sql("CREATE TABLE t (id INT, name STRING, pictures MAP<STRING, BINARY>)
TBLPROPERTIES " +
"('row-tracking.enabled'='true', 'data-evolution.enabled'='true',
'blob-field'='pictures')")
- sql("INSERT INTO t VALUES (1, 'name1', map('old', X'48656C6C6F'))")
+ sql(
+ "INSERT INTO t VALUES " +
+ "(1, 'name1', map('old', X'48656C6C6F')), " +
+ "(2, 'name2', map('keep', X'414243')), " +
+ "(3, 'name3', CAST(NULL AS MAP<STRING, BINARY>))")
sql("CREATE TABLE s (id INT, pictures MAP<STRING, BINARY>)")
- sql("INSERT INTO s VALUES (1, map('new', X'4E4557'))")
-
- val e = intercept[UnsupportedOperationException] {
- sql("""
- |MERGE INTO t
- |USING s
- |ON t.id = s.id
- |WHEN MATCHED THEN UPDATE SET t.pictures = s.pictures
- |""".stripMargin)
- }
- assert(e.getMessage.contains("MAP<X, BLOB>"))
+ sql(
+ "INSERT INTO s VALUES " +
+ "(1, map('new', X'4E4557', 'missing', CAST(NULL AS BINARY))), " +
+ "(3, CAST(map() AS MAP<STRING, BINARY>))")
+
+ sql("""
+ |MERGE INTO t
+ |USING s
+ |ON t.id = s.id
+ |WHEN MATCHED AND s.id = 1 THEN UPDATE SET
+ | t.pictures = map_concat(t.pictures, s.pictures)
+ |WHEN MATCHED THEN UPDATE SET t.pictures = s.pictures
+ |""".stripMargin)
+
+ val updated = sql(
+ "SELECT size(pictures), pictures['old'], pictures['new'], " +
+ "pictures['missing'] IS NULL FROM t WHERE id = 1").collect()(0)
+ assert(updated.getInt(0) == 3)
+ assert(
+ java.util.Arrays.equals(updated.getAs[Array[Byte]](1), Array[Byte](72,
101, 108, 108, 111)))
+ assert(java.util.Arrays.equals(updated.getAs[Array[Byte]](2),
Array[Byte](78, 69, 87)))
+ assert(updated.getBoolean(3))
+
+ val unchanged =
+ sql("SELECT size(pictures), pictures['keep'] FROM t WHERE id =
2").collect()(0)
+ assert(unchanged.getInt(0) == 1)
+ assert(java.util.Arrays.equals(unchanged.getAs[Array[Byte]](1),
Array[Byte](65, 66, 67)))
+
+ checkAnswer(
+ sql("SELECT pictures IS NOT NULL, size(pictures) FROM t WHERE id = 3"),
+ Seq(Row(true, 0)))
}
}
@@ -196,6 +245,107 @@ class BlobUpdateTestBase extends PaimonSparkTestBase {
}
}
+ test("Blob: self merge updates descriptor array blob column") {
+ withTable("t") {
+ sql(
+ "CREATE TABLE t (id INT, pictures ARRAY<BINARY>) TBLPROPERTIES (" +
+ "'row-tracking.enabled'='true', 'data-evolution.enabled'='true', " +
+ "'blob-field'='pictures', 'blob-as-descriptor'='true')")
+ sql(
+ "INSERT INTO t VALUES " +
+ "(1, array(X'48656C6C6F', X'5945')), " +
+ "(2, array(X'414243'))")
+
+ val sourceDescriptor =
+ sql("SELECT pictures[0] FROM t WHERE id =
1").collect()(0).getAs[Array[Byte]](0)
+ assert(BlobDescriptor.deserialize(sourceDescriptor).length() == 5)
+
+ sql("""
+ |MERGE INTO t
+ |USING t AS source
+ |ON t._ROW_ID = source._ROW_ID
+ |WHEN MATCHED AND source.id = 1 THEN UPDATE SET
+ | t.pictures = concat(source.pictures, array(source.pictures[0]))
+ |""".stripMargin)
+
+ val descriptorRow =
+ sql(
+ "SELECT size(pictures), pictures[0], pictures[1], pictures[2] " +
+ "FROM t WHERE id = 1")
+ .collect()(0)
+ assert(descriptorRow.getInt(0) == 3)
+
assert(BlobDescriptor.deserialize(descriptorRow.getAs[Array[Byte]](1)).length()
== 5)
+
assert(BlobDescriptor.deserialize(descriptorRow.getAs[Array[Byte]](2)).length()
== 2)
+
assert(BlobDescriptor.deserialize(descriptorRow.getAs[Array[Byte]](3)).length()
== 5)
+
+ sql("ALTER TABLE t SET TBLPROPERTIES ('blob-as-descriptor'='false')")
+ val rawRow =
+ sql(
+ "SELECT pictures[0], pictures[1], pictures[2] " +
+ "FROM t WHERE id = 1")
+ .collect()(0)
+ assert(
+ java.util.Arrays.equals(rawRow.getAs[Array[Byte]](0), Array[Byte](72,
101, 108, 108, 111)))
+ assert(java.util.Arrays.equals(rawRow.getAs[Array[Byte]](1),
Array[Byte](89, 69)))
+ assert(
+ java.util.Arrays.equals(rawRow.getAs[Array[Byte]](2), Array[Byte](72,
101, 108, 108, 111)))
+
+ val unchanged = sql("SELECT pictures[0] FROM t WHERE id =
2").collect()(0)
+ assert(java.util.Arrays.equals(unchanged.getAs[Array[Byte]](0),
Array[Byte](65, 66, 67)))
+ }
+ }
+
+ test("Blob: self merge updates descriptor map blob column") {
+ withTable("t") {
+ sql(
+ "CREATE TABLE t (id INT, pictures MAP<STRING, BINARY>) TBLPROPERTIES
(" +
+ "'row-tracking.enabled'='true', 'data-evolution.enabled'='true', " +
+ "'blob-field'='pictures', 'blob-as-descriptor'='true')")
+ sql(
+ "INSERT INTO t VALUES " +
+ "(1, map('old', X'48656C6C6F', 'small', X'5945')), " +
+ "(2, map('keep', X'414243'))")
+
+ val sourceDescriptor =
+ sql("SELECT pictures['old'] FROM t WHERE id =
1").collect()(0).getAs[Array[Byte]](0)
+ assert(BlobDescriptor.deserialize(sourceDescriptor).length() == 5)
+
+ sql("""
+ |MERGE INTO t
+ |USING t AS source
+ |ON t._ROW_ID = source._ROW_ID
+ |WHEN MATCHED AND source.id = 1 THEN UPDATE SET
+ | t.pictures = map_concat(
+ | source.pictures, map('copied', source.pictures['old']))
+ |""".stripMargin)
+
+ val descriptorRow =
+ sql(
+ "SELECT size(pictures), pictures['old'], pictures['small'],
pictures['copied'] " +
+ "FROM t WHERE id = 1")
+ .collect()(0)
+ assert(descriptorRow.getInt(0) == 3)
+
assert(BlobDescriptor.deserialize(descriptorRow.getAs[Array[Byte]](1)).length()
== 5)
+
assert(BlobDescriptor.deserialize(descriptorRow.getAs[Array[Byte]](2)).length()
== 2)
+
assert(BlobDescriptor.deserialize(descriptorRow.getAs[Array[Byte]](3)).length()
== 5)
+
+ sql("ALTER TABLE t SET TBLPROPERTIES ('blob-as-descriptor'='false')")
+ val rawRow =
+ sql(
+ "SELECT pictures['old'], pictures['small'], pictures['copied'] " +
+ "FROM t WHERE id = 1")
+ .collect()(0)
+ assert(
+ java.util.Arrays.equals(rawRow.getAs[Array[Byte]](0), Array[Byte](72,
101, 108, 108, 111)))
+ assert(java.util.Arrays.equals(rawRow.getAs[Array[Byte]](1),
Array[Byte](89, 69)))
+ assert(
+ java.util.Arrays.equals(rawRow.getAs[Array[Byte]](2), Array[Byte](72,
101, 108, 108, 111)))
+
+ val unchanged = sql("SELECT pictures['keep'] FROM t WHERE id =
2").collect()(0)
+ assert(java.util.Arrays.equals(unchanged.getAs[Array[Byte]](0),
Array[Byte](65, 66, 67)))
+ }
+ }
+
test("Blob: merge-into updates multiple raw-data BLOB columns with split
blob files") {
withTable("s", "t") {
def bytesHex(value: Int, length: Int): String = {
diff --git
a/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/DataEvolutionDeletionTestBase.scala
b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/DataEvolutionDeletionTestBase.scala
index 24a271c755..9ea8962407 100644
---
a/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/DataEvolutionDeletionTestBase.scala
+++
b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/DataEvolutionDeletionTestBase.scala
@@ -155,10 +155,18 @@ abstract class DataEvolutionDeletionTestBase extends
PaimonSparkTestBase {
| (5, 5, X'05'), (6, 6, X'06'), (7, 7, X'07'), (8, 8, X'08'), (9,
9, X'09')
| AS v(id, b, picture)
|""".stripMargin)
- sql("DELETE FROM t WHERE id IN (0, 1, 2, 3, 4, 6, 9)")
+ sql("""
+ |INSERT INTO t SELECT /*+ REPARTITION(1) */ id, b, picture FROM
VALUES
+ | (10, 10, X'0A'), (11, 11, X'0B'), (12, 12, X'0C'),
+ | (13, 13, X'0D'), (14, 14, X'0E')
+ | AS v(id, b, picture)
+ |""".stripMargin)
+ sql("DELETE FROM t WHERE id IN (0, 1, 2, 3, 4, 6, 9, 14)")
sql("CREATE TABLE s (id INT, picture BINARY)")
- sql("INSERT INTO s VALUES (2, X'22'), (6, X'66'), (7, X'4D'), (9,
X'79')")
+ sql(
+ "INSERT INTO s VALUES " +
+ "(2, X'22'), (6, X'66'), (7, X'4D'), (9, X'79'), (12, X'7A')")
sql("""
|MERGE INTO t
@@ -172,7 +180,13 @@ abstract class DataEvolutionDeletionTestBase extends
PaimonSparkTestBase {
Seq(
Row(5, 5, Array[Byte](5), 5L),
Row(7, 7, Array[Byte](77), 7L),
- Row(8, 8, Array[Byte](8), 8L)))
+ Row(8, 8, Array[Byte](8), 8L),
+ Row(10, 10, Array[Byte](10), 10L),
+ Row(11, 11, Array[Byte](11), 11L),
+ Row(12, 12, Array[Byte](122), 12L),
+ Row(13, 13, Array[Byte](13), 13L)
+ )
+ )
}
}