diff --git a/iceberg/common/src/main/java/com/nvidia/spark/rapids/iceberg/IcebergShimUtils.java b/iceberg/common/src/main/java/com/nvidia/spark/rapids/iceberg/IcebergShimUtils.java index 1979970fa7b..4ce46b6bf60 100644 --- a/iceberg/common/src/main/java/com/nvidia/spark/rapids/iceberg/IcebergShimUtils.java +++ b/iceberg/common/src/main/java/com/nvidia/spark/rapids/iceberg/IcebergShimUtils.java @@ -24,18 +24,28 @@ import org.apache.hadoop.fs.Path; import org.apache.iceberg.ContentFile; import org.apache.iceberg.DeleteFile; +import org.apache.iceberg.FileFormat; import org.apache.iceberg.FileScanTask; import org.apache.iceberg.Schema; import org.apache.iceberg.Table; +import org.apache.iceberg.deletes.PositionDelete; +import org.apache.iceberg.io.DataWriteResult; +import org.apache.iceberg.io.DeleteWriteResult; import org.apache.iceberg.io.FileIO; +import org.apache.iceberg.io.OutputFileFactory; +import org.apache.iceberg.io.PartitioningWriter; +import org.apache.iceberg.io.WriteResult; import org.apache.iceberg.parquet.GpuParquetIO; import org.apache.iceberg.shaded.org.apache.parquet.ParquetReadOptions; import org.apache.iceberg.shaded.org.apache.parquet.hadoop.ParquetFileReader; import org.apache.iceberg.spark.source.GpuSparkScan; +import org.apache.spark.sql.catalyst.InternalRow; import org.apache.spark.sql.connector.read.Scan; +import org.apache.spark.sql.connector.write.DeltaBatchWrite; import scala.Option; import java.io.IOException; +import java.io.Serializable; import java.util.Map; /** @@ -61,6 +71,45 @@ public interface IcebergShimUtils { /** Returns whether a positional delete is an Iceberg Puffin deletion vector. */ boolean isDeletionVector(DeleteFile deleteFile); + /** Returns whether a resolved delete-file format writes Puffin deletion vectors. */ + boolean isPuffinFormat(FileFormat fileFormat); + + /** + * Returns delete files that Iceberg requires a deletion-vector write to replace. + * + * @return an opaque handle, or {@code null} when there are no existing deletes to rewrite + */ + RewritableDeletes broadcastRewritableDeletes( + DeltaBatchWrite write); + + /** + * Opaque, serializable handle for version-specific rewritable-delete state. + * + *

The underlying broadcast value uses Iceberg's {@code DeleteFileSet}, which is absent + * from Iceberg 1.6. Iceberg 1.9 and later share an implementation outside the common module + * so this interface does not expose an API that is unavailable in older Iceberg versions. + */ + interface RewritableDeletes extends Serializable {} + + /** + * Creates Iceberg's version-specific deletion-vector writer. + * + * @param rewritableDeletes existing deletes to merge and replace, or {@code null} when the + * data file has no existing deletes + */ + PartitioningWriter, DeleteWriteResult> + newDeletionVectorWriter( + Table table, OutputFileFactory fileFactory, + RewritableDeletes rewritableDeletes); + + /** Combines data and delete results, including rewritten deletes when supported. */ + WriteResult positionDeltaWriteResult( + DataWriteResult dataResult, DeleteWriteResult deleteResult); + + /** Populates a reusable position-delete record across Iceberg API versions. */ + void setPositionDelete( + PositionDelete delete, CharSequence path, long position); + /** * Reads exactly the recorded deletion-vector byte range and returns its compressed bitmap. * diff --git a/iceberg/common/src/main/java/com/nvidia/spark/rapids/iceberg/ShimUtils.java b/iceberg/common/src/main/java/com/nvidia/spark/rapids/iceberg/ShimUtils.java index adee3b490b3..70ec4300242 100644 --- a/iceberg/common/src/main/java/com/nvidia/spark/rapids/iceberg/ShimUtils.java +++ b/iceberg/common/src/main/java/com/nvidia/spark/rapids/iceberg/ShimUtils.java @@ -25,14 +25,23 @@ import org.apache.hadoop.fs.Path; import org.apache.iceberg.ContentFile; import org.apache.iceberg.DeleteFile; +import org.apache.iceberg.FileFormat; import org.apache.iceberg.FileScanTask; import org.apache.iceberg.Schema; import org.apache.iceberg.Table; +import org.apache.iceberg.deletes.PositionDelete; +import org.apache.iceberg.io.DataWriteResult; +import org.apache.iceberg.io.DeleteWriteResult; import org.apache.iceberg.io.FileIO; +import org.apache.iceberg.io.OutputFileFactory; +import org.apache.iceberg.io.PartitioningWriter; +import org.apache.iceberg.io.WriteResult; import org.apache.iceberg.shaded.org.apache.parquet.ParquetReadOptions; import org.apache.iceberg.shaded.org.apache.parquet.hadoop.ParquetFileReader; import org.apache.iceberg.spark.source.GpuSparkScan; +import org.apache.spark.sql.catalyst.InternalRow; import org.apache.spark.sql.connector.read.Scan; +import org.apache.spark.sql.connector.write.DeltaBatchWrite; import java.io.IOException; import java.util.Map; @@ -81,6 +90,32 @@ public static boolean isDeletionVector(DeleteFile deleteFile) { return IMPL.isDeletionVector(deleteFile); } + public static boolean isPuffinFormat(FileFormat fileFormat) { + return IMPL.isPuffinFormat(fileFormat); + } + + public static IcebergShimUtils.RewritableDeletes broadcastRewritableDeletes( + DeltaBatchWrite write) { + return IMPL.broadcastRewritableDeletes(write); + } + + public static PartitioningWriter, DeleteWriteResult> + newDeletionVectorWriter( + Table table, OutputFileFactory fileFactory, + IcebergShimUtils.RewritableDeletes rewritableDeletes) { + return IMPL.newDeletionVectorWriter(table, fileFactory, rewritableDeletes); + } + + public static WriteResult positionDeltaWriteResult( + DataWriteResult dataResult, DeleteWriteResult deleteResult) { + return IMPL.positionDeltaWriteResult(dataResult, deleteResult); + } + + public static void setPositionDelete( + PositionDelete delete, CharSequence path, long position) { + IMPL.setPositionDelete(delete, path, position); + } + public static IcebergDeletionVector readDeletionVector(DeleteFile deleteFile, RapidsInputFile inputFile, boolean validateCrc) diff --git a/iceberg/common/src/main/scala/org/apache/iceberg/spark/source/GpuSparkFileWriterFactory.scala b/iceberg/common/src/main/scala/org/apache/iceberg/spark/source/GpuSparkFileWriterFactory.scala index 1a1422f00eb..091a5f7045f 100644 --- a/iceberg/common/src/main/scala/org/apache/iceberg/spark/source/GpuSparkFileWriterFactory.scala +++ b/iceberg/common/src/main/scala/org/apache/iceberg/spark/source/GpuSparkFileWriterFactory.scala @@ -18,6 +18,7 @@ package org.apache.iceberg.spark.source import com.nvidia.spark.rapids.{ColumnarOutputWriterFactory, GpuParquetWriter, SpillableColumnarBatch} import com.nvidia.spark.rapids.fileio.iceberg.IcebergFileIO +import com.nvidia.spark.rapids.iceberg.ShimUtils import com.nvidia.spark.rapids.iceberg.parquet.GpuIcebergParquetAppender import org.apache.hadoop.conf.Configuration import org.apache.hadoop.mapreduce.TaskAttemptContext @@ -45,8 +46,10 @@ class GpuSparkFileWriterFactory(val table: Table, ) extends FileWriterFactory[SpillableColumnarBatch] { require(dataFileFormat == FileFormat.PARQUET, s"GpuSparkFileWriterFactory only supports PARQUET file format, but got $dataFileFormat") - require(deleteFileFormat == FileFormat.PARQUET, - s"GpuSparkFileWriterFactory only supports PARQUET file format, but got $deleteFileFormat") + require(deleteFileFormat == FileFormat.PARQUET || + ShimUtils.isPuffinFormat(deleteFileFormat), + s"GpuSparkFileWriterFactory only supports PARQUET or Puffin deletion vectors, " + + s"but got $deleteFileFormat") private def newTaskAttemptContext(sparkType: StructType): TaskAttemptContext = { val conf = new Configuration(hadoopConf.value) @@ -111,4 +114,4 @@ class GpuSparkFileWriterFactory(val table: Table, fileIO = new IcebergFileIO(table.io()) ) } -} \ No newline at end of file +} diff --git a/iceberg/common/src/main/scala/org/apache/iceberg/spark/source/GpuSparkPositionDeltaWrite.scala b/iceberg/common/src/main/scala/org/apache/iceberg/spark/source/GpuSparkPositionDeltaWrite.scala index fd1b24a880a..ba5d8ea8319 100644 --- a/iceberg/common/src/main/scala/org/apache/iceberg/spark/source/GpuSparkPositionDeltaWrite.scala +++ b/iceberg/common/src/main/scala/org/apache/iceberg/spark/source/GpuSparkPositionDeltaWrite.scala @@ -29,11 +29,11 @@ import com.nvidia.spark.rapids.RmmRapidsRetryIterator.withRetryNoSplit import com.nvidia.spark.rapids.SpillPriorities.ACTIVE_ON_DECK_PRIORITY import com.nvidia.spark.rapids.fileio.iceberg.IcebergFileIO import com.nvidia.spark.rapids.iceberg.{ColumnarBatchWithPartition, GpuIcebergPartitioner, - GpuIcebergSpecPartitioner, IcebergFormatVersionSupport} + GpuIcebergSpecPartitioner, IcebergFormatVersionSupport, IcebergShimUtils, ShimUtils} import com.nvidia.spark.rapids.iceberg.utils.GpuStructProjection import org.apache.hadoop.mapreduce.Job import org.apache.iceberg._ -import org.apache.iceberg.deletes.DeleteGranularity +import org.apache.iceberg.deletes.{DeleteGranularity, PositionDelete} import org.apache.iceberg.io._ import org.apache.iceberg.io.DeleteSchemaUtil import org.apache.iceberg.spark.GpuTypeToSparkType @@ -87,11 +87,17 @@ class GpuSparkPositionDeltaWrite(cpu: DeltaWrite) override def advisoryPartitionSizeInBytes(): Long = writeRequirements.advisoryPartitionSizeInBytes() - private[source] def createDeltaWriterFactory: DeltaWriterFactory = { + private[source] def createDeltaWriterFactory( + cpuBatchWrite: DeltaBatchWrite): DeltaWriterFactory = { val sparkContext: JavaSparkContext = GpuSparkWriteAccess.deltaSparkContext(cpu) val tableBroadcast = sparkContext.broadcast(SerializableTable.copyOf(table)) val command = GpuSparkWriteAccess.deltaCommand(cpu) val context = GpuWriteContext(GpuSparkWriteAccess.deltaContext(cpu)) + val rewritableDeletes = if (context.useDVs) { + Option(ShimUtils.broadcastRewritableDeletes(cpuBatchWrite)) + } else { + None + } val writeProps = GpuSparkWriteAccess.deltaWriteProperties(cpu) .asScala .toMap @@ -101,9 +107,10 @@ class GpuSparkPositionDeltaWrite(cpu: DeltaWrite) s"GpuSparkWrite only supports Parquet, but data format got: ${context.dataFileFormat}") } - if (!context.deleteFileFormat.equals(FileFormat.PARQUET)) { + if (!context.deleteFileFormat.equals(FileFormat.PARQUET) && !context.useDVs) { throw new UnsupportedOperationException( - s"GpuSparkWrite only supports Parquet, but delete format got: ${context.deleteFileFormat}") + s"GpuSparkWrite only supports Parquet or Puffin deletion vectors, " + + s"but delete format got: ${context.deleteFileFormat}") } val hadoopConf = sparkContext.hadoopConfiguration @@ -131,6 +138,7 @@ class GpuSparkPositionDeltaWrite(cpu: DeltaWrite) new GpuPositionDeltaWriterFactory( tableBroadcast, + rewritableDeletes, command, context, writeProps, @@ -191,7 +199,7 @@ object GpuSparkPositionDeltaWrite { // delete codec matters too. GpuSparkWrite.tagParquetCompressionForGpu( GpuSparkWriteAccess.deltaWriteProperties(deltaWrite), - hasDeleteFiles = true, meta) + hasDeleteFiles = !context.useDVs, meta) } def convert(deltaWrite: DeltaWrite): GpuSparkPositionDeltaWrite = { @@ -201,6 +209,7 @@ object GpuSparkPositionDeltaWrite { class GpuPositionDeltaWriterFactory( val tableSer: Broadcast[Table], + val rewritableDeletesSer: Option[IcebergShimUtils.RewritableDeletes], val command: Command, val context: GpuWriteContext, val writeProps: Map[String, String], @@ -210,7 +219,6 @@ class GpuPositionDeltaWriterFactory( override def createWriter(partitionId: Int, taskId: Long): DeltaWriter[InternalRow] = { val table = tableSer.value - val deleteFileFactory = OutputFileFactory.builderFor(table, partitionId, taskId) .format(context.deleteFileFormat) .operationId(context.queryId) @@ -235,16 +243,17 @@ class GpuPositionDeltaWriterFactory( new IcebergFileIO(table.io())) if (command == Command.DELETE) { - new GpuDeleteOnlyDeltaWriter(table, writerFactory, deleteFileFactory, context) + new GpuDeleteOnlyDeltaWriter( + table, rewritableDeletesSer, writerFactory, deleteFileFactory, context) .asInstanceOf[DeltaWriter[InternalRow]] } else { if (table.spec().isUnpartitioned) { - new GpuUnpartitionedDeltaWriter(table, writerFactory, dataFileFactory, - deleteFileFactory, context) + new GpuUnpartitionedDeltaWriter(table, rewritableDeletesSer, writerFactory, + dataFileFactory, deleteFileFactory, context) .asInstanceOf[DeltaWriter[InternalRow]] } else { - new GpuPartitionedDeltaWriter(table, writerFactory, dataFileFactory, - deleteFileFactory, context) + new GpuPartitionedDeltaWriter(table, rewritableDeletesSer, writerFactory, + dataFileFactory, deleteFileFactory, context) .asInstanceOf[DeltaWriter[InternalRow]] } } @@ -256,6 +265,9 @@ trait GpuDeltaWriter extends DeltaWriter[ColumnarBatch] { def context: GpuWriteContext + protected def rewritableDeletesSer: + Option[IcebergShimUtils.RewritableDeletes] + protected def buildPartitionProjections( partitionType: IcebergTypes.StructType, specs: collection.Map[Integer, PartitionSpec]): Map[Int, GpuStructProjection] = { @@ -272,7 +284,11 @@ trait GpuDeltaWriter extends DeltaWriter[ColumnarBatch] { val inputOrdered = context.inputOrdered val targetFileSize = context.targetDeleteFileSize - if (inputOrdered) { + if (context.useDVs) { + new GpuBatchPositionDeleteWriter( + ShimUtils.newDeletionVectorWriter( + table, outputFileFactory, rewritableDeletesSer.orNull)) + } else if (inputOrdered) { new GpuClusteredPositionDeleteWriter(writerFactory, outputFileFactory, io, targetFileSize) } else { new GpuFanoutPositionDeleteWriter(writerFactory, outputFileFactory, io, targetFileSize) @@ -458,14 +474,50 @@ class GpuBasePositionDeltaWriter( def result(): WriteResult = { val dataResult = dataWriter.result() val deleteResult = deleteWriter.result() - WriteResult.builder() - .addDataFiles(dataResult.dataFiles()) - .addDeleteFiles(deleteResult.deleteFiles()) - .addReferencedDataFiles(deleteResult.referencedDataFiles()) - .build() + ShimUtils.positionDeltaWriteResult(dataResult, deleteResult) } } +/** Converts GPU-produced batches of file paths and positions for Iceberg's DV encoder. */ +class GpuBatchPositionDeleteWriter( + private val delegate: PartitioningWriter[PositionDelete[InternalRow], DeleteWriteResult]) + extends PartitioningWriter[SpillableColumnarBatch, DeleteWriteResult] { + + require(delegate != null, "delegate must not be null") + + private val positionDelete = PositionDelete.create[InternalRow]() + + override def write( + spillableBatch: SpillableColumnarBatch, + spec: PartitionSpec, + partition: StructLike): Unit = { + withResource(spillableBatch) { spillable => + val (paths, positions, numRows) = withResource(spillable.getColumnarBatch()) { batch => + closeOnExcept(batch.column(0).asInstanceOf[GpuColumnVector].copyToHost()) { paths => + closeOnExcept(batch.column(1).asInstanceOf[GpuColumnVector].copyToHost()) { positions => + (paths, positions, batch.numRows()) + } + } + } + withResource(paths) { pathColumn => + withResource(positions) { positionColumn => + for (row <- 0 until numRows) { + ShimUtils.setPositionDelete( + positionDelete, + pathColumn.getUTF8String(row).toString, + positionColumn.getLong(row)) + delegate.write(positionDelete, spec, partition) + } + } + } + } + } + + override def result(): DeleteWriteResult = delegate.result() + + override def close(): Unit = delegate.close() +} + /** * Base trait for delta writers that handle both deletes and data writes. * This is the GPU equivalent of Java's DeleteAndDataDeltaWriter. @@ -567,6 +619,7 @@ trait GpuDeleteAndDataDeltaWriter extends GpuDeltaWriter { */ class GpuDeleteOnlyDeltaWriter( table: Table, + override protected val rewritableDeletesSer: Option[IcebergShimUtils.RewritableDeletes], writerFactory: GpuSparkFileWriterFactory, deleteFileFactory: OutputFileFactory, override val context: GpuWriteContext) extends GpuDeltaWriter { @@ -668,6 +721,7 @@ class GpuDeleteOnlyDeltaWriter( */ class GpuUnpartitionedDeltaWriter( protected val table: Table, + override protected val rewritableDeletesSer: Option[IcebergShimUtils.RewritableDeletes], writerFactory: GpuSparkFileWriterFactory, dataFileFactory: OutputFileFactory, deleteFileFactory: OutputFileFactory, @@ -718,6 +772,7 @@ class GpuUnpartitionedDeltaWriter( */ class GpuPartitionedDeltaWriter( protected val table: Table, + override protected val rewritableDeletesSer: Option[IcebergShimUtils.RewritableDeletes], writerFactory: GpuSparkFileWriterFactory, dataFileFactory: OutputFileFactory, deleteFileFactory: OutputFileFactory, @@ -778,7 +833,8 @@ case class GpuWriteContext( deleteGranularity: DeleteGranularity, queryId: String, useFanoutWriter: Boolean, - inputOrdered: Boolean) { + inputOrdered: Boolean, + useDVs: Boolean) { /** * Returns the ordinal of the spec ID column in the delete schema. @@ -840,6 +896,7 @@ object GpuWriteContext { val queryId = GpuSparkWriteAccess.contextQueryId(cpu) val useFanoutWriter = GpuSparkWriteAccess.contextUseFanoutWriter(cpu) val inputOrdered = GpuSparkWriteAccess.contextInputOrdered(cpu) + val useDVs = ShimUtils.isPuffinFormat(deleteFileFormat) GpuWriteContext( dataSchema, @@ -853,6 +910,7 @@ object GpuWriteContext { deleteGranularity, queryId, useFanoutWriter, - inputOrdered) + inputOrdered, + useDVs) } } diff --git a/iceberg/common/src/main/scala/org/apache/iceberg/spark/source/GpuSparkWrite.scala b/iceberg/common/src/main/scala/org/apache/iceberg/spark/source/GpuSparkWrite.scala index f2659d5ddea..59942bf06a4 100644 --- a/iceberg/common/src/main/scala/org/apache/iceberg/spark/source/GpuSparkWrite.scala +++ b/iceberg/common/src/main/scala/org/apache/iceberg/spark/source/GpuSparkWrite.scala @@ -26,7 +26,8 @@ import com.nvidia.spark.rapids.Arm.closeOnExcept import com.nvidia.spark.rapids.RapidsPluginImplicits.AutoCloseableSeq import com.nvidia.spark.rapids.SpillPriorities.ACTIVE_ON_DECK_PRIORITY import com.nvidia.spark.rapids.fileio.iceberg.IcebergFileIO -import com.nvidia.spark.rapids.iceberg.{GpuIcebergSpecPartitioner, IcebergFormatVersionSupport} +import com.nvidia.spark.rapids.iceberg.{GpuIcebergSpecPartitioner, IcebergFormatVersionSupport, + ShimUtils} import com.nvidia.spark.rapids.shims.parquet.ParquetFieldIdShims import org.apache.hadoop.mapreduce.Job import org.apache.iceberg._ @@ -240,8 +241,11 @@ object GpuSparkWrite { meta.willNotWorkOnGpu(s"GpuSparkWrite only supports Parquet, but got: ${dataFormat.get}") } - if (deleteFormat.exists(!_.equals(FileFormat.PARQUET))) { - meta.willNotWorkOnGpu(s"GpuSparkWrite only supports Parquet, but got: ${deleteFormat.get}") + if (deleteFormat.exists(format => + !format.equals(FileFormat.PARQUET) && !ShimUtils.isPuffinFormat(format))) { + meta.willNotWorkOnGpu( + s"GpuSparkWrite only supports Parquet or Puffin deletion vectors, " + + s"but got: ${deleteFormat.get}") } // Check partition transform support diff --git a/iceberg/common/src/main/scala/org/apache/iceberg/spark/source/write.scala b/iceberg/common/src/main/scala/org/apache/iceberg/spark/source/write.scala index 4861d846e73..dce9bcb805f 100644 --- a/iceberg/common/src/main/scala/org/apache/iceberg/spark/source/write.scala +++ b/iceberg/common/src/main/scala/org/apache/iceberg/spark/source/write.scala @@ -98,6 +98,6 @@ class GpuPositionDeltaBatchWrite(write: GpuSparkPositionDeltaWrite, override def useCommitCoordinator(): Boolean = cpuBatchWrite.useCommitCoordinator() override def createBatchWriterFactory(info: PhysicalWriteInfo): DeltaWriterFactory = { - write.createDeltaWriterFactory + write.createDeltaWriterFactory(cpuBatchWrite) } -} \ No newline at end of file +} diff --git a/iceberg/iceberg-1-10-x/pom.xml b/iceberg/iceberg-1-10-x/pom.xml index 0b3b2926c33..32ed835abf2 100644 --- a/iceberg/iceberg-1-10-x/pom.xml +++ b/iceberg/iceberg-1-10-x/pom.xml @@ -84,6 +84,7 @@ ${spark.rapids.source.basedir}/iceberg/common/src/main/java ${spark.rapids.source.basedir}/iceberg/common/src/main/scala + ${spark.rapids.source.basedir}/iceberg/iceberg-19plus-common/src/main/java diff --git a/iceberg/iceberg-1-10-x/src/main/java/com/nvidia/spark/rapids/iceberg/iceberg110x/ShimUtilsImpl.java b/iceberg/iceberg-1-10-x/src/main/java/com/nvidia/spark/rapids/iceberg/iceberg110x/ShimUtilsImpl.java index e95dbe353d7..6ca849f9487 100644 --- a/iceberg/iceberg-1-10-x/src/main/java/com/nvidia/spark/rapids/iceberg/iceberg110x/ShimUtilsImpl.java +++ b/iceberg/iceberg-1-10-x/src/main/java/com/nvidia/spark/rapids/iceberg/iceberg110x/ShimUtilsImpl.java @@ -20,7 +20,7 @@ import com.nvidia.spark.rapids.RapidsConf; import com.nvidia.spark.rapids.fileio.iceberg.IcebergInputFile; import com.nvidia.spark.rapids.iceberg.IcebergDeletionVector; -import com.nvidia.spark.rapids.iceberg.IcebergShimUtils; +import com.nvidia.spark.rapids.iceberg.Iceberg19PlusShimUtils; import com.nvidia.spark.rapids.jni.fileio.RapidsInputFile; import org.apache.hadoop.fs.Path; import org.apache.iceberg.*; @@ -42,7 +42,7 @@ import java.util.Map; /** Iceberg 1.10.x shim: uses {@code SparkUtil::internalToSpark} and a cache-aware footer path. */ -public class ShimUtilsImpl implements IcebergShimUtils { +public class ShimUtilsImpl extends Iceberg19PlusShimUtils { @Override public int formatVersion(Table table) { return TableUtil.formatVersion(table); @@ -53,11 +53,6 @@ public String locationOf(ContentFile f) { return f.location(); } - @Override - public boolean isDeletionVector(DeleteFile deleteFile) { - return deleteFile.format() == FileFormat.PUFFIN; - } - @Override public IcebergDeletionVector readDeletionVector( DeleteFile deleteFile, RapidsInputFile inputFile, boolean validateCrc) diff --git a/iceberg/iceberg-1-11-x/pom.xml b/iceberg/iceberg-1-11-x/pom.xml index e7e3c290f0c..012af9cc3b8 100644 --- a/iceberg/iceberg-1-11-x/pom.xml +++ b/iceberg/iceberg-1-11-x/pom.xml @@ -84,6 +84,7 @@ ${spark.rapids.source.basedir}/iceberg/common/src/main/java ${spark.rapids.source.basedir}/iceberg/common/src/main/scala + ${spark.rapids.source.basedir}/iceberg/iceberg-19plus-common/src/main/java diff --git a/iceberg/iceberg-1-11-x/src/main/java/com/nvidia/spark/rapids/iceberg/iceberg111x/ShimUtilsImpl.java b/iceberg/iceberg-1-11-x/src/main/java/com/nvidia/spark/rapids/iceberg/iceberg111x/ShimUtilsImpl.java index 7598f6ddda5..a5c50f70160 100644 --- a/iceberg/iceberg-1-11-x/src/main/java/com/nvidia/spark/rapids/iceberg/iceberg111x/ShimUtilsImpl.java +++ b/iceberg/iceberg-1-11-x/src/main/java/com/nvidia/spark/rapids/iceberg/iceberg111x/ShimUtilsImpl.java @@ -20,7 +20,7 @@ import com.nvidia.spark.rapids.RapidsConf; import com.nvidia.spark.rapids.fileio.iceberg.IcebergInputFile; import com.nvidia.spark.rapids.iceberg.IcebergDeletionVector; -import com.nvidia.spark.rapids.iceberg.IcebergShimUtils; +import com.nvidia.spark.rapids.iceberg.Iceberg19PlusShimUtils; import com.nvidia.spark.rapids.jni.fileio.RapidsInputFile; import org.apache.hadoop.fs.Path; import org.apache.iceberg.*; @@ -42,7 +42,7 @@ import java.util.Map; /** Iceberg 1.11.x shim: uses {@code SparkUtil::internalToSpark} and a cache-aware footer path. */ -public class ShimUtilsImpl implements IcebergShimUtils { +public class ShimUtilsImpl extends Iceberg19PlusShimUtils { @Override public int formatVersion(Table table) { return TableUtil.formatVersion(table); @@ -53,11 +53,6 @@ public String locationOf(ContentFile f) { return f.location(); } - @Override - public boolean isDeletionVector(DeleteFile deleteFile) { - return deleteFile.format() == FileFormat.PUFFIN; - } - @Override public IcebergDeletionVector readDeletionVector( DeleteFile deleteFile, RapidsInputFile inputFile, boolean validateCrc) diff --git a/iceberg/iceberg-1-6-x/src/main/java/com/nvidia/spark/rapids/iceberg/iceberg16x/ShimUtilsImpl.java b/iceberg/iceberg-1-6-x/src/main/java/com/nvidia/spark/rapids/iceberg/iceberg16x/ShimUtilsImpl.java index 1312b3c8ff9..69333d74322 100644 --- a/iceberg/iceberg-1-6-x/src/main/java/com/nvidia/spark/rapids/iceberg/iceberg16x/ShimUtilsImpl.java +++ b/iceberg/iceberg-1-6-x/src/main/java/com/nvidia/spark/rapids/iceberg/iceberg16x/ShimUtilsImpl.java @@ -21,14 +21,22 @@ import com.nvidia.spark.rapids.iceberg.IcebergShimUtils; import com.nvidia.spark.rapids.jni.fileio.RapidsInputFile; import org.apache.iceberg.*; +import org.apache.iceberg.deletes.PositionDelete; +import org.apache.iceberg.io.DataWriteResult; +import org.apache.iceberg.io.DeleteWriteResult; import org.apache.iceberg.io.FileIO; +import org.apache.iceberg.io.OutputFileFactory; +import org.apache.iceberg.io.PartitioningWriter; +import org.apache.iceberg.io.WriteResult; import org.apache.iceberg.relocated.com.google.common.base.Preconditions; import org.apache.iceberg.spark.source.GpuBaseReader; import org.apache.iceberg.spark.source.GpuSparkCopyOnWriteV1Scan; import org.apache.iceberg.spark.source.GpuSparkScan; import org.apache.iceberg.types.Types; import org.apache.iceberg.util.PartitionUtil; +import org.apache.spark.sql.catalyst.InternalRow; import org.apache.spark.sql.connector.read.Scan; +import org.apache.spark.sql.connector.write.DeltaBatchWrite; import java.io.IOException; import java.util.Collections; @@ -67,6 +75,44 @@ public boolean isDeletionVector(DeleteFile deleteFile) { return false; } + @Override + public boolean isPuffinFormat(FileFormat fileFormat) { + return false; + } + + @Override + public RewritableDeletes broadcastRewritableDeletes( + DeltaBatchWrite write) { + throw new UnsupportedOperationException( + "Iceberg 1.6 does not support Puffin deletion vectors"); + } + + @Override + public PartitioningWriter, DeleteWriteResult> + newDeletionVectorWriter( + Table table, OutputFileFactory fileFactory, + RewritableDeletes rewritableDeletes) { + throw new UnsupportedOperationException( + "Iceberg 1.6 does not support Puffin deletion vectors"); + } + + @Override + public WriteResult positionDeltaWriteResult( + DataWriteResult dataResult, DeleteWriteResult deleteResult) { + return WriteResult.builder() + .addDataFiles(dataResult.dataFiles()) + .addDeleteFiles(deleteResult.deleteFiles()) + .addReferencedDataFiles(deleteResult.referencedDataFiles()) + .build(); + } + + @Override + @SuppressWarnings("deprecation") + public void setPositionDelete( + PositionDelete delete, CharSequence path, long position) { + delete.set(path, position, null); + } + @Override public IcebergDeletionVector readDeletionVector( DeleteFile deleteFile, RapidsInputFile inputFile, boolean validateCrc) diff --git a/iceberg/iceberg-1-9-x/pom.xml b/iceberg/iceberg-1-9-x/pom.xml index 1064213fcc0..6361f5d72b4 100644 --- a/iceberg/iceberg-1-9-x/pom.xml +++ b/iceberg/iceberg-1-9-x/pom.xml @@ -74,6 +74,7 @@ ${spark.rapids.source.basedir}/iceberg/common/src/main/java ${spark.rapids.source.basedir}/iceberg/common/src/main/scala + ${spark.rapids.source.basedir}/iceberg/iceberg-19plus-common/src/main/java diff --git a/iceberg/iceberg-1-9-x/src/main/java/com/nvidia/spark/rapids/iceberg/iceberg19x/ShimUtilsImpl.java b/iceberg/iceberg-1-9-x/src/main/java/com/nvidia/spark/rapids/iceberg/iceberg19x/ShimUtilsImpl.java index 77ab23b5b11..ae58ba9f99a 100644 --- a/iceberg/iceberg-1-9-x/src/main/java/com/nvidia/spark/rapids/iceberg/iceberg19x/ShimUtilsImpl.java +++ b/iceberg/iceberg-1-9-x/src/main/java/com/nvidia/spark/rapids/iceberg/iceberg19x/ShimUtilsImpl.java @@ -18,10 +18,10 @@ import com.nvidia.spark.rapids.RapidsConf; import com.nvidia.spark.rapids.iceberg.IcebergDeletionVector; -import com.nvidia.spark.rapids.iceberg.IcebergShimUtils; +import com.nvidia.spark.rapids.iceberg.Iceberg19PlusShimUtils; import com.nvidia.spark.rapids.jni.fileio.RapidsInputFile; import org.apache.iceberg.*; -import org.apache.iceberg.io.FileIO; +import org.apache.iceberg.io.*; import org.apache.iceberg.io.StorageCredential; import org.apache.iceberg.io.SupportsStorageCredentials; import org.apache.iceberg.spark.SparkUtil; @@ -37,7 +37,7 @@ import java.util.Map; /** Iceberg 1.9.x shim: uses {@code SparkUtil::internalToSpark}. */ -public class ShimUtilsImpl implements IcebergShimUtils { +public class ShimUtilsImpl extends Iceberg19PlusShimUtils { @Override public int formatVersion(Table table) { return TableUtil.formatVersion(table); @@ -48,11 +48,6 @@ public String locationOf(ContentFile f) { return f.location(); } - @Override - public boolean isDeletionVector(DeleteFile deleteFile) { - return deleteFile.format() == FileFormat.PUFFIN; - } - @Override public IcebergDeletionVector readDeletionVector( DeleteFile deleteFile, RapidsInputFile inputFile, boolean validateCrc) diff --git a/iceberg/iceberg-19plus-common/src/main/java/com/nvidia/spark/rapids/iceberg/Iceberg19PlusShimUtils.java b/iceberg/iceberg-19plus-common/src/main/java/com/nvidia/spark/rapids/iceberg/Iceberg19PlusShimUtils.java new file mode 100644 index 00000000000..6dd4b6fc939 --- /dev/null +++ b/iceberg/iceberg-19plus-common/src/main/java/com/nvidia/spark/rapids/iceberg/Iceberg19PlusShimUtils.java @@ -0,0 +1,116 @@ +/* + * Copyright (c) 2026, NVIDIA CORPORATION. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package com.nvidia.spark.rapids.iceberg; + +import org.apache.iceberg.DeleteFile; +import org.apache.iceberg.FileFormat; +import org.apache.iceberg.Table; +import org.apache.iceberg.data.BaseDeleteLoader; +import org.apache.iceberg.deletes.PositionDelete; +import org.apache.iceberg.deletes.PositionDeleteIndex; +import org.apache.iceberg.encryption.EncryptingFileIO; +import org.apache.iceberg.io.DataWriteResult; +import org.apache.iceberg.io.DeleteWriteResult; +import org.apache.iceberg.io.OutputFileFactory; +import org.apache.iceberg.io.PartitioningDVWriter; +import org.apache.iceberg.io.PartitioningWriter; +import org.apache.iceberg.io.WriteResult; +import org.apache.iceberg.spark.source.GpuSparkPositionDeltaWriteAccess; +import org.apache.iceberg.util.DeleteFileSet; +import org.apache.spark.broadcast.Broadcast; +import org.apache.spark.sql.catalyst.InternalRow; +import org.apache.spark.sql.connector.write.DeltaBatchWrite; + +import java.util.Map; +import java.util.Set; +import java.util.function.Function; + +/** Shared deletion-vector shim implementation for Iceberg 1.9 and later. */ +public abstract class Iceberg19PlusShimUtils implements IcebergShimUtils { + @Override + public boolean isDeletionVector(DeleteFile deleteFile) { + return deleteFile.format() == FileFormat.PUFFIN; + } + + @Override + public boolean isPuffinFormat(FileFormat fileFormat) { + return fileFormat == FileFormat.PUFFIN; + } + + @Override + public RewritableDeletes broadcastRewritableDeletes(DeltaBatchWrite write) { + Broadcast> rewritableDeletes = + GpuSparkPositionDeltaWriteAccess.broadcastRewritableDeletes(write); + return rewritableDeletes != null ? new RewritableDeletesImpl(rewritableDeletes) : null; + } + + @Override + public PartitioningWriter, DeleteWriteResult> + newDeletionVectorWriter( + Table table, OutputFileFactory fileFactory, + RewritableDeletes rewritableDeletes) { + Map deleteFiles = rewritableDeletes == null + ? null + : ((RewritableDeletesImpl) rewritableDeletes).value(); + return new PartitioningDVWriter<>( + fileFactory, previousDeleteLoader(table, deleteFiles)); + } + + @Override + public WriteResult positionDeltaWriteResult( + DataWriteResult dataResult, DeleteWriteResult deleteResult) { + return WriteResult.builder() + .addDataFiles(dataResult.dataFiles()) + .addDeleteFiles(deleteResult.deleteFiles()) + .addReferencedDataFiles(deleteResult.referencedDataFiles()) + .addRewrittenDeleteFiles(deleteResult.rewrittenDeleteFiles()) + .build(); + } + + @Override + public void setPositionDelete( + PositionDelete delete, CharSequence path, long position) { + delete.set(path, position); + } + + private static Function previousDeleteLoader( + Table table, Map rewritableDeletes) { + if (rewritableDeletes == null) { + return path -> null; + } + + BaseDeleteLoader deleteLoader = new BaseDeleteLoader( + deleteFile -> EncryptingFileIO.combine(table.io(), table.encryption()) + .newInputFile(deleteFile)); + return path -> { + Set files = rewritableDeletes.get(path.toString()); + return files != null ? deleteLoader.loadPositionDeletes(files, path) : null; + }; + } + + private static final class RewritableDeletesImpl implements RewritableDeletes { + private final Broadcast> delegate; + + private RewritableDeletesImpl(Broadcast> delegate) { + this.delegate = delegate; + } + + private Map value() { + return delegate.value(); + } + } +} diff --git a/iceberg/iceberg-19plus-common/src/main/java/org/apache/iceberg/spark/source/GpuSparkPositionDeltaWriteAccess.java b/iceberg/iceberg-19plus-common/src/main/java/org/apache/iceberg/spark/source/GpuSparkPositionDeltaWriteAccess.java new file mode 100644 index 00000000000..3e77027119c --- /dev/null +++ b/iceberg/iceberg-19plus-common/src/main/java/org/apache/iceberg/spark/source/GpuSparkPositionDeltaWriteAccess.java @@ -0,0 +1,75 @@ +/* + * Copyright (c) 2026, NVIDIA CORPORATION. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.iceberg.spark.source; + +import java.lang.reflect.Method; +import java.util.Map; + +import org.apache.iceberg.util.DeleteFileSet; +import org.apache.spark.broadcast.Broadcast; +import org.apache.spark.sql.connector.write.DeltaBatchWrite; + +/** Access to position-delta batch-write internals shared by Iceberg 1.9 and later. */ +public final class GpuSparkPositionDeltaWriteAccess { + private static final ClassValue BROADCAST_REWRITABLE_DELETES_METHOD = + new ClassValue() { + @Override + protected Method computeValue(Class type) { + Method method = findMethod(type, "broadcastRewritableDeletes"); + method.setAccessible(true); + return method; + } + }; + + private GpuSparkPositionDeltaWriteAccess() { + } + + /** + * Returns the delete files that Iceberg's CPU batch write selected for replacement. + * + *

Iceberg keeps {@code broadcastRewritableDeletes()} private on its position-delta batch + * writer, but the GPU deletion-vector writer must use the same selection when merging an + * existing deletion vector. Package placement cannot access a private member, so this helper + * uses reflection and caches the resolved method per runtime class and class loader. + * + * @return the broadcast delete-file map, or {@code null} when there are no existing deletes + * to rewrite, such as the first deletion-vector write for a data file + */ + @SuppressWarnings("unchecked") + public static Broadcast> broadcastRewritableDeletes( + DeltaBatchWrite write) { + try { + Method method = BROADCAST_REWRITABLE_DELETES_METHOD.get(write.getClass()); + return (Broadcast>) method.invoke(write); + } catch (ReflectiveOperationException e) { + throw new IllegalStateException( + "Unable to broadcast rewritable deletes from " + write.getClass().getName(), e); + } + } + + private static Method findMethod(Class targetClass, String methodName) { + Class current = targetClass; + while (current != null) { + try { + return current.getDeclaredMethod(methodName); + } catch (NoSuchMethodException e) { + current = current.getSuperclass(); + } + } + throw new IllegalStateException("No method " + methodName + " in " + targetClass.getName()); + } +} diff --git a/integration_tests/src/main/python/iceberg/iceberg_delete_test.py b/integration_tests/src/main/python/iceberg/iceberg_delete_test.py index ce680ba15c4..b7d3418a537 100644 --- a/integration_tests/src/main/python/iceberg/iceberg_delete_test.py +++ b/integration_tests/src/main/python/iceberg/iceberg_delete_test.py @@ -14,8 +14,8 @@ import pytest -from asserts import assert_equal_with_local_sort, assert_gpu_and_cpu_are_equal_collect, \ - assert_gpu_fallback_write_sql +from asserts import (assert_equal_with_local_sort, assert_gpu_and_cpu_are_equal_collect, + assert_gpu_fallback_write_sql) from conftest import is_iceberg_remote_catalog from data_gen import * from iceberg import (create_iceberg_table, get_full_table_name, iceberg_write_enabled_conf, @@ -170,6 +170,86 @@ def delete_data(spark, table_name): conf=iceberg_delete_cow_enabled_conf) +@iceberg +@ignore_order(local=True) +@pytest.mark.skipif(not supports_iceberg_v3, reason=ICEBERG_V3_UNSUPPORTED_REASON) +@pytest.mark.parametrize('fanout_enabled', [False, True], ids=['clustered', 'fanout']) +def test_iceberg_delete_v3_gpu_writes_and_merges_deletion_vectors( + spark_tmp_table_factory, fanout_enabled): + base_table_name = get_full_table_name(spark_tmp_table_factory) + cpu_table_name = f"{base_table_name}_cpu" + gpu_table_name = f"{base_table_name}_gpu" + data_gen_func = lambda spark: gen_df( + spark, [ + ('id', LongGen(nullable=False, min_val=0, max_val=127)), + ('value', LongGen()), + ('data', StringGen()) + ], seed=0) + table_properties = { + 'format-version': '3', + 'write.spark.fanout.enabled': fanout_enabled + } + create_iceberg_table_with_data( + cpu_table_name, "bucket(2, id)", data_gen_func, table_properties, + delete_mode='merge-on-read') + create_iceberg_table_with_data( + gpu_table_name, "bucket(2, id)", data_gen_func, table_properties, + delete_mode='merge-on-read') + + def delete_data(spark): + is_gpu = spark.conf.get('spark.rapids.sql.enabled') == 'true' + table_name = gpu_table_name if is_gpu else cpu_table_name + spark.sql(f"DELETE FROM {table_name} WHERE id % 3 = 0") + spark.sql(f"DELETE FROM {table_name} WHERE id % 5 = 0") + return spark.table(table_name) + + write_conf = copy_and_update(iceberg_write_enabled_conf, { + 'spark.rapids.sql.format.iceberg.v3.enabled': 'true' + }) + assert_gpu_and_cpu_are_equal_collect(delete_data, conf=write_conf) + + +@iceberg +@ignore_order(local=True) +@pytest.mark.skipif(not supports_iceberg_v3, reason=ICEBERG_V3_UNSUPPORTED_REASON) +def test_iceberg_delete_v3_gpu_upgrades_position_deletes(spark_tmp_table_factory): + base_table_name = get_full_table_name(spark_tmp_table_factory) + cpu_table_name = f"{base_table_name}_cpu" + gpu_table_name = f"{base_table_name}_gpu" + data_gen_func = lambda spark: gen_df( + spark, [ + ('id', LongGen(nullable=False, min_val=0, max_val=63)), + ('value', LongGen()), + ('data', StringGen()) + ], seed=0) + table_properties = {'format-version': '2'} + create_iceberg_table_with_data( + cpu_table_name, data_gen_func=data_gen_func, table_properties=table_properties, + delete_mode='merge-on-read') + create_iceberg_table_with_data( + gpu_table_name, data_gen_func=data_gen_func, table_properties=table_properties, + delete_mode='merge-on-read') + + def create_position_deletes(spark, table_name): + spark.sql(f"DELETE FROM {table_name} WHERE id % 4 = 0") + spark.sql( + f"ALTER TABLE {table_name} SET TBLPROPERTIES ('format-version' = '3')") + + with_cpu_session(lambda spark: create_position_deletes(spark, cpu_table_name)) + with_cpu_session(lambda spark: create_position_deletes(spark, gpu_table_name)) + write_conf = copy_and_update(iceberg_write_enabled_conf, { + 'spark.rapids.sql.format.iceberg.v3.enabled': 'true' + }) + + def delete_data(spark): + is_gpu = spark.conf.get('spark.rapids.sql.enabled') == 'true' + table_name = gpu_table_name if is_gpu else cpu_table_name + spark.sql(f"DELETE FROM {table_name} WHERE id % 5 = 0") + return spark.table(table_name) + + assert_gpu_and_cpu_are_equal_collect(delete_data, conf=write_conf) + + @iceberg @ignore_order(local=True) @pytest.mark.skipif( diff --git a/integration_tests/src/main/python/iceberg/iceberg_merge_test.py b/integration_tests/src/main/python/iceberg/iceberg_merge_test.py index 0d26fd341aa..425a785bdce 100644 --- a/integration_tests/src/main/python/iceberg/iceberg_merge_test.py +++ b/integration_tests/src/main/python/iceberg/iceberg_merge_test.py @@ -14,8 +14,8 @@ import pytest -from asserts import assert_equal_with_local_sort, assert_gpu_and_cpu_are_equal_collect, \ - assert_gpu_fallback_write_sql +from asserts import (assert_equal_with_local_sort, assert_gpu_and_cpu_are_equal_collect, + assert_gpu_fallback_write_sql) from conftest import is_iceberg_remote_catalog from data_gen import * from iceberg import (create_iceberg_table, get_full_table_name, iceberg_write_enabled_conf, @@ -227,6 +227,59 @@ def merge_data(spark, target_table): conf=iceberg_merge_enabled_conf) +@iceberg +@ignore_order(local=True) +@pytest.mark.skipif(not supports_iceberg_v3, reason=ICEBERG_V3_UNSUPPORTED_REASON) +def test_iceberg_merge_v3_gpu_writes_deletion_vectors(spark_tmp_table_factory): + base_table_name = get_full_table_name(spark_tmp_table_factory) + cpu_table_name = f"{base_table_name}_cpu" + gpu_table_name = f"{base_table_name}_gpu" + source_table = f"{base_table_name}_source" + target_data_gen = lambda spark: gen_df(spark, [ + ('id', LongGen(nullable=False, min_val=0, max_val=63)), + ('value', LongGen(nullable=False, min_val=-1000, max_val=1000)), + ('data', StringGen()), + ('flag', BooleanGen()) + ], seed=0) + source_data_gen = lambda spark: gen_df(spark, [ + ('id', UniqueLongGen()), + ('value', LongGen(nullable=False, min_val=1001, max_val=2000)), + ('data', StringGen()), + ('flag', BooleanGen()) + ], seed=1) + table_properties = { + 'format-version': '3', + 'write.merge.mode': 'merge-on-read' + } + create_iceberg_table(cpu_table_name, table_prop=table_properties, df_gen=target_data_gen) + create_iceberg_table(gpu_table_name, table_prop=table_properties, df_gen=target_data_gen) + create_iceberg_table(source_table, table_prop={'format-version': '2'}, + df_gen=source_data_gen) + + def insert_data(spark, table_name, data_gen_func): + data_gen_func(spark).writeTo(table_name).append() + + with_cpu_session(lambda spark: insert_data(spark, cpu_table_name, target_data_gen)) + with_cpu_session(lambda spark: insert_data(spark, gpu_table_name, target_data_gen)) + with_cpu_session(lambda spark: insert_data(spark, source_table, source_data_gen)) + + def merge_data(spark): + is_gpu = spark.conf.get('spark.rapids.sql.enabled') == 'true' + table_name = gpu_table_name if is_gpu else cpu_table_name + spark.sql(f""" + MERGE INTO {table_name} t + USING {source_table} s + ON t.id = s.id + WHEN MATCHED THEN UPDATE SET value = s.value + """) + return spark.table(table_name) + + write_conf = copy_and_update(iceberg_write_enabled_conf, { + 'spark.rapids.sql.format.iceberg.v3.enabled': 'true' + }) + assert_gpu_and_cpu_are_equal_collect(merge_data, conf=write_conf) + + @iceberg @ignore_order(local=True) @pytest.mark.skipif( diff --git a/integration_tests/src/main/python/iceberg/iceberg_update_test.py b/integration_tests/src/main/python/iceberg/iceberg_update_test.py index c30f1384783..11c612dd8ef 100644 --- a/integration_tests/src/main/python/iceberg/iceberg_update_test.py +++ b/integration_tests/src/main/python/iceberg/iceberg_update_test.py @@ -14,8 +14,8 @@ import pytest -from asserts import assert_equal_with_local_sort, assert_gpu_and_cpu_are_equal_collect, \ - assert_gpu_fallback_write_sql +from asserts import (assert_equal_with_local_sort, assert_gpu_and_cpu_are_equal_collect, + assert_gpu_fallback_write_sql) from conftest import is_iceberg_remote_catalog from data_gen import * from iceberg import (create_iceberg_table, get_full_table_name, iceberg_write_enabled_conf, @@ -162,6 +162,39 @@ def update_data(spark, table_name): conf=iceberg_update_cow_enabled_conf) +@iceberg +@ignore_order(local=True) +@pytest.mark.skipif(not supports_iceberg_v3, reason=ICEBERG_V3_UNSUPPORTED_REASON) +def test_iceberg_update_v3_gpu_writes_deletion_vectors(spark_tmp_table_factory): + base_table_name = get_full_table_name(spark_tmp_table_factory) + cpu_table_name = f"{base_table_name}_cpu" + gpu_table_name = f"{base_table_name}_gpu" + data_gen_func = lambda spark: gen_df(spark, [ + ('id', LongGen(nullable=False, min_val=0, max_val=63)), + ('value', LongGen(nullable=False, min_val=-1000, max_val=1000)), + ('data', StringGen()), + ('flag', BooleanGen()) + ], seed=0) + table_properties = {'format-version': '3'} + create_iceberg_table_with_data( + cpu_table_name, data_gen_func=data_gen_func, table_properties=table_properties, + update_mode='merge-on-read') + create_iceberg_table_with_data( + gpu_table_name, data_gen_func=data_gen_func, table_properties=table_properties, + update_mode='merge-on-read') + + def update_data(spark): + is_gpu = spark.conf.get('spark.rapids.sql.enabled') == 'true' + table_name = gpu_table_name if is_gpu else cpu_table_name + spark.sql(f"UPDATE {table_name} SET value = value + 1000 WHERE id % 4 = 0") + return spark.table(table_name) + + write_conf = copy_and_update(iceberg_write_enabled_conf, { + 'spark.rapids.sql.format.iceberg.v3.enabled': 'true' + }) + assert_gpu_and_cpu_are_equal_collect(update_data, conf=write_conf) + + @iceberg @ignore_order(local=True) @pytest.mark.skipif( diff --git a/scala2.13/iceberg/iceberg-1-10-x/pom.xml b/scala2.13/iceberg/iceberg-1-10-x/pom.xml index 5ef88591020..42a2e6cd6a3 100644 --- a/scala2.13/iceberg/iceberg-1-10-x/pom.xml +++ b/scala2.13/iceberg/iceberg-1-10-x/pom.xml @@ -84,6 +84,7 @@ ${spark.rapids.source.basedir}/iceberg/common/src/main/java ${spark.rapids.source.basedir}/iceberg/common/src/main/scala + ${spark.rapids.source.basedir}/iceberg/iceberg-19plus-common/src/main/java diff --git a/scala2.13/iceberg/iceberg-1-11-x/pom.xml b/scala2.13/iceberg/iceberg-1-11-x/pom.xml index 73a3ff63484..c17c4a510c5 100644 --- a/scala2.13/iceberg/iceberg-1-11-x/pom.xml +++ b/scala2.13/iceberg/iceberg-1-11-x/pom.xml @@ -84,6 +84,7 @@ ${spark.rapids.source.basedir}/iceberg/common/src/main/java ${spark.rapids.source.basedir}/iceberg/common/src/main/scala + ${spark.rapids.source.basedir}/iceberg/iceberg-19plus-common/src/main/java diff --git a/scala2.13/iceberg/iceberg-1-9-x/pom.xml b/scala2.13/iceberg/iceberg-1-9-x/pom.xml index 989a536f6d5..ba311489ca7 100644 --- a/scala2.13/iceberg/iceberg-1-9-x/pom.xml +++ b/scala2.13/iceberg/iceberg-1-9-x/pom.xml @@ -74,6 +74,7 @@ ${spark.rapids.source.basedir}/iceberg/common/src/main/java ${spark.rapids.source.basedir}/iceberg/common/src/main/scala + ${spark.rapids.source.basedir}/iceberg/iceberg-19plus-common/src/main/java