aglinxinyuan opened a new pull request, #6882: URL: https://github.com/apache/texera/pull/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]) -- 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]
