From bb7e7978a576a6b0f4becb65597cf2d74a5f4786 Mon Sep 17 00:00:00 2001 From: jackylee-ch Date: Tue, 11 Aug 2026 14:59:10 +0800 Subject: [PATCH] [spark] Respect blob-compaction.enabled in Spark compact procedure --- .../spark/procedure/CompactProcedure.java | 6 ++++- .../paimon/spark/sql/BlobTestBase.scala | 25 +++++++++++++++++++ 2 files changed, 30 insertions(+), 1 deletion(-) diff --git a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/CompactProcedure.java b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/CompactProcedure.java index b73c0b189a41..8fd213da1b97 100644 --- a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/CompactProcedure.java +++ b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/CompactProcedure.java @@ -526,7 +526,11 @@ private void compactDataEvolutionTable( } DataEvolutionCompactCoordinator compactCoordinator = new DataEvolutionCompactCoordinator( - table, partitionPredicate, false, false, snapshot); + table, + partitionPredicate, + table.coreOptions().blobCompactionEnabled(), + false, + snapshot); CommitMessageSerializer messageSerializerser = new CommitMessageSerializer(); String commitUser = createCommitUser(table.coreOptions().toConfiguration()); try { diff --git a/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/BlobTestBase.scala b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/BlobTestBase.scala index 72d80690b38b..b09a243587bd 100644 --- a/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/BlobTestBase.scala +++ b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/BlobTestBase.scala @@ -704,6 +704,31 @@ class BlobTestBase extends PaimonSparkTestBase { } } + test("Blob: test compaction with blob compaction enabled") { + withTable("t") { + sql( + "CREATE TABLE t (id INT, data STRING, picture BINARY) TBLPROPERTIES ('row-tracking.enabled'='true', 'data-evolution.enabled'='true', 'blob-field'='picture', 'blob-compaction.enabled'='true')") + for (i <- 1 to 10) { + sql("INSERT INTO t VALUES (" + i + ", 'paimon', X'48656C6C6F')") + } + sql("INSERT INTO t VALUES (1, 'paimon', X'48656C6C6F')") + + checkAnswer( + sql("SELECT COUNT(*) FROM `t$files`"), + Seq(Row(22)) + ) + sql("CALL paimon.sys.compact('t')").collect() + checkAnswer( + sql("SELECT COUNT(*) FROM `t$files`"), + Seq(Row(2)) + ) + checkAnswer( + sql("SELECT *, _ROW_ID, _SEQUENCE_NUMBER FROM t LIMIT 1"), + Seq(Row(1, "paimon", Array[Byte](72, 101, 108, 108, 111), 0, 11)) + ) + } + } + test("Blob: merge-into updates non-blob column on raw blob table with split blob files") { withTable("s", "t") { sql(