sunchao commented on code in PR #3726:
URL: https://github.com/apache/parquet-java/pull/3726#discussion_r3887224876


##########
parquet-hadoop/src/main/java/org/apache/parquet/hadoop/ParquetFileReader.java:
##########
@@ -2343,32 +2396,39 @@ private void setReadMetrics(long startNs, long len) {
     }
 
     /**
-     * Populate data in a parquet file range from a vectored range; will block 
for up
-     * to {@link #HADOOP_VECTORED_READ_TIMEOUT_SECONDS} seconds.
-     * @param currRange range to populated.
+     * Populate data in a parquet file range from one or more bounded vectored 
ranges; together
+     * they may block for up to {@link #HADOOP_VECTORED_READ_TIMEOUT_SECONDS} 
seconds.
+     * @param ranges bounded ranges containing this part.
      * @param builder used to build chunk list to read the pages for the 
different columns.
      * @throws IOException if there is an error while reading from the stream, 
including a timeout.
      */
-    public void readFromVectoredRange(ParquetFileRange currRange, 
ChunkListBuilder builder) throws IOException {
-      ByteBuffer buffer;
+    public void readFromVectoredRanges(List<ParquetFileRange> ranges, 
ChunkListBuilder builder) throws IOException {
+      List<ByteBuffer> buffers = new ArrayList<>(ranges.size());
+      ParquetFileRange currentRange = null;
       final long timeoutSeconds = HADOOP_VECTORED_READ_TIMEOUT_SECONDS;
+      final long timeoutNanos = TimeUnit.SECONDS.toNanos(timeoutSeconds);
       long readStart = System.nanoTime();
       try {
-        LOG.debug(
-            "Waiting for vectored read to finish for range {} with timeout {} 
seconds",
-            currRange,
-            timeoutSeconds);
-        buffer = FutureIO.awaitFuture(currRange.getDataReadFuture(), 
timeoutSeconds, TimeUnit.SECONDS);
-        setReadMetrics(readStart, currRange.getLength());
+        for (ParquetFileRange range : ranges) {
+          currentRange = range;
+          LOG.debug(
+              "Waiting for vectored read to finish for range {} with timeout 
{} seconds",
+              range,
+              timeoutSeconds);
+          long remainingNanos = Math.max(timeoutNanos - (System.nanoTime() - 
readStart), 0L);
+          buffers.add(FutureIO.awaitFuture(range.getDataReadFuture(), 
remainingNanos, TimeUnit.NANOSECONDS));

Review Comment:
   Fixed in 
[f82eaa42](https://github.com/apache/parquet-java/pull/3726/commits/f82eaa428528849cf019c99410fe879c4d641875).
 Vectored reads now track the original allocations and transfer them to the row 
group's releaser after successful completion. This includes extra checksum 
buffers and buffers whose returned results are slices, so the allocator 
receives the original objects it created. The row group is also closed on 
read/parsing failures.
   
   Added ownership tests for filtered and unfiltered reads, slices, checksum 
allocations, and parse failures. Real Hadoop 3.4.1 raw/checksum filesystem 
probes decoded all 128 rows with balanced allocations/releases and a clean 
tracking allocator.



##########
parquet-hadoop/src/main/java/org/apache/parquet/hadoop/ParquetFileReader.java:
##########
@@ -1357,32 +1334,105 @@ private boolean 
arePartsValidForVectoredIo(List<ConsecutivePartList> allParts) {
    * If directly implemented by a Filesystem then it is likely to be a more 
efficient
    * operation such as a scatter-gather read (native IO) or set of parallel
    * GET requests against an object store.
+   * The allocation limit applies to filesystem buffers; decoders can still 
require a
+   * contiguous buffer for an individual logical value larger than that limit.
    * @param allParts all parts to be read.
    * @param builder used to build chunk list to read the pages for the 
different columns.
-   * @throws IOException any IOE.
-   * @throws IllegalArgumentException arguments are invalid.
-   * @throws UnsupportedOperationException if the filesystem does not support 
vectored IO.
+   * @throws IOException if submitting or consuming the vectored reads fails.
+   * @throws IllegalArgumentException if range preparation fails before any 
reads are submitted.
    */
   private void readVectored(List<ConsecutivePartList> allParts, 
ChunkListBuilder builder) throws IOException {
-
+    final int maximumAllocation = options.getMaxAllocationSize();
+    Preconditions.checkArgument(maximumAllocation > 0, "Invalid maximum 
allocation size %s", maximumAllocation);
+    if (vectoredReadFileLength < 0) {
+      vectoredReadFileLength = file.getLength();
+    }
+    final long fileLength = vectoredReadFileLength;
     List<ParquetFileRange> ranges = new ArrayList<>(allParts.size());
+    List<Integer> partRangeCounts = new ArrayList<>(allParts.size());
     long totalSize = 0;
     for (ConsecutivePartList consecutiveChunks : allParts) {
       final long len = consecutiveChunks.length;
-      Preconditions.checkArgument(
-          len < Integer.MAX_VALUE,
-          "Invalid length %s for vectored read operation. It must be less than 
max integer value.",
-          len);
-      ranges.add(new ParquetFileRange(consecutiveChunks.offset, (int) len));
+      final long start = consecutiveChunks.offset;
+      if (start < 0 || len < 0 || start > fileLength || len > fileLength - 
start) {
+        throw new IOException(String.format(
+            "Invalid vectored read range (offset %d, length %d) for file 
length %d",
+            start, len, fileLength));
+      }
+      final int firstRange = ranges.size();
+      long remaining = len;
+      long offset = start;
+      do {
+        int rangeLength = (int) Math.min(remaining, maximumAllocation);
+        ranges.add(new ParquetFileRange(offset, rangeLength));
+        offset += rangeLength;
+        remaining -= rangeLength;
+      } while (remaining > 0);
+      partRangeCounts.add(ranges.size() - firstRange);
       totalSize += len;
     }
     LOG.debug("Reading {} bytes of data with vectored IO in {} ranges", 
totalSize, ranges.size());
-    // Request a vectored read;
-    f.readVectored(ranges, options.getAllocator());
-    int k = 0;
-    for (ConsecutivePartList consecutivePart : allParts) {
-      ParquetFileRange currRange = ranges.get(k++);
-      consecutivePart.readFromVectoredRange(currRange, builder);
+    final long readStart = System.nanoTime();
+    try {
+      // Even a synchronous rejection can follow partial submission. The 
Hadoop bridge
+      // publishes futures only after submission returns, so missing futures 
do not prove
+      // that no reads started. Once this call is entered, normal-read 
fallback is unsafe.
+      f.readVectored(ranges, options.getAllocator());
+      int firstRange = 0;
+      for (int partIndex = 0; partIndex < allParts.size(); partIndex++) {
+        int endRange = firstRange + partRangeCounts.get(partIndex);
+        
allParts.get(partIndex).readFromVectoredRanges(ranges.subList(firstRange, 
endRange), builder);
+        firstRange = endRange;
+      }
+    } catch (IllegalArgumentException | UnsupportedOperationException e) {
+      // Consumption may also have populated the builder. Do not replay those 
chunks.
+      IOException failure =
+          new IOException("Vectored read failed after asynchronous reads may 
have been submitted", e);
+      awaitRemainingVectoredReads(ranges, readStart, failure);

Review Comment:
   Addressed in 
[f82eaa42](https://github.com/apache/parquet-java/pull/3726/commits/f82eaa428528849cf019c99410fe879c4d641875).
 The bridge now publishes available futures in a `finally` block, without 
masking the original submission failure. A failed submission still fails the 
read without replaying ordinary I/O.
   
   Stream closure is queued behind the actual submission call. Buffer release 
requires every range to have a published, completed, non-cancelled future. 
Hadoop can publish futures for work it never schedules, so cleanup does not 
block indefinitely waiting for those; the batch retains its remaining 
allocations if completion cannot be established.
   
   Added regressions for partial submission and a future-publication failure 
that must not replace the submission exception.



##########
parquet-hadoop/src/main/java/org/apache/parquet/hadoop/ParquetFileReader.java:
##########
@@ -1357,32 +1334,105 @@ private boolean 
arePartsValidForVectoredIo(List<ConsecutivePartList> allParts) {
    * If directly implemented by a Filesystem then it is likely to be a more 
efficient
    * operation such as a scatter-gather read (native IO) or set of parallel
    * GET requests against an object store.
+   * The allocation limit applies to filesystem buffers; decoders can still 
require a
+   * contiguous buffer for an individual logical value larger than that limit.
    * @param allParts all parts to be read.
    * @param builder used to build chunk list to read the pages for the 
different columns.
-   * @throws IOException any IOE.
-   * @throws IllegalArgumentException arguments are invalid.
-   * @throws UnsupportedOperationException if the filesystem does not support 
vectored IO.
+   * @throws IOException if submitting or consuming the vectored reads fails.
+   * @throws IllegalArgumentException if range preparation fails before any 
reads are submitted.
    */
   private void readVectored(List<ConsecutivePartList> allParts, 
ChunkListBuilder builder) throws IOException {
-
+    final int maximumAllocation = options.getMaxAllocationSize();
+    Preconditions.checkArgument(maximumAllocation > 0, "Invalid maximum 
allocation size %s", maximumAllocation);
+    if (vectoredReadFileLength < 0) {
+      vectoredReadFileLength = file.getLength();
+    }
+    final long fileLength = vectoredReadFileLength;
     List<ParquetFileRange> ranges = new ArrayList<>(allParts.size());
+    List<Integer> partRangeCounts = new ArrayList<>(allParts.size());
     long totalSize = 0;
     for (ConsecutivePartList consecutiveChunks : allParts) {
       final long len = consecutiveChunks.length;
-      Preconditions.checkArgument(
-          len < Integer.MAX_VALUE,
-          "Invalid length %s for vectored read operation. It must be less than 
max integer value.",
-          len);
-      ranges.add(new ParquetFileRange(consecutiveChunks.offset, (int) len));
+      final long start = consecutiveChunks.offset;
+      if (start < 0 || len < 0 || start > fileLength || len > fileLength - 
start) {
+        throw new IOException(String.format(
+            "Invalid vectored read range (offset %d, length %d) for file 
length %d",
+            start, len, fileLength));
+      }
+      final int firstRange = ranges.size();
+      long remaining = len;
+      long offset = start;
+      do {
+        int rangeLength = (int) Math.min(remaining, maximumAllocation);
+        ranges.add(new ParquetFileRange(offset, rangeLength));
+        offset += rangeLength;
+        remaining -= rangeLength;
+      } while (remaining > 0);
+      partRangeCounts.add(ranges.size() - firstRange);
       totalSize += len;
     }
     LOG.debug("Reading {} bytes of data with vectored IO in {} ranges", 
totalSize, ranges.size());
-    // Request a vectored read;
-    f.readVectored(ranges, options.getAllocator());
-    int k = 0;
-    for (ConsecutivePartList consecutivePart : allParts) {
-      ParquetFileRange currRange = ranges.get(k++);
-      consecutivePart.readFromVectoredRange(currRange, builder);
+    final long readStart = System.nanoTime();

Review Comment:
   Fixed in 
[f82eaa42](https://github.com/apache/parquet-java/pull/3726/commits/f82eaa428528849cf019c99410fe879c4d641875).
 Submission now runs on a reader-owned worker, and one 300-second budget per 
vectored-read batch starts before submission and is shared with every range 
wait. It is no longer reset for each part.
   
   On timeout, the reader is invalidated and cleanup is queued on the same 
worker, so cancelling the submission `Future` cannot let cleanup overtake a 
backend call that is still running. The caller can return even if the backend 
ignores interruption; its resources remain owned by deferred cleanup, which may 
itself be delayed indefinitely.
   
   Added tests for blocked submission, the shared deadline, and submission that 
continues after cancellation. These are local regression tests; I have not 
validated this against live S3.



##########
parquet-hadoop/src/main/java/org/apache/parquet/hadoop/ParquetFileReader.java:
##########
@@ -1357,32 +1334,105 @@ private boolean 
arePartsValidForVectoredIo(List<ConsecutivePartList> allParts) {
    * If directly implemented by a Filesystem then it is likely to be a more 
efficient
    * operation such as a scatter-gather read (native IO) or set of parallel
    * GET requests against an object store.
+   * The allocation limit applies to filesystem buffers; decoders can still 
require a
+   * contiguous buffer for an individual logical value larger than that limit.
    * @param allParts all parts to be read.
    * @param builder used to build chunk list to read the pages for the 
different columns.
-   * @throws IOException any IOE.
-   * @throws IllegalArgumentException arguments are invalid.
-   * @throws UnsupportedOperationException if the filesystem does not support 
vectored IO.
+   * @throws IOException if submitting or consuming the vectored reads fails.
+   * @throws IllegalArgumentException if range preparation fails before any 
reads are submitted.
    */
   private void readVectored(List<ConsecutivePartList> allParts, 
ChunkListBuilder builder) throws IOException {
-
+    final int maximumAllocation = options.getMaxAllocationSize();
+    Preconditions.checkArgument(maximumAllocation > 0, "Invalid maximum 
allocation size %s", maximumAllocation);
+    if (vectoredReadFileLength < 0) {
+      vectoredReadFileLength = file.getLength();

Review Comment:
   Fixed in 
[f82eaa42](https://github.com/apache/parquet-java/pull/3726/commits/f82eaa428528849cf019c99410fe879c4d641875).
 The length lookup now has a narrow pre-submission fallback: a metadata 
`IOException` disables vectored I/O for that reader and continues with ordinary 
reads on the already-open stream. This avoids repeated failed lookups on later 
row groups. Interruption still propagates, and failures after submission still 
cannot trigger ordinary-read replay.
   
   Added regression tests with and without a supplied footer, plus an 
interrupted-lookup case. Both fallback regressions fail on the original commit 
and pass with this change.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to