[
https://issues.apache.org/jira/browse/SPARK-58265?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
ASF GitHub Bot updated SPARK-58265:
-----------------------------------
Labels: pull-request-available (was: )
> Reuse projected broadcast values for dynamic partition pruning
> --------------------------------------------------------------
>
> Key: SPARK-58265
> URL: https://issues.apache.org/jira/browse/SPARK-58265
> Project: Spark
> Issue Type: Improvement
> Components: Optimizer, SQL
> Affects Versions: 4.3.0, 5.0.0
> Reporter: Chao Sun
> Priority: Major
> Labels: pull-request-available
>
> h2. Problem
> Dynamic partition pruning can already reuse an existing broadcast when the
> values needed to prune a partitioned table are the broadcast's join keys.
> However, a broadcast hash relation also stores complete build-side rows.
> Spark cannot currently reuse another column from those rows, or a
> deterministic expression over that column, even when that column identifies
> the exact partitions needed by a later join.
> As a result, a query can read every partition of a large table even though an
> earlier broadcast already contains the information needed to avoid that scan.
> When {{spark.sql.optimizer.dynamicPartitionPruning.reuseBroadcastOnly=true}},
> Spark drops the pruning filter because its existing key-based reuse path does
> not match. When standalone pruning is allowed, Spark may instead execute the
> filtering-side plan again to compute information that is already present in
> the broadcast.
> h2. Example
> Suppose {{wide_history}} is partitioned by {{ds}}:
> {code:sql}
> WITH cohorts AS (
> SELECT /*+ BROADCAST(d) */
> d.cohort_date,
> u.user_id
> FROM dates d
> JOIN users u
> ON d.cohort_ds = u.cohort_ds
> AND d.region = u.region
> WHERE u.active
> )
> SELECT c.user_id, h.payload
> FROM cohorts c
> LEFT JOIN wide_history h
> ON h.user_id = c.user_id
> AND h.ds = date_format(
> date_add(c.cohort_date, 6),
> 'yyyy-MM-dd')
> {code}
> Spark broadcasts {{dates}} for the first join, using {{(cohort_ds, region)}}
> as the broadcast hash keys. Every broadcast row also contains
> {{cohort_date}}, but {{cohort_date}} is not itself a hash key. The later join
> needs partition values calculated from {{date_add(cohort_date, 6)}}.
> For matching cohort dates {{2024-01-01}} and {{2024-01-08}}, the history scan
> should read partitions {{2024-01-07}} and {{2024-01-14}}. Today, the existing
> reuse path cannot derive those values from the broadcast. It therefore either
> scans every {{wide_history}} partition or runs a separate filtering-side plan.
> h2. Expected behavior
> Preserve existing direct broadcast-key reuse as the first choice. If that
> path cannot produce the required partition values, allow dynamic partition
> pruning to identify a compatible broadcast already required by an earlier
> inner join and project a supported deterministic expression from its complete
> value rows.
> Reuse only the broadcast of the verified source relation with the complete,
> ordered join keys and a compatible broadcast hash mode. The projected values
> should feed the existing partition-pruning predicate without creating another
> broadcast, rerunning the filtering-side plan, or changing the original join's
> results.
> h2. Correctness
> A projected pruning domain may be a safe superset of the final join keys. For
> example, if the broadcast also contains an unmatched cohort date
> {{2024-01-29}}, projection may include history partition {{2024-02-04}}.
> Reading that additional partition is safe because the original join still
> removes unmatched rows. Excluding either required partition would be
> incorrect.
> Broadcast inspection and projected values must be bounded. If the broadcast
> cannot be verified, an expression is unsupported, source statistics are
> unavailable, a resource limit is exceeded, or projection cannot produce the
> complete domain, pruning must fail open. In particular, an unavailable domain
> must never be interpreted as an empty domain, and a disabled pruning
> predicate must not be passed to a V2 connector as a real partition filter.
> The optimization should preserve existing behavior when disabled and support
> adaptive and nonadaptive planning, ANSI mode, composite broadcast keys,
> duplicate or null values, and both V1 and V2 partition filtering.
> h2. Scope
> This issue covers projected broadcast-value reuse for dynamic partition
> pruning. Broadcast-backed Bloom filters are related future work and are
> separately tracked by SPARK-40909.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]