This is an automated email from the ASF dual-hosted git repository.

github-merge-queue[bot] pushed a commit to branch 
gh-readonly-queue/main/pr-6882-1529ae13cba2bd13d8c8ad6aa68cb04bddd7f545
in repository https://gitbox.apache.org/repos/asf/texera.git

commit 61cb36eab900b694bc00c45e5dd8eebb47266d82
Author: Xinyuan Lin <[email protected]>
AuthorDate: Thu Sep 24 03:42:29 2026 +0000

    fix(amber): close Iceberg reader streams leaked by iterator probes and 
bounded reads (#6882)
    
    ### What changes were proposed in this PR?
    
    **Root cause.** `IcebergDocument.getUsingFileSequenceOrder`'s iterator
    opened the Parquet reader (and the `S3InputStream` beneath it) inside
    `hasNext`. Two consumer shapes then leak the stream until the GC
    finalizer reclaims it — producing the `[S3InputStream] Unclosed input
    stream created by …` warnings all over the amber CI jobs:
    
    | Consumer shape | Why it leaked |
    |---|---|
    | Probes — `isEmpty` / `nonEmpty` / lone `hasNext` (e.g.
    `VirtualDocumentSpec`'s "clear the document" test) | `hasNext` opened a
    stream to answer, the caller abandoned the iterator; no close point ever
    ran |
    | Bounded reads — `getRange(from, until)` consumed to exactly the limit
    | The limit close ran only on a *subsequent* `hasNext` call that bounded
    consumers never make |
    
    **Fix — `hasNext` no longer acquires resources:**
    
    - `hasNext` claims the next data file from `FileScanTask.recordCount`
    metadata alone. This adds no new trust: the whole-file skip in
    `seekToUsableFile` already relies on the same field; `planFiles()` never
    splits files in Iceberg 1.9.2 (splitting is `planTasks()`-only) and
    amber's write paths are strictly append-only (no delete files anywhere
    in-repo), so the count is exact.
    - The Parquet reader opens lazily in `next()` (`openPendingFile()`), the
    only resource-acquisition point.
    - `next()` closes the reader deterministically the moment a bounded read
    has served its last record — before → after: `getRange(a, b).toList`
    used to keep the last file's stream open until GC; it now closes inside
    the final `next()`.
    - All closes funnel through an idempotent `closeCurrentReader()` that
    also resets the record iterator, so a closed reader is never polled (the
    old code did poll one on the exhaustion path — benign with today's
    Iceberg iterator internals, but the hazard is gone).
    
    Record sequences, `hasNext` semantics, `NoSuchElementException`
    behavior, and the incremental-snapshot refresh (`lastSnapshotId`
    bookkeeping) are unchanged — only *where* streams open and close moved.
    
    **Also:** close the `RESTCatalog` in
    `IcebergRestCatalogIntegrationSpec.afterAll`. Iceberg 1.9.2's
    `RESTSessionCatalog` tracks per-table `FileIO` instances
    (`FileIOTracker`) and closes them with the catalog, which removes the
    sibling `Unclosed S3FileIO instance` finalizer warnings.
    
    Residual (pre-existing, untouched) abandonment paths — notably
    `SyncExecutionResource.collectOperatorResult`'s visualization
    early-return — are inventoried in #6881 as follow-ups rather than
    expanded here.
    
    ### Any related issues, documentation, discussions?
    
    Closes #6881
    
    ### How was this PR tested?
    
    The changed paths are pinned by the existing `VirtualDocumentSpec`
    contract suite (probe, range, `getAfter`, incremental second-batch
    arrival, concurrent writes), which `IcebergDocumentSpec` runs against
    real Iceberg storage in the `amber-integration` CI job — those specs
    exercise every branch of the restructured iterator. Verified locally:
    `WorkflowCore/Test/compile`, `WorkflowExecutionService/Test/compile`,
    and `scalafmtCheck` on both source sets all pass.
    
    Additionally verified by an exhaustive old-vs-new state-machine review
    covering: probe-only use, repeated `hasNext`, multi-file drains, bounded
    reads with partial mid-file skip, `from` beyond EOF, empty ranges, empty
    tables, snapshot arrival mid-read, zero-record files, loop termination,
    and closed-reader polling — record sequences and `hasNext` booleans are
    identical in every scenario; the only differences are the intended
    open/close points. The partial-skip invariant (`from -
    numOfSkippedRecords` non-zero only for the first claimed file) follows
    from `seekToUsableFile`'s `dropWhile` guarantee and is documented in the
    code.
    
    No new unit test is added because the leak itself is only observable
    through GC-finalizer instrumentation; the observable regression signal
    is the disappearance of the `Unclosed input stream` warnings from the
    amber CI logs.
    
    ### Was this PR authored or co-authored using generative AI tooling?
    
    Generated-by: Claude Code (Opus 4.8 [1M context])
    
    ---------
    
    Co-authored-by: Yicong Huang 
<[email protected]>
---
 .../IcebergRestCatalogIntegrationSpec.scala        |  14 +++
 .../storage/result/iceberg/IcebergDocument.scala   | 126 +++++++++++++++------
 .../result/iceberg/IcebergDocumentSpec.scala       | 101 ++++++++++++++++-
 3 files changed, 208 insertions(+), 33 deletions(-)

diff --git 
a/amber/src/test/integration/org/apache/texera/amber/storage/iceberg/IcebergRestCatalogIntegrationSpec.scala
 
b/amber/src/test/integration/org/apache/texera/amber/storage/iceberg/IcebergRestCatalogIntegrationSpec.scala
index 807591dde5..e7896305b4 100644
--- 
a/amber/src/test/integration/org/apache/texera/amber/storage/iceberg/IcebergRestCatalogIntegrationSpec.scala
+++ 
b/amber/src/test/integration/org/apache/texera/amber/storage/iceberg/IcebergRestCatalogIntegrationSpec.scala
@@ -47,6 +47,20 @@ class IcebergRestCatalogIntegrationSpec extends AnyFlatSpec 
with BeforeAndAfterA
     )
   }
 
