Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
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
6 changes: 0 additions & 6 deletions dist/unshimmed-common-from-single-shim.txt
Original file line number Diff line number Diff line change
Expand Up @@ -11,15 +11,9 @@ com/nvidia/spark/rapids/optimizer/SQLOptimizerPlugin*
com/nvidia/spark/rapids/ShimLoaderTemp*
com/nvidia/spark/rapids/SparkShims*
com/nvidia/shaded/spark/com/google/flatbuffers/*.class
com/nvidia/spark/rapids/iceberg/data/GpuDeleteFilter2.class
com/nvidia/spark/rapids/iceberg/parquet/GpuIcebergParquetReader$.class
com/nvidia/spark/rapids/iceberg/parquet/SingleFile.class
com/nvidia/spark/rapids/iceberg/parquet/SingleFile$.class
com/nvidia/spark/rapids/iceberg/parquet/ThreadConf.class
com/nvidia/spark/rapids/iceberg/spark/RapidsSparkCatalog.class
com/nvidia/spark/rapids/iceberg/spark/RapidsSparkSessionCatalog.class
com/nvidia/spark/rapids/iceberg/spark/source/RapidsSparkTable.class
org/apache/iceberg/io/GpuPositionDeleteFileWriter$.class
org/apache/iceberg/spark/source/GpuBaseReader.class
org/apache/iceberg/spark/source/GpuSparkPlanningUtil.class
org/apache/spark/sql/rapids/AdaptiveSparkPlanHelperShim*
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -27,7 +27,8 @@
import java.util.zip.CRC32;

/**
* An Iceberg deletion vector kept in its compressed Roaring-bitmap representation.
* An Iceberg deletion vector kept in its compressed Roaring-bitmap representation. This helper
* does not depend on Iceberg runtime classes.
*
* <p>The serialized bytes use the portable 64-bit Roaring format expected by cuDF. This object
* owns its host buffer and must be closed after all borrowed references have been released.
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,72 @@
/*
* Copyright (c) 2025-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.parquet

import scala.collection.JavaConverters._

import org.apache.hadoop.conf.Configuration
import org.apache.iceberg.MetadataColumns
import org.apache.iceberg.hadoop.HadoopInputFile
import org.apache.iceberg.io.InputFile
import org.apache.iceberg.shaded.org.apache.parquet.{HadoopReadOptions, ParquetReadOptions}
import org.apache.iceberg.shaded.org.apache.parquet.schema.{
MessageType => ShadedMessageType, Types => ShadedTypes}
import org.apache.iceberg.shaded.org.apache.parquet.schema.PrimitiveType.{
PrimitiveTypeName => ShadedPrimitiveTypeName}
import org.apache.iceberg.shaded.org.apache.parquet.schema.Type.{Repetition => ShadedRepetition}

object GpuIcebergParquetReaderUtils {
private val READ_PROPERTIES_TO_REMOVE = Set(
"parquet.read.filter",
"parquet.private.read.filter.predicate",
"parquet.read.support.class")

/**
* Adds the leading file-global row index emitted by the cuDF deletion-vector reader to the
* schema consumed by the Iceberg post-processor.
*/
private[iceberg] def withNativeRowIndex(
fileReadSchema: ShadedMessageType): ShadedMessageType = {
val rowPosition = ShadedTypes
.primitive(ShadedPrimitiveTypeName.INT64, ShadedRepetition.REQUIRED)
.id(MetadataColumns.ROW_POSITION.fieldId())
.named(MetadataColumns.ROW_POSITION.name())
new ShadedMessageType(
fileReadSchema.getName,
(rowPosition +: fileReadSchema.getFields.asScala).asJava)
}

def buildReaderOptions(file: InputFile, split: Option[(Long, Long)])
: ParquetReadOptions = {
var optionsBuilder: ParquetReadOptions.Builder = null
file match {
case hadoop: HadoopInputFile =>
// remove read properties already set that may conflict with this read
val conf = new Configuration(hadoop.getConf)
for (property <- READ_PROPERTIES_TO_REMOVE) {
conf.unset(property)
}
optionsBuilder = HadoopReadOptions.builder(conf)
case _ =>
optionsBuilder = ParquetReadOptions.builder()
}
split.foreach { case (start, length) =>
optionsBuilder = optionsBuilder.withRange(start, start + length)
}
optionsBuilder.build
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,21 @@
/*
* Copyright (c) 2025-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.parquet

trait ThreadConf

case object SingleFile extends ThreadConf
Original file line number Diff line number Diff line change
@@ -0,0 +1,24 @@
/*
* Copyright (c) 2025-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.io

import org.apache.iceberg.MetadataColumns.{DELETE_FILE_PATH, DELETE_FILE_POS}

private[io] object GpuPositionDeleteFieldIds {
val FILE_AND_POS_FIELD_IDS: Set[Integer] = Set(
DELETE_FILE_PATH.fieldId(), DELETE_FILE_POS.fieldId())
}
Original file line number Diff line number Diff line change
Expand Up @@ -24,7 +24,7 @@ import com.nvidia.spark.rapids.Arm.{closeOnExcept, withResource}
import com.nvidia.spark.rapids.GpuMetric.{JOIN_TIME, OP_TIME_LEGACY}
import com.nvidia.spark.rapids.fileio.iceberg.{IcebergFileIO, IcebergInputFile}
import com.nvidia.spark.rapids.iceberg.ShimUtils
import com.nvidia.spark.rapids.iceberg.data.GpuDeleteFilter2._
import com.nvidia.spark.rapids.iceberg.data.GpuDeleteFileInfo._
import com.nvidia.spark.rapids.iceberg.fieldIndex
import com.nvidia.spark.rapids.iceberg.parquet.GpuIcebergParquetReaderConf
import org.apache.iceberg.{DeleteFile, FileContent, MetadataColumns, Schema}
Expand Down Expand Up @@ -72,6 +72,68 @@ object GpuDeleteFileInfo {
val effectivePositionDeletes = if (deletionVector.isDefined) Seq.empty else positionDeletes
new GpuDeleteFileInfo(deletionVector, equalityDeletes.toSeq ++ effectivePositionDeletes)
}

private[iceberg] val DELETE_EXTRA_METADATA_COLUMNS: Seq[NestedField] = Seq(
MetadataColumns.FILE_PATH,
MetadataColumns.ROW_POSITION)

private[iceberg] val DELETE_EXTRA_METADATA_COLUMN_IDS: Set[Int] =
DELETE_EXTRA_METADATA_COLUMNS
.map(_.fieldId())
.toSet

private[iceberg] val POS_DELETE_SCHEMA: Schema = new Schema(
MetadataColumns.DELETE_FILE_PATH,
MetadataColumns.DELETE_FILE_POS)


private[iceberg] def mergeColumn(
batch: ColumnarBatch, srcColIdx: Int, destColIdx: Int)
(mergeOp: (GpuColumnVector, GpuColumnVector) => GpuColumnVector): ColumnarBatch = {
require(srcColIdx >= 0 && srcColIdx < batch.numCols(),
s"Invalid src column index: $srcColIdx, numCols: ${batch.numCols()}")
require(destColIdx >= 0 && destColIdx < batch.numCols(),
s"Invalid dest column index: $destColIdx, numCols: ${batch.numCols()}")
require(srcColIdx != destColIdx, "srcColIdx and destColIdx should be different")

val srcVec = batch.column(srcColIdx).asInstanceOf[GpuColumnVector]
val destVec = batch.column(destColIdx).asInstanceOf[GpuColumnVector]

withResource(batch) { _ =>
closeOnExcept(mergeOp(srcVec, destVec)) { mergeVec =>
val newColumns = new Array[ColumnVector](batch.numCols() - 1)
for (i <- 0 until batch.numCols() - 1) {
if (i == destColIdx) {
newColumns(i) = mergeVec
} else {
newColumns(i) = batch.column(i).asInstanceOf[GpuColumnVector].incRefCount()
}
}
new ColumnarBatch(newColumns, batch.numRows())
}
}
}

private[iceberg] def filterAndDrop(batch: ColumnarBatch,
isDeletedColIdx: Int,
outputDataType: Array[DataType],
dropMask: Array[Boolean] = Array.empty): ColumnarBatch = {
withResource(batch) { _ =>
withResource(GpuColumnVector.from(batch)) { table =>
withResource(table.getColumn(isDeletedColIdx).not()) { maskCv =>
withResource(table.filter(maskCv)) { newTable =>
if (dropMask.nonEmpty) {
withResource(GpuColumnVector.from(newTable, outputDataType)) { newBatch =>
GpuColumnVector.dropColumns(newBatch, dropMask)
}
} else {
GpuColumnVector.from(newTable, outputDataType)
}
}
}
}
}
}
}

class GpuDeleteFilter(
Expand Down Expand Up @@ -361,70 +423,6 @@ class GpuDeleteFilter(
}
}

object GpuDeleteFilter2 {
private[iceberg] val DELETE_EXTRA_METADATA_COLUMNS: Seq[NestedField] = Seq(
MetadataColumns.FILE_PATH,
MetadataColumns.ROW_POSITION)

private[iceberg] val DELETE_EXTRA_METADATA_COLUMN_IDS: Set[Int] =
DELETE_EXTRA_METADATA_COLUMNS
.map(_.fieldId())
.toSet

private[iceberg] val POS_DELETE_SCHEMA: Schema = new Schema(
MetadataColumns.DELETE_FILE_PATH,
MetadataColumns.DELETE_FILE_POS)


private[iceberg] def mergeColumn(
batch: ColumnarBatch, srcColIdx: Int, destColIdx: Int)
(mergeOp: (GpuColumnVector, GpuColumnVector) => GpuColumnVector): ColumnarBatch = {
require(srcColIdx >= 0 && srcColIdx < batch.numCols(),
s"Invalid src column index: $srcColIdx, numCols: ${batch.numCols()}")
require(destColIdx >= 0 && destColIdx < batch.numCols(),
s"Invalid dest column index: $destColIdx, numCols: ${batch.numCols()}")
require(srcColIdx != destColIdx, "srcColIdx and destColIdx should be different")

val srcVec = batch.column(srcColIdx).asInstanceOf[GpuColumnVector]
val destVec = batch.column(destColIdx).asInstanceOf[GpuColumnVector]

withResource(batch) { _ =>
closeOnExcept(mergeOp(srcVec, destVec)) { mergeVec =>
val newColumns = new Array[ColumnVector](batch.numCols() - 1)
for (i <- 0 until batch.numCols() - 1) {
if (i == destColIdx) {
newColumns(i) = mergeVec
} else {
newColumns(i) = batch.column(i).asInstanceOf[GpuColumnVector].incRefCount()
}
}
new ColumnarBatch(newColumns, batch.numRows())
}
}
}

private[iceberg] def filterAndDrop(batch: ColumnarBatch,
isDeletedColIdx: Int,
outputDataType: Array[DataType],
dropMask: Array[Boolean] = Array.empty): ColumnarBatch = {
withResource(batch) { _ =>
withResource(GpuColumnVector.from(batch)) { table =>
withResource(table.getColumn(isDeletedColIdx).not()) { maskCv =>
withResource(table.filter(maskCv)) { newTable =>
if (dropMask.nonEmpty) {
withResource(GpuColumnVector.from(newTable, outputDataType)) { newBatch =>
GpuColumnVector.dropColumns(newBatch, dropMask)
}
} else {
GpuColumnVector.from(newTable, outputDataType)
}
}
}
}
}
}
}

private case class DeleteFilterContext(
buildBatch: LazySpillableColumnarBatch,
buildKeys: Seq[GpuExpression],
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -34,18 +34,12 @@ import org.apache.hadoop.conf.Configuration
import org.apache.hadoop.fs.Path
import org.apache.iceberg.{MetadataColumns, Schema}
import org.apache.iceberg.expressions.Expression
import org.apache.iceberg.hadoop.HadoopInputFile
import org.apache.iceberg.io.InputFile
import org.apache.iceberg.mapping.NameMapping
import org.apache.iceberg.parquet._
import org.apache.iceberg.shaded.org.apache.parquet.{HadoopReadOptions, ParquetReadOptions}
import org.apache.iceberg.shaded.org.apache.parquet.ParquetReadOptions
import org.apache.iceberg.shaded.org.apache.parquet.hadoop.ParquetFileReader
import org.apache.iceberg.shaded.org.apache.parquet.hadoop.metadata.{BlockMetaData => ShadedBlockMetaData}
import org.apache.iceberg.shaded.org.apache.parquet.schema.{
MessageType => ShadedMessageType, Types => ShadedTypes}
import org.apache.iceberg.shaded.org.apache.parquet.schema.PrimitiveType.{
PrimitiveTypeName => ShadedPrimitiveTypeName}
import org.apache.iceberg.shaded.org.apache.parquet.schema.Type.{Repetition => ShadedRepetition}
import org.apache.iceberg.shaded.org.apache.parquet.schema.{MessageType => ShadedMessageType}
import org.apache.parquet.hadoop.metadata.BlockMetaData

import org.apache.spark.internal.Logging
Expand All @@ -64,7 +58,7 @@ case class IcebergPartitionedFile(
lazy val path: Path = new Path(new URI(urlEncodedPath))

def parquetReadOptions: ParquetReadOptions = {
GpuIcebergParquetReader.buildReaderOptions(file.getDelegate, split)
GpuIcebergParquetReaderUtils.buildReaderOptions(file.getDelegate, split)
}

def newReader(metrics: Map[String, GpuMetric] = Map.empty): ParquetFileReader = {
Expand Down Expand Up @@ -116,10 +110,6 @@ case class IcebergPartitionedFile(
}
}

sealed trait ThreadConf

case object SingleFile extends ThreadConf

case class MultiThread(
poolConfBuilder: ThreadPoolConfBuilder,
maxNumFilesProcessed: Int,
Expand Down Expand Up @@ -302,53 +292,11 @@ trait GpuIcebergParquetReader extends Iterator[ColumnarBatch] with AutoCloseable
)

val postProcessorReadSchema = if (hasDeletionVector) {
GpuIcebergParquetReader.withNativeRowIndex(fileReadSchema)
GpuIcebergParquetReaderUtils.withNativeRowIndex(fileReadSchema)
} else {
fileReadSchema
}
(parquetFileInfo, postProcessorReadSchema)
}
}
}

object GpuIcebergParquetReader {
private val READ_PROPERTIES_TO_REMOVE = Set(
"parquet.read.filter",
"parquet.private.read.filter.predicate",
"parquet.read.support.class")

/**
* Adds the leading file-global row index emitted by the cuDF deletion-vector reader to the
* schema consumed by the Iceberg post-processor.
*/
private[iceberg] def withNativeRowIndex(
fileReadSchema: ShadedMessageType): ShadedMessageType = {
val rowPosition = ShadedTypes
.primitive(ShadedPrimitiveTypeName.INT64, ShadedRepetition.REQUIRED)
.id(MetadataColumns.ROW_POSITION.fieldId())
.named(MetadataColumns.ROW_POSITION.name())
new ShadedMessageType(
fileReadSchema.getName,
(rowPosition +: fileReadSchema.getFields.asScala).asJava)
}

def buildReaderOptions(file: InputFile, split: Option[(Long, Long)])
: ParquetReadOptions = {
var optionsBuilder: ParquetReadOptions.Builder = null
file match {
case hadoop: HadoopInputFile =>
// remove read properties already set that may conflict with this read
val conf = new Configuration(hadoop.getConf)
for (property <- READ_PROPERTIES_TO_REMOVE) {
conf.unset(property)
}
optionsBuilder = HadoopReadOptions.builder(conf)
case _ =>
optionsBuilder = ParquetReadOptions.builder()
}
split.foreach { case (start, length) =>
optionsBuilder = optionsBuilder.withRange(start, start + length)
}
optionsBuilder.build
}
}
Loading
Loading