Skip to content

Commit 2cc32d2

Browse files
committed
Address comments
1 parent 2d9ec08 commit 2cc32d2

11 files changed

Lines changed: 19 additions & 49 deletions

File tree

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

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -65,7 +65,7 @@ public IcebergFileIO(FileIO delegate, Map<String, GpuMetric> metrics) {
6565

6666

6767
@Override
68-
public RapidsInputFile newInputFile(String path) throws IOException {
68+
public IcebergInputFile newInputFile(String path) throws IOException {
6969
return newInputFile(delegate.newInputFile(path));
7070
}
7171

@@ -74,7 +74,7 @@ public RapidsInputFile newInputFile(String path) throws IOException {
7474
*
7575
* @param inputFile the Iceberg input file to wrap
7676
*/
77-
public RapidsInputFile newInputFile(InputFile inputFile) {
77+
public IcebergInputFile newInputFile(InputFile inputFile) {
7878
Objects.requireNonNull(inputFile, "inputFile can't be null");
7979
return IcebergS3InputFile.maybeCreate(inputFile, delegate, metrics);
8080
}

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

Lines changed: 0 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -20,7 +20,6 @@
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;
2423
import org.apache.hadoop.fs.Path;
2524
import org.apache.iceberg.ContentFile;
2625
import org.apache.iceberg.FileScanTask;
@@ -94,7 +93,6 @@ public interface IcebergShimUtils {
9493
*/
9594
default ParquetFileReader openParquetReader(
9695
IcebergInputFile inputFile,
97-
RapidsInputFile footerInputFile,
9896
Path filePath,
9997
ParquetReadOptions options,
10098
scala.collection.immutable.Map<String, GpuMetric> metrics) throws IOException {

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

Lines changed: 1 addition & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -20,7 +20,6 @@
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;
2423

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

7170
public static ParquetFileReader openParquetReader(
7271
IcebergInputFile inputFile,
73-
RapidsInputFile footerInputFile,
7472
Path filePath,
7573
ParquetReadOptions options,
7674
scala.collection.immutable.Map<String, GpuMetric> metrics) throws IOException {
77-
return IMPL.openParquetReader(inputFile, footerInputFile, filePath, options, metrics);
75+
return IMPL.openParquetReader(inputFile, filePath, options, metrics);
7876
}
7977

8078
public static GpuSparkScan newCopyOnWriteScan(

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

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -47,7 +47,7 @@
4747
*
4848
* <p>The package-private S3 file access is isolated in {@link IcebergS3InputFileAccess}.
4949
*/
50-
public final class IcebergS3InputFile implements RapidsInputFile {
50+
public final class IcebergS3InputFile extends IcebergInputFile {
5151
private static final Logger LOG = LoggerFactory.getLogger(IcebergS3InputFile.class);
5252

5353
private final IcebergInputFile delegate;
@@ -62,14 +62,15 @@ private IcebergS3InputFile(
6262
IcebergS3Client icebergS3Client,
6363
GpuMetric tailReadCount,
6464
GpuMetric tailReadBytes) {
65+
super(delegate.getDelegate());
6566
this.delegate = delegate;
6667
this.s3Uri = s3Uri;
6768
this.icebergS3Client = icebergS3Client;
6869
this.tailReadCount = tailReadCount;
6970
this.tailReadBytes = tailReadBytes;
7071
}
7172

72-
public static RapidsInputFile maybeCreate(
73+
public static IcebergInputFile maybeCreate(
7374
InputFile inputFile, FileIO fileIO, Map<String, GpuMetric> metrics) {
7475
// When the gating conf is off (or the file is not an S3 file), return the
7576
// default IcebergInputFile so the standard Iceberg SeekableInputStream path is used.

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

Lines changed: 4 additions & 19 deletions
Original file line numberDiff line numberDiff line change
@@ -24,9 +24,8 @@ import scala.collection.JavaConverters._
2424

2525
import com.nvidia.spark.rapids.{CombineConf, DateTimeRebaseCorrected, GpuMetric, ThreadPoolConfBuilder}
2626
import com.nvidia.spark.rapids.Arm.withResource
27-
import com.nvidia.spark.rapids.fileio.iceberg.{IcebergFileIO, IcebergInputFile}
27+
import com.nvidia.spark.rapids.fileio.iceberg.IcebergInputFile
2828
import com.nvidia.spark.rapids.iceberg.parquet.converter.FromIcebergShaded._
29-
import com.nvidia.spark.rapids.jni.fileio.RapidsInputFile
3029
import com.nvidia.spark.rapids.parquet.{GpuParquetUtils, ParquetFileInfoWithBlockMeta}
3130
import com.nvidia.spark.rapids.shims.PartitionedFileUtilsShim
3231
import org.apache.hadoop.conf.Configuration
@@ -63,20 +62,8 @@ case class IcebergPartitionedFile(
6362
}
6463

6564
def newReader(metrics: Map[String, GpuMetric] = Map.empty): ParquetFileReader = {
66-
newReader(file, metrics)
67-
}
68-
69-
def newReader(
70-
fileIO: IcebergFileIO,
71-
metrics: Map[String, GpuMetric]): ParquetFileReader = {
72-
newReader(fileIO.newInputFile(file.getDelegate), metrics)
73-
}
74-
75-
private def newReader(
76-
footerInputFile: RapidsInputFile,
77-
metrics: Map[String, GpuMetric]): ParquetFileReader = {
7865
try {
79-
GpuParquetIO.openReader(file, footerInputFile, path, parquetReadOptions, metrics)
66+
GpuParquetIO.openReader(file, path, parquetReadOptions, metrics)
8067
} catch {
8168
case e: IOException =>
8269
throw new UncheckedIOException(s"Failed to newInputFile Parquet file: " +
@@ -159,7 +146,6 @@ case class GpuIcebergParquetReaderConf(
159146
)
160147

161148
trait GpuIcebergParquetReader extends Iterator[ColumnarBatch] with AutoCloseable with Logging {
162-
def rapidsFileIO: IcebergFileIO
163149
def conf: GpuIcebergParquetReaderConf
164150

165151
def projectSchema(fileSchema: ShadedMessageType, requiredSchema: Schema):
@@ -217,7 +203,7 @@ trait GpuIcebergParquetReader extends Iterator[ColumnarBatch] with AutoCloseable
217203

218204
def filterParquetBlocks(file: IcebergPartitionedFile,
219205
requiredSchema: Schema): (ParquetFileInfoWithBlockMeta, ShadedMessageType) = {
220-
withResource(file.newReader(rapidsFileIO, conf.metrics)) { reader =>
206+
withResource(file.newReader(conf.metrics)) { reader =>
221207
val fileSchema = reader.getFileMetaData.getSchema
222208
val (typeWithIds, fileReadSchema) = projectSchema(fileSchema, requiredSchema)
223209
val filteredBlocks = filterRowGroups(reader, requiredSchema, typeWithIds, file.filter)
@@ -244,8 +230,7 @@ trait GpuIcebergParquetReader extends Iterator[ColumnarBatch] with AutoCloseable
244230
}
245231
}
246232
if (file.split.isDefined) {
247-
withResource(
248-
file.copy(split = None).newReader(rapidsFileIO, conf.metrics)) { fullReader =>
233+
withResource(file.copy(split = None).newReader(conf.metrics)) { fullReader =>
249234
populateFromAllBlocks(fullReader.getFooter.getBlocks)
250235
}
251236
} else {

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

Lines changed: 1 addition & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -19,7 +19,6 @@ package org.apache.iceberg.parquet
1919
import com.nvidia.spark.rapids.GpuMetric
2020
import com.nvidia.spark.rapids.fileio.iceberg.IcebergInputFile
2121
import com.nvidia.spark.rapids.iceberg.ShimUtils
22-
import com.nvidia.spark.rapids.jni.fileio.RapidsInputFile
2322
import org.apache.hadoop.fs.Path
2423
import org.apache.iceberg.io.InputFile
2524
import org.apache.iceberg.shaded.org.apache.parquet.ParquetReadOptions
@@ -39,10 +38,9 @@ object GpuParquetIO {
3938
*/
4039
def openReader(
4140
inputFile: IcebergInputFile,
42-
footerInputFile: RapidsInputFile,
4341
filePath: Path,
4442
options: ParquetReadOptions,
4543
metrics: Map[String, GpuMetric]): ParquetFileReader = {
46-
ShimUtils.openParquetReader(inputFile, footerInputFile, filePath, options, metrics)
44+
ShimUtils.openParquetReader(inputFile, filePath, options, metrics)
4745
}
4846
}

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

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -20,7 +20,7 @@ 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
@@ -110,7 +110,7 @@ class GpuIcebergPartitionReader(private val task: GpuSparkInputPartition,
110110
val inputFiles = table.encryption()
111111
.decrypt(encryptedFiles.asJava)
112112
.asScala
113-
.map(f => f.location() -> new IcebergInputFile(f))
113+
.map(f => f.location() -> rapidsFileIO.newInputFile(f))
114114
.toMap
115115

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

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

Lines changed: 1 addition & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -20,7 +20,6 @@
2020
import com.nvidia.spark.rapids.RapidsConf;
2121
import com.nvidia.spark.rapids.fileio.iceberg.IcebergInputFile;
2222
import com.nvidia.spark.rapids.iceberg.IcebergShimUtils;
23-
import com.nvidia.spark.rapids.jni.fileio.RapidsInputFile;
2423
import org.apache.hadoop.fs.Path;
2524
import org.apache.iceberg.*;
2625
import org.apache.iceberg.io.FileIO;
@@ -74,12 +73,10 @@ public Map<String, Map<String, String>> storageCredentialOverlays(FileIO fileIO)
7473
@Override
7574
public ParquetFileReader openParquetReader(
7675
IcebergInputFile inputFile,
77-
RapidsInputFile footerInputFile,
7876
Path filePath,
7977
ParquetReadOptions options,
8078
scala.collection.immutable.Map<String, GpuMetric> metrics) throws IOException {
81-
return GpuParquetIOShim.openReader(
82-
inputFile, footerInputFile, filePath, options, metrics);
79+
return GpuParquetIOShim.openReader(inputFile, filePath, options, metrics);
8380
}
8481

8582
@Override

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

Lines changed: 2 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -20,7 +20,6 @@ import com.nvidia.spark.rapids.Arm.{closeOnExcept, withResource}
2020
import com.nvidia.spark.rapids.GpuMetric
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
2423
import com.nvidia.spark.rapids.parquet.{HMBInputFile, ParquetFooterUtils}
2524
import org.apache.hadoop.fs.Path
2625
import org.apache.iceberg.parquet.GpuParquetIO
@@ -37,13 +36,12 @@ import org.apache.iceberg.shaded.org.apache.parquet.hadoop.ParquetFileReader
3736
object GpuParquetIOShim {
3837
def openReader(
3938
inputFile: IcebergInputFile,
40-
footerInputFile: RapidsInputFile,
4139
filePath: Path,
4240
options: ParquetReadOptions,
4341
metrics: Map[String, GpuMetric]): ParquetFileReader = {
4442
val metadata = withResource(ParquetFooterUtils.getFooterBuffer(
45-
footerInputFile, metrics,
46-
ParquetFooterUtils.readFooterBufferFromInputFile(footerInputFile, filePath))) { hmb =>
43+
inputFile, metrics,
44+
ParquetFooterUtils.readFooterBufferFromInputFile(inputFile, filePath))) { hmb =>
4745
val shadedHmbFile = ToIcebergShaded.shade(new HMBInputFile(hmb))
4846
withResource(shadedHmbFile.newStream()) { hmbStream =>
4947
ParquetFileReader.readFooter(shadedHmbFile, options, hmbStream)

iceberg/iceberg-1-11-x/src/main/java/com/nvidia/spark/rapids/iceberg/iceberg111x/ShimUtilsImpl.java

Lines changed: 1 addition & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -20,7 +20,6 @@
2020
import com.nvidia.spark.rapids.RapidsConf;
2121
import com.nvidia.spark.rapids.fileio.iceberg.IcebergInputFile;
2222
import com.nvidia.spark.rapids.iceberg.IcebergShimUtils;
23-
import com.nvidia.spark.rapids.jni.fileio.RapidsInputFile;
2423
import org.apache.hadoop.fs.Path;
2524
import org.apache.iceberg.*;
2625
import org.apache.iceberg.io.FileIO;
@@ -74,12 +73,10 @@ public Map<String, Map<String, String>> storageCredentialOverlays(FileIO fileIO)
7473
@Override
7574
public ParquetFileReader openParquetReader(
7675
IcebergInputFile inputFile,
77-
RapidsInputFile footerInputFile,
7876
Path filePath,
7977
ParquetReadOptions options,
8078
scala.collection.immutable.Map<String, GpuMetric> metrics) throws IOException {
81-
return GpuParquetIOShim.openReader(
82-
inputFile, footerInputFile, filePath, options, metrics);
79+
return GpuParquetIOShim.openReader(inputFile, filePath, options, metrics);
8380
}
8481

8582
@Override

0 commit comments

Comments
 (0)