aglinxinyuan commented on code in PR #6882:
URL: https://github.com/apache/texera/pull/6882#discussion_r4089367936


##########
common/workflow-core/src/main/scala/org/apache/texera/amber/core/storage/result/iceberg/IcebergDocument.scala:
##########
@@ -277,48 +323,54 @@ private[storage] class IcebergDocument[T >: Null <: 
AnyRef](
             return true
           }
 
-          if (!usableFileIterator.hasNext) {
-            usableFileIterator = seekToUsableFile()
-          }
+          // The active file (if any) is exhausted; release its reader before
+          // deciding from metadata whether more records exist.
+          closeCurrentReader()
 
-          while (!currentRecordIterator.hasNext && usableFileIterator.hasNext) 
{
-            val nextFile = usableFileIterator.next()
-            val schemaToUse = columns match {
-              case Some(cols) => tableSchema.select(cols.asJava)
-              case None       => tableSchema
+          if (pendingFile.isEmpty) {
+            if (!usableFileIterator.hasNext) {
+              usableFileIterator = seekToUsableFile()
             }
-            // Release the prior file's reader before opening the next.
-            currentRecordIteratorCloser.close()
-            val nextIter = IcebergUtil.readDataFileAsIterator(
-              nextFile.file(),
-              schemaToUse,
-              table.get
-            )
-            currentRecordIteratorCloser = nextIter
-            currentRecordIterator = nextIter.asScala
-
-            // Skip records within the file if necessary
-            val recordsToSkipInFile = from - numOfSkippedRecords
-            if (recordsToSkipInFile > 0) {
-              currentRecordIterator = 
currentRecordIterator.drop(recordsToSkipInFile)
-              numOfSkippedRecords += recordsToSkipInFile
+            // Claim the next file that still has usable records, judging by
+            // recordCount metadata alone (the same field the whole-file skip
+            // in seekToUsableFile already relies on). The partial skip
+            // `from - numOfSkippedRecords` can only be non-zero for the first
+            // claimed file: seekToUsableFile's dropWhile guarantees that file
+            // satisfies recordCount > partial skip, and openPendingFile zeroes
+            // the skip before any later file is judged here.
+            while (pendingFile.isEmpty && usableFileIterator.hasNext) {
+              val task = usableFileIterator.next()
+              if (task.file().recordCount() > from - numOfSkippedRecords) {
+                pendingFile = Some(task)
+              }
             }
           }
 
-          val hasMore = currentRecordIterator.hasNext
-          if (!hasMore) {
-            // All files exhausted; release the last file's reader.
-            currentRecordIteratorCloser.close()
-            currentRecordIteratorCloser = () => ()
-          }
-          hasMore
+          pendingFile.nonEmpty
         }
 
         override def next(): T = {
           if (!hasNext) throw new NoSuchElementException("No more records 
available")
 
+          // hasNext only claims files by their metadata; the Parquet reader is
+          // opened lazily here. The loop is defensive: should a claimed file
+          // yield nothing after the partial skip, hasNext claims the next one
+          // (or reports exhaustion).
+          while (!currentRecordIterator.hasNext) {
+            openPendingFile()
+            if (!currentRecordIterator.hasNext && !hasNext) {
+              throw new NoSuchElementException("No more records available")
+            }
+          }
+
           val record = currentRecordIterator.next()
           numOfReturnedRecords += 1
+          if (numOfReturnedRecords >= totalRecordsToReturn) {
+            // Bounded read fully served: release the reader now instead of
+            // waiting for a further hasNext call that bounded consumers
+            // (e.g. getRange) rarely make.
+            closeCurrentReader()

Review Comment:
   Done in 67e2ee2061. Reader creation now goes through an overridable 
`openDataFile`. `IcebergDocumentSpec` overrides it in a 
`ReaderCountingDocument` that counts opens and every `close()` call, and pins:
   
   | Case | opens | closes |
   |---|---|---|
   | `hasNext` / `nonEmpty` / `isEmpty` probes on `get`, `getRange`, `getAfter` 
| 0 | 0 |
   | `getRange(0, 2)`: closed right after the 2nd `next()`, no `hasNext` 
needed; later `hasNext` calls don't close again | 1 | 1 |
   | `getRange(2, 5).toList` | 1 | 1 |
   | full `get()` across 2 files | 2 | 2 |
   | `getAfter(4)` skips file 1 by metadata only | 1 | 1 |
   
   I checked that these tests catch the leak: putting back the eager open in 
`hasNext` and removing the limit close in `next()` fails both of the first two 
cases.
   
   This commit also picks up the two suppressed suggestions: 
`closeCurrentReader` resets its state before calling `close()`, and 
`IcebergRestCatalogIntegrationSpec.afterAll` calls `super.afterAll()` in a 
`finally`.
   



-- 
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]

Reply via email to