+  override def afterAll(): Unit = {
+    // RESTCatalog is Closeable: closing releases its HTTP client and the
+    // S3FileIO instances it created for table operations. Left unclosed,
+    // the finalizer reclaims them and logs "Unclosed S3FileIO instance"
+    // warnings with full stack traces after the spec finishes.
+    try {
+      if (restCatalog != null) {
+        restCatalog.close()
+      }
+    } finally {
+      super.afterAll()
+    }
+  }
+
   behavior of "Iceberg REST catalog"
 
   it should "round-trip table metadata via the REST catalog" in {
diff --git 
a/common/workflow-core/src/main/scala/org/apache/texera/amber/core/storage/result/iceberg/IcebergDocument.scala
 
b/common/workflow-core/src/main/scala/org/apache/texera/amber/core/storage/result/iceberg/IcebergDocument.scala
index 3f3f131ada..2de20c9f43 100644
--- 
a/common/workflow-core/src/main/scala/org/apache/texera/amber/core/storage/result/iceberg/IcebergDocument.scala
+++ 
b/common/workflow-core/src/main/scala/org/apache/texera/amber/core/storage/result/iceberg/IcebergDocument.scala
@@ -27,6 +27,7 @@ import org.apache.commons.io.IOUtils
 import org.apache.iceberg.catalog.{Catalog, TableIdentifier}
 import org.apache.iceberg.data.Record
 import org.apache.iceberg.exceptions.NoSuchTableException
+import org.apache.iceberg.io.CloseableIterator
 import org.apache.iceberg.types.{Conversions, Types}
 import org.apache.iceberg.{FileScanTask, Table}
 
@@ -155,6 +156,18 @@ private[storage] class IcebergDocument[T >: Null <: 
AnyRef](
     )
   }
 
