-
Notifications
You must be signed in to change notification settings - Fork 28.3k
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
[SPARK-47094][SQL] SPJ : Dynamically rebalance number of buckets when…
… they are not equal ### What changes were proposed in this pull request? -- Allow SPJ between 'compatible' bucket funtions -- Add a mechanism to define 'reducible' functions, one function whose output can be 'reduced' to another for all inputs. ### Why are the changes needed? -- SPJ currently applies only if the partition transform expressions on both sides are identifical. ### Does this PR introduce _any_ user-facing change? No ### How was this patch tested? Added new tests in KeyGroupedPartitioningSuite ### Was this patch authored or co-authored using generative AI tooling? No Closes #45267 from szehon-ho/spj-uneven-buckets. Authored-by: Szehon Ho <szehon.apache@gmail.com> Signed-off-by: Chao Sun <chao@openai.com>
- Loading branch information
Showing
9 changed files
with
821 additions
and
15 deletions.
There are no files selected for viewing
42 changes: 42 additions & 0 deletions
42
sql/catalyst/src/main/java/org/apache/spark/sql/connector/catalog/functions/Reducer.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,42 @@ | ||
/* | ||
* Licensed to the Apache Software Foundation (ASF) under one or more | ||
* contributor license agreements. See the NOTICE file distributed with | ||
* this work for additional information regarding copyright ownership. | ||
* The ASF licenses this file to You under the Apache License, Version 2.0 | ||
* (the "License"); you may not use this file except in compliance with | ||
* the License. You may obtain a copy of the License at | ||
* | ||
* http://www.apache.org/licenses/LICENSE-2.0 | ||
* | ||
* Unless required by applicable law or agreed to in writing, software | ||
* distributed under the License is distributed on an "AS IS" BASIS, | ||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. | ||
* See the License for the specific language governing permissions and | ||
* limitations under the License. | ||
*/ | ||
package org.apache.spark.sql.connector.catalog.functions; | ||
|
||
import org.apache.spark.annotation.Evolving; | ||
|
||
/** | ||
* A 'reducer' for output of user-defined functions. | ||
* | ||
* @see ReducibleFunction | ||
* | ||
* A user defined function f_source(x) is 'reducible' on another user_defined function | ||
* f_target(x) if | ||
* <ul> | ||
* <li> There exists a reducer function r(x) such that r(f_source(x)) = f_target(x) for | ||
* all input x, or </li> | ||
* <li> More generally, there exists reducer functions r1(x) and r2(x) such that | ||
* r1(f_source(x)) = r2(f_target(x)) for all input x. </li> | ||
* </ul> | ||
* | ||
* @param <I> reducer input type | ||
* @param <O> reducer output type | ||
* @since 4.0.0 | ||
*/ | ||
@Evolving | ||
public interface Reducer<I, O> { | ||
O reduce(I arg); | ||
} |
106 changes: 106 additions & 0 deletions
106
...yst/src/main/java/org/apache/spark/sql/connector/catalog/functions/ReducibleFunction.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,106 @@ | ||
/* | ||
* Licensed to the Apache Software Foundation (ASF) under one or more | ||
* contributor license agreements. See the NOTICE file distributed with | ||
* this work for additional information regarding copyright ownership. | ||
* The ASF licenses this file to You under the Apache License, Version 2.0 | ||
* (the "License"); you may not use this file except in compliance with | ||
* the License. You may obtain a copy of the License at | ||
* | ||
* http://www.apache.org/licenses/LICENSE-2.0 | ||
* | ||
* Unless required by applicable law or agreed to in writing, software | ||
* distributed under the License is distributed on an "AS IS" BASIS, | ||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. | ||
* See the License for the specific language governing permissions and | ||
* limitations under the License. | ||
*/ | ||
package org.apache.spark.sql.connector.catalog.functions; | ||
|
||
import org.apache.spark.annotation.Evolving; | ||
|
||
/** | ||
* Base class for user-defined functions that can be 'reduced' on another function. | ||
* | ||
* A function f_source(x) is 'reducible' on another function f_target(x) if | ||
* <ul> | ||
* <li> There exists a reducer function r(x) such that r(f_source(x)) = f_target(x) | ||
* for all input x, or </li> | ||
* <li> More generally, there exists reducer functions r1(x) and r2(x) such that | ||
* r1(f_source(x)) = r2(f_target(x)) for all input x. </li> | ||
* </ul> | ||
* <p> | ||
* Examples: | ||
* <ul> | ||
* <li>Bucket functions where one side has reducer | ||
* <ul> | ||
* <li>f_source(x) = bucket(4, x)</li> | ||
* <li>f_target(x) = bucket(2, x)</li> | ||
* <li>r(x) = x % 2</li> | ||
* </ul> | ||
* | ||
* <li>Bucket functions where both sides have reducer | ||
* <ul> | ||
* <li>f_source(x) = bucket(16, x)</li> | ||
* <li>f_target(x) = bucket(12, x)</li> | ||
* <li>r1(x) = x % 4</li> | ||
* <li>r2(x) = x % 4</li> | ||
* </ul> | ||
* | ||
* <li>Date functions | ||
* <ul> | ||
* <li>f_source(x) = days(x)</li> | ||
* <li>f_target(x) = hours(x)</li> | ||
* <li>r(x) = x / 24</li> | ||
* </ul> | ||
* </ul> | ||
* @param <I> reducer function input type | ||
* @param <O> reducer function output type | ||
* @since 4.0.0 | ||
*/ | ||
@Evolving | ||
public interface ReducibleFunction<I, O> { | ||
|
||
/** | ||
* This method is for the bucket function. | ||
* | ||
* If this bucket function is 'reducible' on another bucket function, | ||
* return the {@link Reducer} function. | ||
* <p> | ||
* For example, to return reducer for reducing f_source = bucket(4, x) on f_target = bucket(2, x) | ||
* <ul> | ||
* <li>thisBucketFunction = bucket</li> | ||
* <li>thisNumBuckets = 4</li> | ||
* <li>otherBucketFunction = bucket</li> | ||
* <li>otherNumBuckets = 2</li> | ||
* </ul> | ||
* | ||
* @param thisNumBuckets parameter for this function | ||
* @param otherBucketFunction the other parameterized function | ||
* @param otherNumBuckets parameter for the other function | ||
* @return a reduction function if it is reducible, null if not | ||
*/ | ||
default Reducer<I, O> reducer( | ||
int thisNumBuckets, | ||
ReducibleFunction<?, ?> otherBucketFunction, | ||
int otherNumBuckets) { | ||
throw new UnsupportedOperationException(); | ||
} | ||
|
||
/** | ||
* This method is for all other functions. | ||
* | ||
* If this function is 'reducible' on another function, return the {@link Reducer} function. | ||
* <p> | ||
* Example of reducing f_source = days(x) on f_target = hours(x) | ||
* <ul> | ||
* <li>thisFunction = days</li> | ||
* <li>otherFunction = hours</li> | ||
* </ul> | ||
* | ||
* @param otherFunction the other function | ||
* @return a reduction function if it is reducible, null if not. | ||
*/ | ||
default Reducer<I, O> reducer(ReducibleFunction<?, ?> otherFunction) { | ||
throw new UnsupportedOperationException(); | ||
} | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Oops, something went wrong.