Skip to content

[spark] Add opt-in repartitioning for oversized Blob table scans - #10457

Open
gavin9402 wants to merge 7 commits into
apache:masterfrom
gavin9402:repartition_large_scan
Open

gavin9402 wants to merge 7 commits into
apache:masterfrom
gavin9402:repartition_large_scan

Conversation

@gavin9402

@gavin9402 gavin9402 commented Oct 9, 2026 •

Copy link
Copy Markdown
Contributor

Motivation

While investigating a Spark query over a Paimon data-evolution table containing Blob columns, we observed that the scan stage launched only one task despite reading a large amount of data.

Data-evolution split planning preserves row-id groups. Files belonging to the same group can therefore remain in a single large split, which Spark-side bin packing cannot divide into independent scan partitions. Operators pipelined with that scan also run with limited parallelism.

This PR adds an optional shuffle immediately after the scan so downstream operators can run across more partitions. It preserves the source split boundaries required by the reader.

Changes

Add the RepartitionLargePaimonScan optimizer rule:

  • Apply only to tables whose full schema contains a BLOB, ARRAY<BLOB>, or MAP<..., BLOB> column, including scans where the Blob column is pruned. Skip other tables before planning splits.

  • Run after V2 scan pushdown, when the Paimon scan and its input partitions are available.

  • Reuse BinPackingSplits.filesMaxPartitionBytes to determine the target shuffle partition size, preserving Paimon and Spark configuration precedence.

  • Insert Repartition(..., shuffle = true) when any input partition exceeds twice that size. An input partition exactly at the threshold does not trigger repartitioning.

  • Calculate the partition count as:

    ceil(total bytes across all input partitions / filesMaxPartitionBytes)
    
  • Skip streaming scans and preserve an existing shuffle directly above the scan.

  • Avoid inserting duplicate repartitions when the rule runs repeatedly.

File-open costs and minimum partition counts do not affect the trigger threshold or shuffle partition count.

Add a connector option, disabled by default:

SET spark.paimon.read.repartition-large-scan.enabled = true;

When disabled, the rule returns the original plan without planning scan splits or input partitions.

Scope

This change increases downstream parallelism; it does not increase the number of source scan tasks or split row-id groups.

Partition sizes come from file metadata, including Blob file sizes. They do not represent decoded row sizes or measured shuffle bytes.

Validation

The Paimon unit tests cover:

  • Default and explicit disablement without scan planning.
  • Non-Blob tables skipped without scan planning, and Blob tables remaining eligible after column pruning.
  • Oversized-partition detection and exact-threshold behavior.
  • Total-byte accounting and rounding up.
  • Configuration precedence and independence from file-open costs and minimum partition counts.
  • Overflow-safe threshold doubling.
  • Existing shuffles, empty scans, repeated application, and partition-count overflow.
  • Rule registration after scan pushdown and unchanged non-Paimon plans.

@gavin9402 gavin9402 changed the title Repartition large scan [spark] Add opt-in repartitioning for oversized Paimon scan partitions Oct 9, 2026
@gavin9402 gavin9402 changed the title [spark] Add opt-in repartitioning for oversized Paimon scan partitions [spark] Add opt-in repartitioning for oversized Blob table scans Oct 9, 2026
@JingsongLi

Copy link
Copy Markdown
Contributor

[P2] Avoid reducing existing downstream parallelism when adding the shuffle

At RepartitionLargePaimonScan.scala:75-81, an oversized input partition triggers a new shuffle even when ceil(total bytes / targetSize) is smaller than the current input-partition count. Since this option applies to the whole session, an otherwise well-parallelized scan can lose downstream concurrency when one large Blob partition is mixed with many small partitions.

I reproduced this with a real Spark 3 Paimon data-evolution table partitioned by id, a Blob payload, 21 table partitions (20 one-byte payloads and one 256 KiB payload), source.split.target-size=64kb, default file-open cost, and AQE disabled. For SELECT id, length(payload), the option disabled gives 21 RDD partitions; enabled gives 5. All row results are preserved, but all 21 source tasks remain and the added shuffle reduces the downstream stage to 5 tasks. For expensive, roughly uniform per-row UDF/map work on a cluster with more than five available cores, that caps concurrency below the existing plan. I have not measured a runtime slowdown; the observed regression is the reduction in available downstream parallelism.

Please floor the computed count at the existing input-partition count, or skip insertion when it would not increase parallelism. The 14 existing rule tests pass, and a real 100-row Blob query also demonstrated the intended 2-to-15 downstream increase with identical results; adding the skewed-table case would cover this missing guard.

Reviewed head: 2f0c354a75b0946ad01a5bdfbeaff40da0b4f9e4.

@gavin9402

Copy link
Copy Markdown
Contributor Author

[P2] Avoid reducing existing downstream parallelism when adding the shuffle

