JustinBinber opened a new issue, #10021:
URL: https://github.com/apache/paimon/issues/10021

   ### Search before asking
   
   - [x] I searched in the [issues](https://github.com/apache/paimon/issues) 
and found nothing similar.
   
   
   ### Motivation
   
   ### Problem observed in practice
   
   While evaluating large-scale Spark `MERGE INTO` workloads on a Paimon table 
with many data files, I repeatedly saw the following pattern:
   
   - The same SQL could complete when the target partition range was small, but 
expanding the range caused the Driver to restart or fail with an out-of-memory 
error.
   - Some failures happened after the data-processing stages had already 
completed, around target split discovery, Spark task construction, or task 
serialization. This made the completed computation unusable and made retries 
expensive.
   - Reducing the number of target partitions per application and compacting 
small files improved stability, but significantly reduced throughput and did 
not remove the underlying capacity boundary.
   - Increasing shuffle parallelism, changing source split target size, or 
relying on Spark spill did not consistently solve the problem. Metadata used 
during planning and task serialization is not the same as spillable shuffle 
data.
   
   These experiments suggest that file metadata can become the limiting 
resource independently of the amount of row data processed by each task. The 
immediate workaround is to submit many smaller jobs, but this increases 
scheduling overhead, operational complexity, and total recovery cost.
   
   I am opening this issue because the problem appears to be structural rather 
than workload-specific: users with a sufficiently large number of files or 
overlapping file groups can encounter the same Driver/task-metadata boundary 
even when their SQL is otherwise valid. A bounded metadata path would make 
large scans and merges degrade predictably instead of requiring trial-and-error 
partition sizing.
   
   ### Memory boundaries
   
   Spark scans can retain and serialize an unbounded amount of file metadata 
when a Paimon table contains many files. The pressure appears at three 
different boundaries:
   
   1. **Driver planning:** `AbstractFileStoreScan.plan()` materializes manifest 
entries; `TableScan.Plan.splits()` exposes a `List<Split>`, and Spark bin 
packing consumes an `Array[Split]`.
   2. **Task transport:** a Spark `InputPartition` embeds all of its `Split` 
objects, so a large partition increases Driver serialization/RPC pressure and 
is retransmitted on retry or speculation.
   3. **Executor reconstruction:** the Executor reconstructs the complete split 
list before `PaimonPartitionReader` iterates it.
   
   A local synthetic probe using public Paimon `DataFileMeta` / `DataSplit` 
shapes built 350 splits with 1,000 file descriptors per split:
   
   - inline `InputPartition`: 136,489,062 bytes
   - external-reference descriptor: about 1.4 KB
   - external container: 136,585,309 bytes
   - current-thread cumulative allocation: about 852 MiB for inline encoding, 
518 MiB for external encoding, and 513 MiB for full external decoding
   
   This is a metadata-only capacity probe, not an end-to-end file scan. 
Cumulative allocation is not peak heap, and these numbers do not claim a 
speedup.
   
   Issue #5773 reports a Driver OOM around `AbstractFileStoreScan.plan()`, but 
it also includes separate aggregation/type-conversion behavior. This proposal 
focuses on bounding scan-planning and task-metadata memory.
   
   ### Solution
   
   I propose a staged, opt-in design so each change remains independently 
reviewable:
   
   1. **Externalize oversized Spark InputPartition metadata.** Keep small 
partitions inline. Store large metadata in a seekable shared `FileIO` path 
using a versioned framed format. Send only a descriptor containing path, 
offset, length, counts, and checksum. Use range reads and query-scoped cleanup. 
Keep the feature disabled by default initially.
   2. **Stream split decoding on Executors.** Expose a closeable iterator so a 
Reader does not first build a complete `List<Split>` for one InputPartition.
   3. **Add a compatible streaming/page-based scan path.** Plan manifest 
entries and splits incrementally and feed Spark bin packing without retaining 
the complete scan in Driver heap. Preserve the current `Plan.splits()` behavior 
for existing consumers and fall back to eager planning where streaming is 
unsafe.
   4. **Bound a single huge DataSplit.** Ordinary scans may split independent 
files. Data Evolution scans must split only independent row-id-range overlap 
components, and must fall back when independence cannot be proven.
   
   Questions for maintainers:
   
   - Is an opt-in shared `FileIO` staging path acceptable for oversized Spark 
InputPartition metadata?
   - What API shape is preferred for streaming scan planning while preserving 
`TableScan.Plan.splits()` compatibility?
   - For Data Evolution, is row-id-range connected-component splitting an 
acceptable direction, or should it be discussed separately?
   
   Compatibility and non-goals:
   
   - Default behavior remains unchanged initially.
   - No change to merge or deduplication semantics.
   - No claim of end-to-end performance improvement until a real scan benchmark 
is available.
   - The first PR would cover only task-metadata transport; later PRs would 
depend on review and consensus.
   
   ### Anything else?
   
   I have a proof of concept and focused tests for:
   
   - small-value inline fallback
   - shared-container range reads
   - failure cleanup
   - Spark Executor reads
   - bucket preservation
   - a 128 MiB threshold probe
   
   No proprietary code, internal logs, table names, cluster configuration, or 
private infrastructure details are included. I would like to submit the 
implementation as small PRs after agreeing on the design.
   
   ### Are you willing to submit a PR?
   
   - [x] I'm willing to submit a PR!


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

Reply via email to