Skip to content

Commit 8460dcd

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

12 files changed

Lines changed: 60 additions & 16 deletions

File tree

iceberg/common/src/main/java/com/nvidia/spark/rapids/iceberg/IcebergShimUtils.java

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -20,6 +20,7 @@
2020
import com.nvidia.spark.rapids.NoopMetric$;
2121
import com.nvidia.spark.rapids.RapidsConf;
2222
import com.nvidia.spark.rapids.fileio.iceberg.IcebergInputFile;
23+
import com.nvidia.spark.rapids.jni.fileio.RapidsInputFile;
2324
import org.apache.hadoop.fs.Path;
2425
import org.apache.iceberg.ContentFile;
2526
import org.apache.iceberg.FileScanTask;
@@ -93,6 +94,7 @@ public interface IcebergShimUtils {
9394
*/
9495
default ParquetFileReader openParquetReader(
9596
IcebergInputFile inputFile,
97+
RapidsInputFile footerFile,
9698
Path filePath,
9799
ParquetReadOptions options,
98100
scala.collection.immutable.Map<String, GpuMetric> metrics) throws IOException {

iceberg/common/src/main/java/com/nvidia/spark/rapids/iceberg/ShimUtils.java

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -20,6 +20,7 @@
2020
import com.nvidia.spark.rapids.RapidsConf;
2121
import com.nvidia.spark.rapids.ShimLoader;
2222
import com.nvidia.spark.rapids.fileio.iceberg.IcebergInputFile;
23+
import com.nvidia.spark.rapids.jni.fileio.RapidsInputFile;
2324

2425
import org.apache.hadoop.fs.Path;
2526
import org.apache.iceberg.ContentFile;
@@ -69,10 +70,11 @@ public static Map<String, Map<String, String>> storageCredentialOverlays(FileIO
6970

7071
public static ParquetFileReader openParquetReader(
7172
IcebergInputFile inputFile,
73+
RapidsInputFile footerFile,
7274
Path filePath,
7375
ParquetReadOptions options,
7476
scala.collection.immutable.Map<String, GpuMetric> metrics) throws IOException {
75-
return IMPL.openParquetReader(inputFile, filePath, options, metrics);
77+
return IMPL.openParquetReader(inputFile, footerFile, filePath, options, metrics);
7678
}
7779

7880
public static GpuSparkScan newCopyOnWriteScan(

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

Lines changed: 8 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -127,12 +127,18 @@ public void readVectored(HostMemoryBuffer output, List<CopyRange> copyRanges)
127127
*/
128128
@Override
129129
public void readTail(long length, HostMemoryBuffer output) throws IOException {
130+
readTail(length, output, 0L);
131+
}
132+
133+
public long readTail(long length, HostMemoryBuffer output, long outputOffset)
134+
throws IOException {
130135
if (length == 0) {
131-
return;
136+
return 0;
132137
}
133138
if (length < 0) {
134139
throw new IllegalArgumentException("length must be non-negative");
135140
}
136-
IcebergS3RangeCopier.copyTailToHMB(icebergS3Client, output, s3Uri, length, /*dstOffset*/ 0L);
141+
return IcebergS3RangeCopier.copyTailToHMB(
142+
icebergS3Client, output, s3Uri, length, outputOffset);
137143
}
138144
}

iceberg/common/src/main/scala/com/nvidia/spark/rapids/iceberg/data/GpuDeleteFilter.scala

Lines changed: 0 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -415,4 +415,3 @@ private case class DeleteFilterContext(
415415
joinTime)
416416
}
417417
}
418-

iceberg/common/src/main/scala/com/nvidia/spark/rapids/iceberg/data/GpuDeleteLoader.scala

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -117,4 +117,4 @@ class DefaultDeleteLoader(
117117
case SingleFile => SingleFile
118118
}
119119
}
120-
}
120+
}

iceberg/common/src/main/scala/com/nvidia/spark/rapids/iceberg/parquet/reader.scala

Lines changed: 6 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -25,6 +25,7 @@ import scala.collection.JavaConverters._
2525
import com.nvidia.spark.rapids.{CombineConf, DateTimeRebaseCorrected, GpuMetric, ThreadPoolConfBuilder}
2626
import com.nvidia.spark.rapids.Arm.withResource
2727
import com.nvidia.spark.rapids.fileio.iceberg.IcebergInputFile
28+
import com.nvidia.spark.rapids.jni.fileio.RapidsInputFile
2829
import com.nvidia.spark.rapids.iceberg.parquet.converter.FromIcebergShaded._
2930
import com.nvidia.spark.rapids.parquet.{GpuParquetUtils, ParquetFileInfoWithBlockMeta}
3031
import com.nvidia.spark.rapids.shims.PartitionedFileUtilsShim
@@ -52,8 +53,10 @@ import org.apache.spark.sql.vectorized.ColumnarBatch
5253
case class IcebergPartitionedFile(
5354
file: IcebergInputFile,
5455
split: Option[(Long, Long)] = None,
55-
filter: Option[Expression] = None) {
56+
filter: Option[Expression] = None,
57+
footerFile: RapidsInputFile = null) {
5658

59+
private lazy val footerInputFile = Option(footerFile).getOrElse(file)
5760
lazy val urlEncodedPath: String = new Path(file.getDelegate.location()).toUri.toString
5861
lazy val path: Path = new Path(new URI(urlEncodedPath))
5962

@@ -63,7 +66,7 @@ case class IcebergPartitionedFile(
6366

6467
def newReader(metrics: Map[String, GpuMetric] = Map.empty): ParquetFileReader = {
6568
try {
66-
GpuParquetIO.openReader(file, path, parquetReadOptions, metrics)
69+
GpuParquetIO.openReader(file, footerInputFile, path, parquetReadOptions, metrics)
6770
} catch {
6871
case e: IOException =>
6972
throw new UncheckedIOException(s"Failed to newInputFile Parquet file: " +
@@ -304,4 +307,4 @@ object GpuIcebergParquetReader {
304307
}
305308
optionsBuilder.build
306309
}
307-
}
310+
}

iceberg/common/src/main/scala/org/apache/iceberg/parquet/GpuParquetIO.scala

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,7 @@
1717
package org.apache.iceberg.parquet
1818

1919
import com.nvidia.spark.rapids.GpuMetric
20+
import com.nvidia.spark.rapids.jni.fileio.RapidsInputFile
2021
import com.nvidia.spark.rapids.fileio.iceberg.IcebergInputFile
2122
import com.nvidia.spark.rapids.iceberg.ShimUtils
2223
import org.apache.hadoop.fs.Path
@@ -38,9 +39,10 @@ object GpuParquetIO {
3839
*/
3940
def openReader(
4041
inputFile: IcebergInputFile,
42+
footerFile: RapidsInputFile,
4143
filePath: Path,
4244
options: ParquetReadOptions,
4345
metrics: Map[String, GpuMetric]): ParquetFileReader = {
44-
ShimUtils.openParquetReader(inputFile, filePath, options, metrics)
46+
ShimUtils.openParquetReader(inputFile, footerFile, filePath, options, metrics)
4547
}
4648
}

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

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -26,6 +26,7 @@ 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

@@ -117,7 +118,8 @@ class GpuIcebergPartitionReader(private val task: GpuSparkInputPartition,
117118
val file = inputFiles(locationOf(t.file()))
118119
val icebergFile = IcebergPartitionedFile(file,
119120
Some((t.start(), t.length())),
120-
Some(t.residual()))
121+
Some(t.residual()),
122+
IcebergS3InputFile.maybeCreate(file.getDelegate, fileIO))
121123

122124
icebergFile -> t
123125
}))

iceberg/iceberg-1-10-x/src/main/java/com/nvidia/spark/rapids/iceberg/iceberg110x/ShimUtilsImpl.java

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,7 @@
1919
import com.nvidia.spark.rapids.GpuMetric;
2020
import com.nvidia.spark.rapids.RapidsConf;
2121
import com.nvidia.spark.rapids.fileio.iceberg.IcebergInputFile;
22+
import com.nvidia.spark.rapids.jni.fileio.RapidsInputFile;
2223
import com.nvidia.spark.rapids.iceberg.IcebergShimUtils;
2324
import org.apache.hadoop.fs.Path;
2425
import org.apache.iceberg.*;
@@ -73,10 +74,11 @@ public Map<String, Map<String, String>> storageCredentialOverlays(FileIO fileIO)
7374
@Override
7475
public ParquetFileReader openParquetReader(
7576
IcebergInputFile inputFile,
77+
RapidsInputFile footerFile,
7678
Path filePath,
7779
ParquetReadOptions options,
7880
scala.collection.immutable.Map<String, GpuMetric> metrics) throws IOException {
79-
return GpuParquetIOShim.openReader(inputFile, filePath, options, metrics);
81+
return GpuParquetIOShim.openReader(inputFile, footerFile, filePath, options, metrics);
8082
}
8183

8284
@Override

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

Lines changed: 14 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -17,11 +17,13 @@
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
23+
import com.nvidia.spark.rapids.jni.fileio.RapidsInputFile
2324
import com.nvidia.spark.rapids.parquet.{HMBInputFile, ParquetFooterUtils}
2425
import org.apache.hadoop.fs.Path
26+
import org.apache.iceberg.aws.s3.IcebergS3InputFile
2527
import org.apache.iceberg.parquet.GpuParquetIO
2628
import org.apache.iceberg.shaded.org.apache.parquet.ParquetReadOptions
2729
import org.apache.iceberg.shaded.org.apache.parquet.hadoop.ParquetFileReader
@@ -36,12 +38,22 @@ import org.apache.iceberg.shaded.org.apache.parquet.hadoop.ParquetFileReader
3638
object GpuParquetIOShim {
3739
def openReader(
3840
inputFile: IcebergInputFile,
41+
footerFile: RapidsInputFile,
3942
filePath: Path,
4043
options: ParquetReadOptions,
4144
metrics: Map[String, GpuMetric]): ParquetFileReader = {
4245
val metadata = withResource(ParquetFooterUtils.getFooterBuffer(
4346
inputFile, metrics,
44-
ParquetFooterUtils.readFooterBufferFromInputFile(inputFile, filePath))) { hmb =>
47+
footerFile match {
48+
case s3File: IcebergS3InputFile =>
49+
PerfIO.readParquetFooterBufferFromTail(
50+
filePath,
51+
(length, output, outputOffset) =>
52+
s3File.readTail(length, output, outputOffset),
53+
ParquetFooterUtils.verifyParquetMagic)
54+
case _ =>
55+
ParquetFooterUtils.readFooterBufferFromInputFile(footerFile, filePath)
56+
})) { hmb =>
4557
val shadedHmbFile = ToIcebergShaded.shade(new HMBInputFile(hmb))
4658
withResource(shadedHmbFile.newStream()) { hmbStream =>
4759
ParquetFileReader.readFooter(shadedHmbFile, options, hmbStream)

0 commit comments

Comments
 (0)