At RepartitionLargePaimonScan.scala:75-81, an oversized input partition triggers a new shuffle even when ceil(total bytes / targetSize) is smaller than the current input-partition count. Since this option applies to the whole session, an otherwise well-parallelized scan can lose downstream concurrency when one large Blob partition is mixed with many small partitions.

I reproduced this with a real Spark 3 Paimon data-evolution table partitioned by id, a Blob payload, 21 table partitions (20 one-byte payloads and one 256 KiB payload), source.split.target-size=64kb, default file-open cost, and AQE disabled. For SELECT id, length(payload), the option disabled gives 21 RDD partitions; enabled gives 5. All row results are preserved, but all 21 source tasks remain and the added shuffle reduces the downstream stage to 5 tasks. For expensive, roughly uniform per-row UDF/map work on a cluster with more than five available cores, that caps concurrency below the existing plan. I have not measured a runtime slowdown; the observed regression is the reduction in available downstream parallelism.

Please floor the computed count at the existing input-partition count, or skip insertion when it would not increase parallelism. The 14 existing rule tests pass, and a real 100-row Blob query also demonstrated the intended 2-to-15 downstream increase with identical results; adding the skewed-table case would cover this missing guard.

Reviewed head: 2f0c354a75b0946ad01a5bdfbeaff40da0b4f9e4.

Thanks for the detailed reproduction. Updated the rule to skip repartitioning when the computed partition count is less than or equal to the existing input-partition count. In that case, the original scan is preserved without an additional shuffle.

@JingsongLi JingsongLi left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Reviewed e529202fa975409a7aded5c8a478159e059edae6 in an isolated worktree. The 16 existing rule tests and the standard Maven checks passed under Spark 3.5.8. The two inline findings were reproduced with real Paimon tables.

if (numPartitions <= partitions.size) {
relation
} else {
Repartition(numPartitions.toInt, shuffle = true, relation)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[P1] Evaluate input-file expressions before the inserted shuffle

When this option is enabled and a scan exceeds the threshold, inserting Repartition directly around the relation leaves an existing Project or Filter containing input_file_name() above the shuffle. The expression then runs in a downstream task without the source reader's InputFileBlockHolder context, so it returns an empty string and can change query results.

I reproduced this on Spark 3.5.8 with AQE disabled: create a data-evolution/row-tracking table containing id INT, set source.split.target-size=4kb, write 5,000 rows into one normal Parquet file, then run ALTER TABLE t ADD COLUMN payload BINARY COMMENT '__BLOB_FIELD'. The full schema makes this scan eligible even though payload is pruned. With the option disabled, SELECT id, input_file_name() FROM t returns the actual file path and SELECT id FROM t WHERE length(input_file_name()) > 0 returns all 5,000 rows. With the option enabled, the filename becomes empty and the same filtering query returns zero rows.

Please evaluate/materialize the file-dependent expressions below the added shuffle, or skip affected scans. Add a regression test for both the projection and filtering cases.

if (numPartitions <= partitions.size) {
relation
} else {
Repartition(numPartitions.toInt, shuffle = true, relation)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[P2] Keep dynamic partition pruning adjacent to the scan

The late optimizer batch runs after Spark inserts dynamic partition pruning. Rewriting an existing Filter(dynamicpruning, scan) as Filter(dynamicpruning, Repartition(scan)) prevents DataSourceV2Strategy from extracting that condition into BatchScanExec.runtimeFilters: its PhysicalOperation extractor cannot traverse Repartition. The condition instead becomes a regular FilterExec above the exchange, so Paimon's runtime partition filtering is never invoked and all source partitions are read and shuffled first.

I reproduced this on Spark 3.5.8 with AQE disabled and DPP enabled: a Blob fact table has three pt partitions, each with 100 rows and 8 KiB payloads, and source.split.target-size=64kb. Join it on pt to a broadcast dimension containing (pt, keep) = (0,1), (1,0), (2,0), with WHERE d.keep = 1, and select f.id, length(f.payload). With this option disabled, the fact scan has a runtime partition filter, reads 4 files and emits 100 rows. Enabled, its runtime filters are empty, it reads 12 files and emits 300 rows before the ordinary filter reduces the result. Final rows are identical, but source pruning is lost.

Please place the repartition above the existing dynamic-pruning filter while preserving Spark's scan extraction, and add an integration assertion that the fact scan retains its runtime filters.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for the detailed reproductions. Both issues are now addressed:

  • Input-file expressions: The rule conservatively skips scan subtrees beneath input_file_name(), input_file_block_start(), or input_file_block_length(), preserving the source reader’s file context.
  • Dynamic partition pruning: The repartition is inserted above the existing pruning filter, keeping the filter and scan together so Spark can extract BatchScanExec.runtimeFilters.

Added real-query regression tests for the input-file projection and filtering cases, plus a Blob join that verifies runtime filters remain present and the fact scan emits only the selected partition’s 100 rows.

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

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants