Skip to content

Commit 498b8fe

Browse files
committed
Prefetch Iceberg Parquet footers with a suffix read
Signed-off-by: Zach Puller <zpuller@nvidia.com>
1 parent 8e360c1 commit 498b8fe

5 files changed

Lines changed: 65 additions & 8 deletions

File tree

iceberg/common/src/main/java/com/nvidia/spark/rapids/fileio/iceberg/IcebergInputFile.java

Lines changed: 24 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,7 @@
1616

1717
package com.nvidia.spark.rapids.fileio.iceberg;
1818

19+
import ai.rapids.cudf.HostMemoryBuffer;
1920
import com.nvidia.spark.rapids.jni.fileio.RapidsInputFile;
2021
import com.nvidia.spark.rapids.jni.fileio.SeekableInputStream;
2122
import org.apache.iceberg.io.InputFile;
@@ -31,10 +32,16 @@
3132
*/
3233
public class IcebergInputFile implements RapidsInputFile {
3334
private final InputFile delegate;
35+
private final SuffixReader suffixReader;
3436

3537
public IcebergInputFile(InputFile delegate) {
38+
this(delegate, null);
39+
}
40+
41+
public IcebergInputFile(InputFile delegate, SuffixReader suffixReader) {
3642
Objects.requireNonNull(delegate, "delegate can't be null");
3743
this.delegate = delegate;
44+
this.suffixReader = suffixReader;
3845
}
3946

4047
@Override
@@ -60,4 +67,21 @@ public SeekableInputStream open() throws IOException {
6067
public InputFile getDelegate() {
6168
return delegate;
6269
}
70+
71+
public boolean supportsSuffixReads() {
72+
return suffixReader != null;
73+
}
74+
75+
public long readSuffix(long length, HostMemoryBuffer output, long outputOffset)
76+
throws IOException {
77+
if (suffixReader == null) {
78+
throw new UnsupportedOperationException("Suffix reads are not supported");
79+
}
80+
return suffixReader.read(length, output, outputOffset);
81+
}
82+
83+
@FunctionalInterface
84+
public interface SuffixReader {
85+
long read(long length, HostMemoryBuffer output, long outputOffset) throws IOException;
86+
}
6387
}

iceberg/common/src/main/java/org/apache/iceberg/aws/s3/IcebergS3InputFile.java

Lines changed: 18 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -85,6 +85,16 @@ public static RapidsInputFile maybeCreate(InputFile inputFile, FileIO fileIO) {
8585
return new IcebergS3InputFile(delegate, s3Uri, icebergS3Client);
8686
}
8787

88+
public static IcebergInputFile maybeCreateIcebergInputFile(
89+
InputFile inputFile, FileIO fileIO) {
90+
RapidsInputFile optimizedFile = maybeCreate(inputFile, fileIO);
91+
if (optimizedFile instanceof IcebergS3InputFile) {
92+
IcebergS3InputFile s3File = (IcebergS3InputFile) optimizedFile;
93+
return new IcebergInputFile(inputFile, s3File::readTail);
94+
}
95+
return (IcebergInputFile) optimizedFile;
96+
}
97+
8898
@Override
8999
public String path() {
90100
return delegate.path();
@@ -127,12 +137,18 @@ public void readVectored(HostMemoryBuffer output, List<CopyRange> copyRanges)
127137
*/
128138
@Override
129139
public void readTail(long length, HostMemoryBuffer output) throws IOException {
140+
readTail(length, output, 0L);
141+
}
142+
143+
public long readTail(long length, HostMemoryBuffer output, long outputOffset)
144+
throws IOException {
130145
if (length == 0) {
131-
return;
146+
return 0;
132147
}
133148
if (length < 0) {
134149
throw new IllegalArgumentException("length must be non-negative");
135150
}
136-
IcebergS3RangeCopier.copyTailToHMB(icebergS3Client, output, s3Uri, length, /*dstOffset*/ 0L);
151+
return IcebergS3RangeCopier.copyTailToHMB(
152+
icebergS3Client, output, s3Uri, length, outputOffset);
137153
}
138154
}

iceberg/common/src/main/scala/org/apache/iceberg/spark/source/GpuIcebergPartitionReader.scala

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -20,12 +20,13 @@ import scala.collection.JavaConverters._
2020

2121
import com.nvidia.spark.rapids.GpuMetric
2222
import com.nvidia.spark.rapids.MapUtil.toMapStrict
23-
import com.nvidia.spark.rapids.fileio.iceberg.{IcebergFileIO, IcebergInputFile}
23+
import com.nvidia.spark.rapids.fileio.iceberg.IcebergFileIO
2424
import com.nvidia.spark.rapids.iceberg.ShimUtils
2525
import com.nvidia.spark.rapids.iceberg.ShimUtils.locationOf
2626
import com.nvidia.spark.rapids.iceberg.data.GpuDeleteFilter
2727
import com.nvidia.spark.rapids.iceberg.parquet._
2828
import org.apache.iceberg._
29+
import org.apache.iceberg.aws.s3.IcebergS3InputFile
2930
import org.apache.iceberg.encryption.EncryptedFiles
3031
import org.apache.iceberg.mapping.NameMappingParser
3132

@@ -110,7 +111,7 @@ class GpuIcebergPartitionReader(private val task: GpuSparkInputPartition,
110111
val inputFiles = table.encryption()
111112
.decrypt(encryptedFiles.asJava)
112113
.asScala
113-
.map(f => f.location() -> new IcebergInputFile(f))
114+
.map(f => f.location() -> IcebergS3InputFile.maybeCreateIcebergInputFile(f, fileIO))
114115
.toMap
115116

116117
val taskMap = toMapStrict(tasks.map(t => {

iceberg/iceberg-1-10-x/src/main/scala/com/nvidia/spark/rapids/iceberg/iceberg110x/GpuParquetIOShim.scala

Lines changed: 10 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -17,7 +17,7 @@
1717
package com.nvidia.spark.rapids.iceberg.iceberg110x
1818

1919
import com.nvidia.spark.rapids.Arm.{closeOnExcept, withResource}
20-
import com.nvidia.spark.rapids.GpuMetric
20+
import com.nvidia.spark.rapids.{GpuMetric, PerfIO}
2121
import com.nvidia.spark.rapids.fileio.iceberg.IcebergInputFile
2222
import com.nvidia.spark.rapids.iceberg.parquet.converter.ToIcebergShaded
2323
import com.nvidia.spark.rapids.parquet.{HMBInputFile, ParquetFooterUtils}
@@ -41,7 +41,15 @@ object GpuParquetIOShim {
4141
metrics: Map[String, GpuMetric]): ParquetFileReader = {
4242
val metadata = withResource(ParquetFooterUtils.getFooterBuffer(
4343
inputFile, metrics,
44-
ParquetFooterUtils.readFooterBufferFromInputFile(inputFile, filePath))) { hmb =>
44+
if (inputFile.supportsSuffixReads) {
45+
PerfIO.readParquetFooterBufferFromTail(
46+
filePath,
47+
(length, output, outputOffset) =>
48+
inputFile.readSuffix(length, output, outputOffset),
49+
ParquetFooterUtils.verifyParquetMagic)
50+
} else {
51+
ParquetFooterUtils.readFooterBufferFromInputFile(inputFile, filePath)
52+
})) { hmb =>
4553
val shadedHmbFile = ToIcebergShaded.shade(new HMBInputFile(hmb))
4654
withResource(shadedHmbFile.newStream()) { hmbStream =>
4755
ParquetFileReader.readFooter(shadedHmbFile, options, hmbStream)

iceberg/iceberg-1-11-x/src/main/scala/com/nvidia/spark/rapids/iceberg/iceberg111x/GpuParquetIOShim.scala

Lines changed: 10 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -17,7 +17,7 @@
1717
package com.nvidia.spark.rapids.iceberg.iceberg111x
1818

1919
import com.nvidia.spark.rapids.Arm.{closeOnExcept, withResource}
20-
import com.nvidia.spark.rapids.GpuMetric
20+
import com.nvidia.spark.rapids.{GpuMetric, PerfIO}
2121
import com.nvidia.spark.rapids.fileio.iceberg.IcebergInputFile
2222
import com.nvidia.spark.rapids.iceberg.parquet.converter.ToIcebergShaded
2323
import com.nvidia.spark.rapids.parquet.{HMBInputFile, ParquetFooterUtils}
@@ -41,7 +41,15 @@ object GpuParquetIOShim {
4141
metrics: Map[String, GpuMetric]): ParquetFileReader = {
4242
val metadata = withResource(ParquetFooterUtils.getFooterBuffer(
4343
inputFile, metrics,
44-
ParquetFooterUtils.readFooterBufferFromInputFile(inputFile, filePath))) { hmb =>
44+
if (inputFile.supportsSuffixReads) {
45+
PerfIO.readParquetFooterBufferFromTail(
46+
filePath,
47+
(length, output, outputOffset) =>
48+
inputFile.readSuffix(length, output, outputOffset),
49+
ParquetFooterUtils.verifyParquetMagic)
50+
} else {
51+
ParquetFooterUtils.readFooterBufferFromInputFile(inputFile, filePath)
52+
})) { hmb =>
4553
val shadedHmbFile = ToIcebergShaded.shade(new HMBInputFile(hmb))
4654
withResource(shadedHmbFile.newStream()) { hmbStream =>
4755
ParquetFileReader.readFooter(shadedHmbFile, options, hmbStream)

0 commit comments

Comments
 (0)