Copilot commented on code in PR #6882:
URL: https://github.com/apache/texera/pull/6882#discussion_r4058804179
##########
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:
The resource-lifetime regression is not asserted by the existing tests:
their probe/range result checks also passed with the old leak. Please add
deterministic instrumentation (for example, an injectable/counting reader or
FileIO) that verifies a probe opens no reader and consuming the final
`getRange` element closes its reader exactly once, so this fix cannot silently
regress.
--
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]