fix(amber): bound sync result reads so visualization fetches release their reader - #6883
fix(amber): bound sync result reads so visualization fetches release their reader#6883aglinxinyuan wants to merge 3 commits into
Conversation
…their reader SyncExecutionResource.collectOperatorResult read operator results through an unbounded document.get() iterator; the visualization and oversized-first-tuple branches returned after a single next(), and the catch-all swallowed exceptions mid-drain, leaving the Iceberg Parquet reader / S3 input stream open until GC finalization. Size each read to what the method actually consumes: getRange(0, 1) for the first tuple, getRange(1, totalCount) for the truncation loops (which always drain), and drain the bounded remainder in the catch path. Bounded getRange iterators release their reader inside the next() that serves their last record (apache#6882), so no close API is needed.
Automated Reviewer SuggestionsBased on the
|
There was a problem hiding this comment.
Pull request overview
This PR updates Amber’s synchronous result fetch path (SyncExecutionResource.collectOperatorResult) to avoid leaking Iceberg/Parquet/S3 readers when the method returns early (visualization payload or oversized-first-tuple) or when an exception interrupts iteration, by switching from an unbounded document.get() iterator to bounded document.getRange(from, until) reads and adding a catch-path drain.
Changes:
- Replace the unbounded
document.get()read with boundedgetRange(0, 1)for the first tuple andgetRange(1, totalCount)for the remainder. - Hoist the remainder iterator so the catch path can drain it and trigger reader release on failures.
- Cap truncation-loop reads to
totalCountfor response internal consistency.
Comments suppressed due to low confidence (2)
amber/src/main/scala/org/apache/texera/web/resource/SyncExecutionResource.scala:572
- The visualization early-return consumes exactly one element from
getRange(0, 1)and returns without any further iterator interaction. On the current Iceberg iterator implementation, the limit-close is triggered by a subsequenthasNextcall, so this branch can still leak resources unless you explicitly trigger that close point before returning.
// the frontend renders it as an iframe rather than a table.
val firstTuple = firstTupleIterator.next()
if (totalCount == 1 && isVisualizationTuple(firstTuple)) {
val jsonResults =
ExecutionResultService.convertTuplesToJson(List(firstTuple), isVisualization = true)
amber/src/main/scala/org/apache/texera/web/resource/SyncExecutionResource.scala:605
- The oversized-first-tuple early-return has the same issue as the visualization return: it consumes one record from
getRange(0, 1)and returns without hitting the iterator’s limit-close point. Drain/probe the bounded iterator once before returning so resources are released deterministically.
// Records 1 until totalCount; the truncation loops below drain this
// iterator fully, releasing its reader at the bound.
tupleIterator = document.getRange(1, totalCount)
💡 Add Copilot custom instructions for smarter, more guided reviews. Learn how to get started.
|
| config | throughput | MB/s | latency | max Δ latest / 7d | |
|---|---|---|---|---|---|
| 🔴 | bs=10 sw=10 sl=64 | 463 | 0.283 | 20,066/34,011/34,011 us | 🔴 +9.9% / 🔴 +115.2% |
| 🔴 | bs=100 sw=10 sl=64 | 945 | 0.577 | 103,189/140,570/140,570 us | 🔴 +8.2% / 🔴 +31.0% |
| 🟢 | bs=1000 sw=10 sl=64 | 1,136 | 0.693 | 882,872/902,396/902,396 us | 🟢 -9.4% / 🟢 -14.5% |
Baseline details
Latest main 06f8e42 from same runner
| config | metric | PR | latest main | 7d avg | Δ latest | Δ 7d |
|---|---|---|---|---|---|---|
| bs=10 sw=10 sl=64 | throughput | 463 tuples/sec | 496 tuples/sec | 787.55 tuples/sec | -6.7% | -41.2% |
| bs=10 sw=10 sl=64 | MB/s | 0.283 MB/s | 0.303 MB/s | 0.481 MB/s | -6.6% | -41.1% |
| bs=10 sw=10 sl=64 | p50 | 20,066 us | 18,323 us | 12,255 us | +9.5% | +63.7% |
| bs=10 sw=10 sl=64 | p95 | 34,011 us | 30,934 us | 15,802 us | +9.9% | +115.2% |
| bs=10 sw=10 sl=64 | p99 | 34,011 us | 30,934 us | 19,008 us | +9.9% | +78.9% |
| bs=100 sw=10 sl=64 | throughput | 945 tuples/sec | 981 tuples/sec | 997.81 tuples/sec | -3.7% | -5.3% |
| bs=100 sw=10 sl=64 | MB/s | 0.577 MB/s | 0.599 MB/s | 0.609 MB/s | -3.7% | -5.3% |
| bs=100 sw=10 sl=64 | p50 | 103,189 us | 99,264 us | 100,690 us | +4.0% | +2.5% |
| bs=100 sw=10 sl=64 | p95 | 140,570 us | 129,971 us | 107,316 us | +8.2% | +31.0% |
| bs=100 sw=10 sl=64 | p99 | 140,570 us | 129,971 us | 113,823 us | +8.2% | +23.5% |
| bs=1000 sw=10 sl=64 | throughput | 1,136 tuples/sec | 1,124 tuples/sec | 1,030 tuples/sec | +1.1% | +10.3% |
| bs=1000 sw=10 sl=64 | MB/s | 0.693 MB/s | 0.686 MB/s | 0.629 MB/s | +1.0% | +10.2% |
| bs=1000 sw=10 sl=64 | p50 | 882,872 us | 884,760 us | 981,213 us | -0.2% | -10.0% |
| bs=1000 sw=10 sl=64 | p95 | 902,396 us | 995,756 us | 1,027,605 us | -9.4% | -12.2% |
| bs=1000 sw=10 sl=64 | p99 | 902,396 us | 995,756 us | 1,055,466 us | -9.4% | -14.5% |
Raw CSV
config_idx,batch_size,schema_width,string_len,num_batches,total_ms,total_tuples,total_bytes,tuples_per_sec,mb_per_sec,lat_p50_us,lat_p95_us,lat_p99_us
0,10,10,64,20,431.94,200,128000,463,0.283,20066.06,34011.19,34011.19
1,100,10,64,20,2116.85,2000,1280000,945,0.577,103188.81,140569.91,140569.91
2,1000,10,64,20,17608.69,20000,12800000,1136,0.693,882872.23,902395.98,902395.98
Codecov Report❌ Patch coverage is
❌ Your patch status has failed because the patch coverage (0.00%) is below the target coverage (60.00%). You can increase the patch coverage or adjust the target coverage. Additional details and impacted files@@ Coverage Diff @@
## main #6883 +/- ##
============================================
+ Coverage 78.98% 78.99% +0.01%
+ Complexity 3785 3783 -2
============================================
Files 1160 1160
Lines 46105 46112 +7
Branches 5115 5118 +3
============================================
+ Hits 36414 36425 +11
+ Misses 8069 8057 -12
- Partials 1622 1630 +8
*This pull request uses carry forward flags. Click here to find out more. ☔ View full report in Codecov by Harness. 🚀 New features to boost your workflow:
|
tupleIterator was hoisted so the catch block could drain a partially-consumed bounded read, but it stayed empty until the second getRange. A throw while probing, reading or converting the first tuple therefore skipped the drain and leaked that read's underlying reader. Assign it immediately after creation.
What changes were proposed in this PR?
SyncExecutionResource.collectOperatorResultread operator results through an unboundeddocument.get()iterator. Three paths abandoned that iterator with its Iceberg Parquet reader / S3 input stream still open, leaving the GC finalizer to reclaim it ([S3InputStream] Unclosed input stream created by ...WARNs — the "SyncExecutionResource" rows of the #6881 inventory):Since #6882, a bounded
getRange(from, until)iterator releases its reader inside thenext()that serves its last record, andhasNextprobes open nothing. So the fix is to size each read to what the method actually consumes — noclose()API added toVirtualDocument:totalCount == 1)getRange(0, 1)next()getRange(0, 1)next()getRange(1, totalCount)next()serving the bound / exhaustionNotes for review:
getCountwas sampled; they are now capped at the sametotalCountused for thetotalCount/skippedRowsfields, keeping the response internally consistent.Any related issues, documentation, discussions?
Follow-up to the residual-leak inventory in #6881 (
SyncExecutionResource.collectOperatorResultrows). Builds on the iterator restructure in #6882. The remaining #6881 rows (ResultExportServiceclient-disconnect,InputPortMaterializationReaderThreadinterrupt) are lower priority and left for separate PRs.How was this PR tested?
Verified by inspection against the bounded-read semantics of #6882: each path's read bound equals its consumption, so every path ends in a limit-close, an exhaustion-close, or nothing opened. The leak itself is only observable as GC-finalizer WARNs, so no unit test pins it;
getRangerange semantics are covered by the existingVirtualDocumentSpec/IcebergDocumentSpecsuites, and this endpoint's behavior is unchanged (same tuples, same truncation, same response shape). LocalWorkflowExecutionService/compileandscalafmtCheckpass.Was this PR authored or co-authored using generative AI tooling?
Generated-by: Claude Code (Fable 5)