Repository navigation
Conversation
|
[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 I reproduced this with a real Spark 3 Paimon data-evolution table partitioned by 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: |
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
left a comment
There was a problem hiding this comment.
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) |
There was a problem hiding this comment.
[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) |
There was a problem hiding this comment.
[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.
There was a problem hiding this comment.
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.
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
RepartitionLargePaimonScanoptimizer rule:Apply only to tables whose full schema contains a
BLOB,ARRAY<BLOB>, orMAP<..., 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.filesMaxPartitionBytesto 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:
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:
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: