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
Original file line number Diff line number Diff line change
Expand Up @@ -63,16 +63,6 @@ public IcebergInputFile newInputFile(InputFile inputFile) {
return IcebergS3InputFile.maybeCreate(inputFile, delegate);
}

/**
* Always returns a plain {@link IcebergInputFile}, bypassing the S3 PerfIO
* fast-path. Use this from internal call sites that need the iceberg
* {@link InputFile} accessor on the read side (e.g. writers re-reading the
* footer of a just-written file) and do not benefit from PerfIO.
*/
public IcebergInputFile newIcebergInputFile(String path) throws IOException {
return new IcebergInputFile(delegate.newInputFile(path));
}

@Override
public IcebergOutputFile newOutputFile(String path) throws IOException {
return new IcebergOutputFile(delegate.newOutputFile(path));
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -22,10 +22,12 @@ import java.util.stream.{Stream => JStream}

import com.nvidia.spark.rapids.{GpuParquetWriter, SpillableColumnarBatch}
import com.nvidia.spark.rapids.Arm.withResource
import com.nvidia.spark.rapids.fileio.iceberg.IcebergFileIO
import com.nvidia.spark.rapids.iceberg.parquet.converter.ToIcebergShaded
import com.nvidia.spark.rapids.parquet.HMBInputFile
import org.apache.iceberg.{FieldMetrics, Metrics, MetricsConfig}
import org.apache.iceberg.io.FileAppender
import org.apache.iceberg.parquet.ParquetUtil
import org.apache.iceberg.shaded.org.apache.parquet.hadoop.ParquetFileReader
import org.apache.iceberg.shaded.org.apache.parquet.hadoop.metadata.ParquetMetadata


Expand All @@ -39,8 +41,7 @@ import org.apache.iceberg.shaded.org.apache.parquet.hadoop.metadata.ParquetMetad
*/
class GpuIcebergParquetAppender(
val inner: GpuParquetWriter,
val metricsConfig: MetricsConfig,
val fileIO: IcebergFileIO) extends FileAppender[SpillableColumnarBatch] {
val metricsConfig: MetricsConfig) extends FileAppender[SpillableColumnarBatch] {
private var closed = false
private var footer: ParquetMetadata = _

Expand All @@ -58,12 +59,11 @@ class GpuIcebergParquetAppender(

override def close(): Unit = {
if (!closed) {
inner.close()
footer = withResource(
IcebergPartitionedFile(fileIO.newIcebergInputFile(inner.path)).newReader()) {
reader =>
// TODO: Remove the read after https://github.com/rapidsai/cudf/issues/18886 got fixed.
footer = withResource(inner.closeAndGetFooter()) { footerBuffer =>
withResource(ParquetFileReader.open(
ToIcebergShaded.shade(new HMBInputFile(footerBuffer)))) { reader =>
reader.getFooter
}
}
closed = true
}
Expand All @@ -74,4 +74,4 @@ class GpuIcebergParquetAppender(

ParquetUtil.getSplitOffsets(footer)
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -107,8 +107,7 @@ class GpuSparkFileWriterFactory(val table: Table,

new GpuIcebergParquetAppender(
gpuWriter,
metricsConfig = MetricsConfig.forTable(table),
fileIO = new IcebergFileIO(table.io())
metricsConfig = MetricsConfig.forTable(table)
)
}
}
}
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
/*
* Copyright (c) 2019-2025, NVIDIA CORPORATION.
* Copyright (c) 2019-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.
Expand Down Expand Up @@ -281,12 +281,14 @@ abstract class ColumnarOutputWriter(context: TaskAttemptContext,
* Closes the [[ColumnarOutputWriter]]. Invoked on the executor side after all columnar batches
* are persisted, before the task output is committed.
*/
def close(): Unit = {
private def prepareToClose(): Unit = {
if (!anythingWritten) {
// This prevents writing out bad files
bufferBatchAndClose(GpuColumnVector.emptyBatch(dataSchema))
}
tableWriter.close()
}

private def finishClose(): Unit = {
GpuSemaphore.releaseIfNecessary(TaskContext.get())
writeBufferedData()
outputStream.close()
Expand All @@ -295,6 +297,20 @@ abstract class ColumnarOutputWriter(context: TaskAttemptContext,
}
}

protected final def closeAndReturn[T <: AutoCloseable](closeWriter: => T): T = {
prepareToClose()
closeOnExcept(closeWriter) { result =>
finishClose()
result
}
}

def close(): Unit = {
prepareToClose()
tableWriter.close()
finishClose()
}

/**
* The file path to write. Invoked on the executor side.
*/
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -478,7 +478,7 @@ class GpuParquetWriter(
}
}

override val tableWriter: TableWriter = {
override val tableWriter: ParquetTableWriter = {
val writeContext = new ParquetWriteSupport().init(conf)
// Use the value resolved in prepareWrite. Reading SQLConf here would discard
// a per-write option.
Expand All @@ -503,4 +503,8 @@ class GpuParquetWriter(
}
Table.writeParquetChunked(builder.build(), this)
}

private[rapids] def closeAndGetFooter(): HostMemoryBuffer = {
closeAndReturn(tableWriter.closeAndGetFooter())
}
}
Loading