Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
23 commits
Select commit Hold shift + click to select a range
a3d1918
Support Iceberg v3 deletion vectors on GPU
liurenjie1024 Aug 12, 2026
15d4185
Use native cuDF Parquet deletion vectors for Iceberg
liurenjie1024 Aug 17, 2026
36db87f
Address Iceberg deletion vector review feedback
liurenjie1024 Aug 17, 2026
aaa1c8d
Refine Iceberg deletion vector integration
liurenjie1024 Aug 18, 2026
08d297c
Update Iceberg reader copyright year
liurenjie1024 Aug 18, 2026
5640f1f
Address Iceberg deletion vector review feedback
liurenjie1024 Aug 19, 2026
85b21db
Address latest Iceberg review feedback
liurenjie1024 Aug 19, 2026
840bd99
Add standalone Iceberg v3 deletion vector coverage
liurenjie1024 Aug 20, 2026
33177e7
Address latest Iceberg review feedback
liurenjie1024 Aug 20, 2026
33096d1
Merge branch 'ray/15440' into ray/15442
liurenjie1024 Aug 20, 2026
5de16a7
Add Iceberg v3 deletion vector writes
liurenjie1024 Aug 20, 2026
a7b2cdb
Merge remote-tracking branch 'upstream/main' into ray/15442
liurenjie1024 Aug 26, 2026
6223f68
Address Iceberg deletion vector review feedback
liurenjie1024 Aug 26, 2026
1133ee1
Address additional deletion vector review feedback
liurenjie1024 Aug 27, 2026
35c576c
Address Iceberg deletion vector review feedback
liurenjie1024 Aug 27, 2026
714fd4a
Address deletion vector shim review feedback
liurenjie1024 Aug 27, 2026
f6c9a3a
Refine Iceberg deletion vector tests
liurenjie1024 Aug 28, 2026
586d2d5
Refine Iceberg deletion vector test data
liurenjie1024 Aug 31, 2026
66d6b94
Use standard assertions for Iceberg v3 DML tests
liurenjie1024 Aug 31, 2026
af70721
Merge upstream/main into ray/15442
liurenjie1024 Sep 1, 2026
4d0fdfc
Merge upstream/main into ray/15442
liurenjie1024 Sep 1, 2026
2d8bd6d
Refactor Iceberg 1.9+ deletion vector shims
liurenjie1024 Sep 2, 2026
4c95fd1
Merge remote-tracking branch 'upstream/main' into ray/15442
liurenjie1024 Sep 4, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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;

/**
Expand All @@ -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);
Comment thread
liurenjie1024 marked this conversation as resolved.

/**
* 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(
Comment thread
liurenjie1024 marked this conversation as resolved.
DeltaBatchWrite write);

/**
* Opaque, serializable handle for version-specific rewritable-delete state.
*
* <p>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<PositionDelete<InternalRow>, 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<InternalRow> delete, CharSequence path, long position);

/**
* Reads exactly the recorded deletion-vector byte range and returns its compressed bitmap.
*
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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<PositionDelete<InternalRow>, 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<InternalRow> delete, CharSequence path, long position) {
IMPL.setPositionDelete(delete, path, position);
}

public static IcebergDeletionVector readDeletionVector(DeleteFile deleteFile,
RapidsInputFile inputFile,
boolean validateCrc)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -111,4 +114,4 @@ class GpuSparkFileWriterFactory(val table: Table,
fileIO = new IcebergFileIO(table.io())
)
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand All @@ -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
Expand Down Expand Up @@ -131,6 +138,7 @@ class GpuSparkPositionDeltaWrite(cpu: DeltaWrite)

new GpuPositionDeltaWriterFactory(
tableBroadcast,
rewritableDeletes,
command,
context,
writeProps,
Expand Down Expand Up @@ -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 = {
Expand All @@ -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],
Expand All @@ -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)
Expand All @@ -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]]
}
}
Expand All @@ -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] = {
Expand All @@ -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)
Expand Down Expand Up @@ -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] {

Comment thread
liurenjie1024 marked this conversation as resolved.
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.
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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,
Expand All @@ -853,6 +910,7 @@ object GpuWriteContext {
deleteGranularity,
queryId,
useFanoutWriter,
inputOrdered)
inputOrdered,
useDVs)
}
}
Loading
Loading