[
https://issues.apache.org/jira/browse/SPARK-59123?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Dongjoon Hyun updated SPARK-59123:
----------------------------------
Parent: SPARK-58540
Issue Type: Sub-task (was: Improvement)
> Avoid per-key intermediate collections in KeyedPartitioning.projectKeys and
> reduceKeys
> --------------------------------------------------------------------------------------
>
> Key: SPARK-59123
> URL: https://issues.apache.org/jira/browse/SPARK-59123
> Project: Spark
> Issue Type: Sub-task
> Components: SQL
> Affects Versions: 4.4.0
> Reporter: Peter Toth
> Priority: Major
> Labels: pull-request-available
>
> `KeyedPartitioning.projectKeys` and `KeyedPartitioning.reduceKeys` rebuild
> several throwaway collections per partition key, to end up with one array
> each.
> {{projectKeys}} materialises an intermediate {{Seq}} per key and then copies
> it into an array, destructuring a {{Tuple2}} per position:
> {code:scala}
> val projectedKey = positionsWithTypes.map {
> case (position, dataType) => key.row.get(position, dataType)
> }.toArray[Any]
> {code}
> {{reduceKeys}} does the same four times over. {{key.row.toSeq(dataTypes)}}
> allocates an array and an {{ArraySeq}} wrapper, {{zip(reducers)}} a sequence
> of tuples, {{map}} a third sequence, and {{toArray}} the array that was
> wanted in the first place. The erased {{Some(reducer: Reducer[Any, Any])}}
> type test also runs once per key per column:
> {code:scala}
> val keyValues = key.row.toSeq(dataTypes)
> val reducedKey = keyValues.zip(reducers).map {
> case (v, Some(reducer: Reducer[Any, Any])) => reducer.reduce(v)
> case (v, _) => v
> }.toArray
> {code}
> Both are O(number of partition keys), which is the number of splits a scan
> reports, so tens of thousands is ordinary. {{projectKeys}} runs over all of
> them on every {{EnsureRequirements}} and {{ValidateRequirements}} pass
> whenever
> {{spark.sql.sources.v2.bucketing.allowKeysSubsetOfPartitionKeys.enabled}} is
> on, not only on the reducer path.
> The fix is local to the two bodies: hoist the positions, data types and
> reducers into arrays once, then fill a single {{Array[Any]}} with an indexed
> loop, and decide the reducer per column outside the key loop. No behaviour
> change, so the existing storage-partitioned-join tests cover it.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]