From 5f8f024ca32a90e69fd9d3ce1efab5b8e7491338 Mon Sep 17 00:00:00 2001 From: weimingdiit Date: Mon, 27 Jul 2026 18:52:35 +0800 Subject: [PATCH 1/3] [AURON #2434] Support full-data-file Iceberg changelog deletes Allow native Iceberg changelog scan to handle DeletedDataFileScanTask when the task represents a DELETE operation and has no existing delete files. Keep existing support for insert changelog tasks and continue falling back for row-level deletes, position/equality deletes, update changelog tasks, mixed file formats, and delete tasks with existing deletes. Add tests for full-data-file delete changelog scans and mixed insert/delete changelog ranges. Signed-off-by: weimingdiit --- .../auron/iceberg/IcebergScanSupport.scala | 62 +++++++++------ .../AuronIcebergIntegrationSuite.scala | 76 +++++++++++++++---- 2 files changed, 102 insertions(+), 36 deletions(-) diff --git a/thirdparty/auron-iceberg/src/main/scala/org/apache/spark/sql/auron/iceberg/IcebergScanSupport.scala b/thirdparty/auron-iceberg/src/main/scala/org/apache/spark/sql/auron/iceberg/IcebergScanSupport.scala index e087e1617..aaf91c9b9 100644 --- a/thirdparty/auron-iceberg/src/main/scala/org/apache/spark/sql/auron/iceberg/IcebergScanSupport.scala +++ b/thirdparty/auron-iceberg/src/main/scala/org/apache/spark/sql/auron/iceberg/IcebergScanSupport.scala @@ -22,7 +22,7 @@ import scala.collection.JavaConverters._ import scala.util.control.NonFatal import org.apache.commons.lang3.reflect.MethodUtils -import org.apache.iceberg.{AddedRowsScanTask, ChangelogOperation, ChangelogScanTask, FileFormat, FileScanTask, MetadataColumns, ScanTask} +import org.apache.iceberg.{AddedRowsScanTask, ChangelogOperation, ChangelogScanTask, DataFile, DeletedDataFileScanTask, FileFormat, FileScanTask, MetadataColumns, ScanTask} import org.apache.iceberg.expressions.{And => IcebergAnd, BoundPredicate, Expression => IcebergExpression, Not => IcebergNot, Or => IcebergOr, UnboundPredicate} import org.apache.iceberg.spark.source.AuronIcebergSourceUtil import org.apache.spark.internal.Logging @@ -267,22 +267,14 @@ object IcebergScanSupport extends Logging { return None } - val addedRowsTasks = changelogTasks.collect { case task: AddedRowsScanTask => task } - // First native changelog support is insert-only. Delete and update images need Iceberg - // delete-file handling, so keep them on Spark's reader for now. - if (addedRowsTasks.size != changelogTasks.size) { + // Native changelog scan can read data-file rows directly for insert tasks and + // full-data-file delete tasks. Row-level deletes still need Iceberg delete-file handling. + val nativeChangelogTasks = changelogTasks.flatMap(toNativeChangelogDataFileTask) + if (nativeChangelogTasks.size != changelogTasks.size) { return None } - if (!addedRowsTasks.forall(_.operation() == ChangelogOperation.INSERT)) { - return None - } - - if (!addedRowsTasks.forall(task => deletesEmpty(task.deletes()))) { - return None - } - - val formats = addedRowsTasks.map(_.file().format()).distinct + val formats = nativeChangelogTasks.map(_.file.format()).distinct if (formats.size > 1) { return None } @@ -298,7 +290,7 @@ object IcebergScanSupport extends Logging { } val pruningPredicates = collectPruningPredicates(scan.asInstanceOf[AnyRef], readSchema) - val nativeTasks = addedRowsTasks.map(task => toNativeScanTask(task, partitionSchema)) + val nativeTasks = nativeChangelogTasks.map(task => toNativeScanTask(task, partitionSchema)) Some( IcebergScanPlan( nativeTasks, @@ -528,6 +520,12 @@ object IcebergScanSupport extends Logging { private case class IcebergPartitionView(tasks: Seq[ScanTask]) + private case class NativeChangelogDataFileTask( + file: DataFile, + start: Long, + length: Long, + changelogTask: ChangelogScanTask) + private def icebergPartition(partition: InputPartition): Option[IcebergPartitionView] = { val className = partition.getClass.getName // Only accept Iceberg SparkInputPartition to access task groups. @@ -559,6 +557,23 @@ object IcebergScanSupport extends Logging { } } + private def toNativeChangelogDataFileTask( + task: ChangelogScanTask): Option[NativeChangelogDataFileTask] = { + task match { + case added: AddedRowsScanTask + if added.operation() == ChangelogOperation.INSERT && + deletesEmpty(added.deletes()) => + Some(NativeChangelogDataFileTask(added.file(), added.start(), added.length(), added)) + case deleted: DeletedDataFileScanTask + if deleted.operation() == ChangelogOperation.DELETE && + deletesEmpty(deleted.existingDeletes()) => + Some( + NativeChangelogDataFileTask(deleted.file(), deleted.start(), deleted.length(), deleted)) + case _ => + None + } + } + private def toNativeScanTask( task: FileScanTask, partitionSchema: StructType): IcebergNativeScanTask = { @@ -572,15 +587,18 @@ object IcebergScanSupport extends Logging { } private def toNativeScanTask( - task: AddedRowsScanTask, + task: NativeChangelogDataFileTask, partitionSchema: StructType): IcebergNativeScanTask = { - val file = task.file() IcebergNativeScanTask( - file.location(), - task.start(), - task.length(), - file.fileSizeInBytes(), - metadataPartitionValues(file.location(), file.specId(), Some(task), partitionSchema)) + task.file.location(), + task.start, + task.length, + task.file.fileSizeInBytes(), + metadataPartitionValues( + task.file.location(), + task.file.specId(), + Some(task.changelogTask), + partitionSchema)) } private def metadataPartitionValues( diff --git a/thirdparty/auron-iceberg/src/test/scala/org/apache/auron/iceberg/AuronIcebergIntegrationSuite.scala b/thirdparty/auron-iceberg/src/test/scala/org/apache/auron/iceberg/AuronIcebergIntegrationSuite.scala index 0032e289e..cecc01698 100644 --- a/thirdparty/auron-iceberg/src/test/scala/org/apache/auron/iceberg/AuronIcebergIntegrationSuite.scala +++ b/thirdparty/auron-iceberg/src/test/scala/org/apache/auron/iceberg/AuronIcebergIntegrationSuite.scala @@ -597,6 +597,48 @@ class AuronIcebergIntegrationSuite } } + test("iceberg native scan supports full-data-file delete changelog scan") { + withTable("local.db.t_changelog_full_file_delete") { + withTempView("t_changelog_full_file_delete_changes") { + sql(""" + |create table local.db.t_changelog_full_file_delete (id int, v string, p int) + |using iceberg + |partitioned by (p) + |tblproperties ('format-version' = '2') + |""".stripMargin) + sql(""" + |insert into local.db.t_changelog_full_file_delete + |values (1, 'a', 1), (2, 'b', 2) + |""".stripMargin) + val startSnapshotId = currentSnapshotId("local.db.t_changelog_full_file_delete") + sql("delete from local.db.t_changelog_full_file_delete where p = 1") + val endSnapshotId = currentSnapshotId("local.db.t_changelog_full_file_delete") + createChangelogView( + "local.db.t_changelog_full_file_delete", + "t_changelog_full_file_delete_changes", + startSnapshotId, + endSnapshotId) + + val query = + """ + |select id, v, p, _change_type + |from t_changelog_full_file_delete_changes + |order by id + |""".stripMargin + val expected = Seq(Row(1, "a", 1, "DELETE")) + withSQLConf("spark.auron.enable" -> "false") { + checkAnswer(sql(query), expected) + } + withSQLConf("spark.auron.enable" -> "true", "spark.auron.enable.iceberg.scan" -> "true") { + val df = sql(query) + checkAnswer(df, expected) + val nativeScan = executedNativeIcebergTableScanExec(df) + assert(nativeScan.staticPlan.scanTasks.size == 1) + } + } + } + } + test("iceberg native changelog scan remains correct in dynamic pruning join") { withTable("local.db.t_changelog_dpp", "local.db.t_changelog_dpp_dim") { withTempView("t_changelog_dpp_changes") { @@ -700,39 +742,45 @@ class AuronIcebergIntegrationSuite } } - test("iceberg changelog scan falls back when delete changes exist") { - withTable("local.db.t_changelog_delete") { - withTempView("t_changelog_delete_changes") { + test("iceberg native scan supports mixed insert and full-data-file delete changelog scan") { + withTable("local.db.t_changelog_mixed_delete") { + withTempView("t_changelog_mixed_delete_changes") { sql(""" - |create table local.db.t_changelog_delete (id int, v string) + |create table local.db.t_changelog_mixed_delete (id int, v string, p int) |using iceberg + |partitioned by (p) |tblproperties ('format-version' = '2') |""".stripMargin) - sql("insert into local.db.t_changelog_delete values (1, 'a'), (2, 'b')") - val startSnapshotId = currentSnapshotId("local.db.t_changelog_delete") - sql("delete from local.db.t_changelog_delete where id = 1") - val endSnapshotId = currentSnapshotId("local.db.t_changelog_delete") + sql(""" + |insert into local.db.t_changelog_mixed_delete + |values (1, 'a', 1), (2, 'b', 2) + |""".stripMargin) + val startSnapshotId = currentSnapshotId("local.db.t_changelog_mixed_delete") + sql("delete from local.db.t_changelog_mixed_delete where p = 1") + sql("insert into local.db.t_changelog_mixed_delete values (3, 'c', 3)") + val endSnapshotId = currentSnapshotId("local.db.t_changelog_mixed_delete") createChangelogView( - "local.db.t_changelog_delete", - "t_changelog_delete_changes", + "local.db.t_changelog_mixed_delete", + "t_changelog_mixed_delete_changes", startSnapshotId, endSnapshotId) val query = """ - |select id, v, _change_type, _change_ordinal, _commit_snapshot_id - |from t_changelog_delete_changes + |select id, v, p, _change_type, _change_ordinal, _commit_snapshot_id + |from t_changelog_mixed_delete_changes |order by id, _change_type |""".stripMargin var expected: Seq[Row] = Nil withSQLConf("spark.auron.enable" -> "false") { expected = sql(query).collect().toSeq } + assert(expected.exists(row => row.getString(3) == "DELETE")) + assert(expected.exists(row => row.getString(3) == "INSERT")) withSQLConf("spark.auron.enable" -> "true", "spark.auron.enable.iceberg.scan" -> "true") { val df = sql(query) checkAnswer(df, expected) - val plan = df.queryExecution.executedPlan.toString() - assert(!plan.contains("NativeIcebergTableScan")) + executedNativeIcebergTableScanExec(df) } } } From c0fea0172e112dc552ce8e32d20f5f60582eaadb Mon Sep 17 00:00:00 2001 From: weimingdiit Date: Thu, 30 Jul 2026 19:52:56 +0800 Subject: [PATCH 2/3] update code Signed-off-by: weimingdiit --- .../apache/spark/sql/auron/iceberg/IcebergScanSupport.scala | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/thirdparty/auron-iceberg/src/main/scala/org/apache/spark/sql/auron/iceberg/IcebergScanSupport.scala b/thirdparty/auron-iceberg/src/main/scala/org/apache/spark/sql/auron/iceberg/IcebergScanSupport.scala index aaf91c9b9..6abbe8e8c 100644 --- a/thirdparty/auron-iceberg/src/main/scala/org/apache/spark/sql/auron/iceberg/IcebergScanSupport.scala +++ b/thirdparty/auron-iceberg/src/main/scala/org/apache/spark/sql/auron/iceberg/IcebergScanSupport.scala @@ -565,8 +565,7 @@ object IcebergScanSupport extends Logging { deletesEmpty(added.deletes()) => Some(NativeChangelogDataFileTask(added.file(), added.start(), added.length(), added)) case deleted: DeletedDataFileScanTask - if deleted.operation() == ChangelogOperation.DELETE && - deletesEmpty(deleted.existingDeletes()) => + if deletesEmpty(deleted.existingDeletes()) => Some( NativeChangelogDataFileTask(deleted.file(), deleted.start(), deleted.length(), deleted)) case _ => From f3db3dfabe51e0515a04c35feac95b53dbe32244 Mon Sep 17 00:00:00 2001 From: weimingdiit Date: Thu, 30 Jul 2026 20:22:49 +0800 Subject: [PATCH 3/3] fix checkstyle Signed-off-by: weimingdiit --- .../apache/spark/sql/auron/iceberg/IcebergScanSupport.scala | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/thirdparty/auron-iceberg/src/main/scala/org/apache/spark/sql/auron/iceberg/IcebergScanSupport.scala b/thirdparty/auron-iceberg/src/main/scala/org/apache/spark/sql/auron/iceberg/IcebergScanSupport.scala index 6abbe8e8c..9521d5a7f 100644 --- a/thirdparty/auron-iceberg/src/main/scala/org/apache/spark/sql/auron/iceberg/IcebergScanSupport.scala +++ b/thirdparty/auron-iceberg/src/main/scala/org/apache/spark/sql/auron/iceberg/IcebergScanSupport.scala @@ -564,8 +564,7 @@ object IcebergScanSupport extends Logging { if added.operation() == ChangelogOperation.INSERT && deletesEmpty(added.deletes()) => Some(NativeChangelogDataFileTask(added.file(), added.start(), added.length(), added)) - case deleted: DeletedDataFileScanTask - if deletesEmpty(deleted.existingDeletes()) => + case deleted: DeletedDataFileScanTask if deletesEmpty(deleted.existingDeletes()) => Some( NativeChangelogDataFileTask(deleted.file(), deleted.start(), deleted.length(), deleted)) case _ =>