Skip to content

[Feature][Spark] Bound split-planning and InputPartition metadata memory for scans with many files #10021

Description

@JustinBinber

Search before asking

  • I searched in the 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?

  • I'm willing to submit a PR!

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    enhancementNew feature or request

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions