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