+  /**
+    * Opens a Parquet reader over one data file. The iterators returned by the
+    * read methods call this only from `next()`, never from `hasNext`, and 
close
+    * each reader exactly once. Overridden in tests to count opens and closes.
+    */
+  protected def openDataFile(
+      task: FileScanTask,
+      schema: org.apache.iceberg.Schema,
+      table: Table
+  ): CloseableIterator[Record] =
+    IcebergUtil.readDataFileAsIterator(task.file(), schema, table)
+
   /**
     * Util iterator to get T in certain range
     *
@@ -195,6 +208,50 @@ private[storage] class IcebergDocument[T >: Null <: 
AnyRef](
         private var currentRecordIterator: Iterator[Record] = Iterator.empty
         private var currentRecordIteratorCloser: AutoCloseable = () => ()
 
+        // Next file to read, claimed from usableFileIterator by hasNext but
+        // not opened until next() actually needs a record. hasNext answers
+        // from the file's recordCount metadata alone, so probes such as
+        // isEmpty/nonEmpty (which abandon the iterator right after hasNext)
+        // never leave a Parquet reader / S3 stream open behind them.
+        private var pendingFile: Option[FileScanTask] = None
+
+        // Idempotent release of the active file's reader (and the S3 stream
+        // beneath it). Also resets the record iterator so a closed reader is
+        // never polled again.
+        private def closeCurrentReader(): Unit = {
+          val closer = currentRecordIteratorCloser
+          currentRecordIteratorCloser = () => ()
+          currentRecordIterator = Iterator.empty
+          closer.close()
+        }
+
+        // Open the file claimed by hasNext and point currentRecordIterator at
+        // its records, applying the one-off partial skip for the first file.
+        // Only next() calls this, keeping hasNext free of any Parquet/S3
+        // resource acquisition.
+        private def openPendingFile(): Unit = {
+          val task = pendingFile.getOrElse(
+            throw new IllegalStateException("no pending file to open")
+          )
+          pendingFile = None
+          val schemaToUse = columns match {
+            case Some(cols) => tableSchema.select(cols.asJava)
+            case None       => tableSchema
+          }
+          // Release the prior file's reader before opening the next.
+          closeCurrentReader()
+          val nextIter = openDataFile(task, 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
+          }
+        }
+
         // Util function to load the table's metadata
         private def loadTableMetadata(): Option[Table] = {
           IcebergUtil.loadTableMetadata(
@@ -276,8 +333,7 @@ private[storage] class IcebergDocument[T >: Null <: AnyRef](
         override def hasNext: Boolean = {
           if (numOfReturnedRecords >= totalRecordsToReturn) {
             // Caller-imposed limit reached; release the active file's reader.
-            currentRecordIteratorCloser.close()
-            currentRecordIteratorCloser = () => ()
+            closeCurrentReader()
             return false
           }
 
@@ -287,48 +343,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()
+          }
           val schemaToUse = columns match {
             case Some(cols) => tableSchema.select(cols.asJava)
             case None       => tableSchema
diff --git 
a/common/workflow-core/src/test/scala/org/apache/texera/amber/core/storage/result/iceberg/IcebergDocumentSpec.scala
 
b/common/workflow-core/src/test/scala/org/apache/texera/amber/core/storage/result/iceberg/IcebergDocumentSpec.scala
index fb2f5b4a15..dd4c2ec897 100644
--- 
a/common/workflow-core/src/test/scala/org/apache/texera/amber/core/storage/result/iceberg/IcebergDocumentSpec.scala
+++ 
b/common/workflow-core/src/test/scala/org/apache/texera/amber/core/storage/result/iceberg/IcebergDocumentSpec.scala
@@ -27,7 +27,8 @@ import org.apache.texera.amber.util.IcebergUtil
 import org.apache.iceberg.catalog.TableIdentifier
 import org.apache.iceberg.data.Record
 import org.apache.iceberg.exceptions.NoSuchTableException
-import org.apache.iceberg.{Schema => IcebergSchema}
+import org.apache.iceberg.io.CloseableIterator
+import org.apache.iceberg.{FileScanTask, Table, Schema => IcebergSchema}
 import org.scalatest.BeforeAndAfterAll
 import org.scalatest.flatspec.AnyFlatSpec
 import org.scalatest.matchers.should.Matchers
@@ -92,6 +93,40 @@ class IcebergDocumentSpec extends AnyFlatSpec with Matchers 
with BeforeAndAfterA
       .add("ts", AttributeType.TIMESTAMP, new Timestamp(1_600_000_000_000L + 
id))
       .build()
 
+  /**
+    * An [[IcebergDocument]] that counts the Parquet readers it opens and every
+    * `close()` call made on them, so a spec can pin reader lifetimes directly
+    * instead of inferring them from finalizer warnings.
+    */
+  private class ReaderCountingDocument(tableName: String)
+      extends IcebergDocument[Tuple](tableNamespace, tableName, icebergSchema, 
serde, deserde) {
+    var opens = 0
+    var closes = 0
+
+    override protected def openDataFile(
+        task: FileScanTask,
+        schema: IcebergSchema,
+        table: Table
+    ): CloseableIterator[Record] = {
+      val reader = super.openDataFile(task, schema, table)
+      opens += 1
+      new CloseableIterator[Record] {
+        override def hasNext: Boolean = reader.hasNext
+        override def next(): Record = reader.next()
+        override def close(): Unit = {
+          closes += 1
+          reader.close()
+        }
+      }
+    }
+  }
+
+  private def newCountingDocument(): ReaderCountingDocument = {
+    val tableName = freshTableName()
+    newDocument(tableName)
+    new ReaderCountingDocument(tableName)
+  }
+
   /** Write the given tuples through a single writer session (one committed 
file). */
   private def write(doc: IcebergDocument[Tuple], tuples: Seq[Tuple]): Unit = {
     val writer = doc.writer(UUID.randomUUID().toString)
@@ -208,6 +243,70 @@ class IcebergDocumentSpec extends AnyFlatSpec with 
Matchers with BeforeAndAfterA
     }
   }
 
+  it should "open no reader when an iterator is only probed with hasNext" in {
+    val doc = newCountingDocument()
+    write(doc, (0 until 5).map(tuple))
+
+    doc.get().hasNext shouldBe true
+    doc.getRange(1, 3).hasNext shouldBe true
+    doc.getAfter(2).hasNext shouldBe true
+    doc.get().nonEmpty shouldBe true
+    doc.get().isEmpty shouldBe false
+
+    doc.opens shouldBe 0
+    doc.closes shouldBe 0
+  }
+
+  it should "close the reader exactly once when the final getRange element is 
consumed" in {
+    val doc = newCountingDocument()
+    write(doc, (0 until 5).map(tuple))
+
+    val it = doc.getRange(0, 2)
+    it.next().getField[Int]("id") shouldBe 0
+    doc.opens shouldBe 1
+    doc.closes shouldBe 0
+
+    // The last in-range record releases the reader at once, with no further
+    // hasNext call, even though the file still holds unread records.
+    it.next().getField[Int]("id") shouldBe 1
+    doc.closes shouldBe 1
+
+    // Later calls see the limit and do not close the reader a second time.
+    it.hasNext shouldBe false
+    it.hasNext shouldBe false
+    doc.opens shouldBe 1
+    doc.closes shouldBe 1
+  }
+
+  it should "close the reader when a getRange result is collected with toList" 
in {
+    val doc = newCountingDocument()
+    write(doc, (0 until 10).map(tuple))
+
+    doc.getRange(2, 5).toList.map(_.getField[Int]("id")) shouldBe List(2, 3, 4)
+    doc.opens shouldBe 1
+    doc.closes shouldBe 1
+  }
+
+  it should "open and close one reader per file on a full read across files" 
in {
+    val doc = newCountingDocument()
+    write(doc, (0 until 3).map(tuple))
+    write(doc, (3 until 6).map(tuple))
+
+    doc.get().toList.map(_.getField[Int]("id")).toSet shouldBe (0 until 
6).toSet
+    doc.opens shouldBe 2
+    doc.closes shouldBe 2
+  }
+
+  it should "skip whole files by metadata without opening them" in {
+    val doc = newCountingDocument()
+    write(doc, (0 until 3).map(tuple))
+    write(doc, (3 until 6).map(tuple))
+
+    doc.getAfter(4).toList.map(_.getField[Int]("id")) shouldBe List(4, 5)
+    doc.opens shouldBe 1
+    doc.closes shouldBe 1
+  }
+
   it should "compute per-field statistics for numeric, string and timestamp 
columns" in {
     val doc = newDocument()
     write(doc, (1 to 5).map(tuple))

Reply via email to