Uh oh!
There was an error while loading. Please reload this page.
[SPARK-54199][SQL] Add DataFrame API support for new KLL quantiles sketch functions - #52900
[SPARK-54199][SQL] Add DataFrame API support for new KLL quantiles sketch functions#52900dtenedor wants to merge 8 commits into
Conversation
dtenedor
commented
Nov 5, 2025
| * | ||
| * @group agg_funcs | ||
| * @since 4.1.0 | ||
| * @since 4.2.0 |
There was a problem hiding this comment.
Do we want to add something in the PR description about these versioning changes?
There was a problem hiding this comment.
My mistake on this version change, it looks like 4.1.0 is the correct version after all. Reverted that.
| kll_sketch_get_quantile_bigint.__doc__ = pysparkfuncs.kll_sketch_get_quantile_bigint.__doc__ | ||
| def kll_sketch_get_quantile_float(sketch: "ColumnOrName", rank: "ColumnOrName") -> Column: |
There was a problem hiding this comment.
noticed that in the line below, we don't check if the ranks column type is not supported, not sure what would happen in that case. https://github.com/apache/spark/blob/master/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/expressions/kllExpressions.scala#L488
There was a problem hiding this comment.
Looking at this, I think the class mixing in ImplicitCastInputTypes should make sure that both arguments have either the same types or are coercible to the required types. I'll add a couple negative tests just to make sure.
abstract class KllSketchGetQuantileBase
extends BinaryExpression
with CodegenFallback
with ImplicitCastInputTypes {
...
override def inputTypes: Seq[AbstractDataType] =
Seq(
BinaryType,
TypeCollection(
DoubleType,
ArrayType(DoubleType, containsNull = false)))
There was a problem hiding this comment.
Ah then the checkInputDataTypes function probably does the check when the expression is resolved. I missed that.
cboumalh
left a comment
There was a problem hiding this comment.
LGTM, just have the comment above, but all looks good. Thank you!
dtenedor
commented
Nov 6, 2025
Thanks for your review! |
dongjoon-hyun
commented
Nov 6, 2025
Hi, @dtenedor . I don't see any committers' approval yet.
![]() |
dongjoon-hyun
left a comment
There was a problem hiding this comment.
Hi, @cloud-fan . It would be great if someone from Databricks gave a post-commit approval for this because the code is targeting Apache Spark 4.1.0 like the following. @dtenedor seems to be unable to get a proper chance of community review yet.
/**
* Aggregate function: returns the compact binary representation of the Datasketches
* KllLongsSketch built with the values in the input column. The optional k parameter controls
* the size and accuracy of the sketch (default 200, range 8-65535).
*
* @group agg_funcs
* @since 4.1.0
*/
def kll_sketch_agg_bigint(e: Column, k: Column): Column =
Column.fn("kll_sketch_agg_bigint", e, k)
dongjoon-hyun
commented
Nov 6, 2025
In addition, if this is not targeting for Apache Spark 4.1.0 technically, please update the wrong |
cloud-fan
commented
Nov 7, 2025
The change LGTM, it's just adding dataframe functions to match SQL. Given the related SQL functions are in 4.1, I think it makes sense to backport this to 4.1 as well. @dtenedor what do you think? |
dongjoon-hyun
commented
Nov 7, 2025
Thank you, @cloud-fan ! |
…etch functions ### What changes were proposed in this pull request? This PR adds DataFrame API support for the KLL quantile sketch functions that were previously added to Spark SQL in #52800. This lets users leverage KLL sketches through both Scala and Python DataFrame APIs in addition to the existing SQL interface. **Key additions:** 1. **Scala DataFrame API** (`sql/api/src/main/scala/org/apache/spark/sql/functions.scala`): - 18 new functions covering aggregate, merge, quantile, and rank operations - Multiple overloads for each function supporting: - `Column` parameters for computed values - `String` parameters for column names - `Int` parameters for literal k values - Optional k parameters with sensible defaults - Functions for all three data type variants: bigint, float, double 2. **Python DataFrame API** (`python/pyspark/sql/functions/builtin.py`): - 18 corresponding Python functions with: - Comprehensive docstrings with usage examples - Proper type hints (`ColumnOrName`, `Optional[Union[int, Column]]`) - Support for both column objects and column name strings - Added to PySpark documentation reference 3. **Python Spark Connect Support** (`python/pyspark/sql/connect/functions/builtin.py`): - Full compatibility with Spark Connect architecture - All 18 functions properly registered ### Why are the changes needed? While the SQL API for KLL sketches was previously added, DataFrame API support is essential for full usability. Without DataFrame API support, users would be forced to use SQL expressions via `expr()` or `selectExpr()`, which is less ergonomic and type-safe. ### Does this PR introduce any user-facing change? Yes, this PR adds DataFrame API support for the 18 KLL sketch functions: **Scala DataFrame API Example:** ```scala import org.apache.spark.sql.functions._ // Create sketch with default k val df = Seq(1, 2, 3, 4, 5).toDF("value") val sketch = df.agg(kll_sketch_agg_bigint($"value")) // Create sketch with custom k value val sketch2 = df.agg(kll_sketch_agg_bigint("value", 400)) // Get median (0.5 quantile) val sketchDf = df.agg(kll_sketch_agg_bigint($"value").alias("sketch")) val median = sketchDf.select(kll_sketch_get_quantile_bigint($"sketch", lit(0.5))) // Get multiple quantiles val quantiles = sketchDf.select( kll_sketch_get_quantile_bigint($"sketch", array(lit(0.25), lit(0.5), lit(0.75))) ) // Merge sketches val merged = sketchDf.select( kll_sketch_merge_bigint($"sketch", $"sketch").alias("merged") ) // Get count of items val count = sketchDf.select(kll_sketch_get_n_bigint($"sketch")) ``` **Python DataFrame API Example:** ```python from pyspark.sql import functions as sf # Create sketch with default k df = spark.createDataFrame([1, 2, 3, 4, 5], "INT") sketch = df.agg(sf.kll_sketch_agg_bigint("value")) # Create sketch with custom k value sketch2 = df.agg(sf.kll_sketch_agg_bigint("value", 400)) # Get median (0.5 quantile) sketch_df = df.agg(sf.kll_sketch_agg_bigint("value").alias("sketch")) median = sketch_df.select(sf.kll_sketch_get_quantile_bigint("sketch", sf.lit(0.5))) # Get multiple quantiles quantiles = sketch_df.select( sf.kll_sketch_get_quantile_bigint("sketch", sf.array(sf.lit(0.25), sf.lit(0.5), sf.lit(0.75))) ) # Merge sketches merged = sketch_df.select( sf.kll_sketch_merge_bigint("sketch", "sketch").alias("merged") ) # Get count of items count = sketch_df.select(sf.kll_sketch_get_n_bigint("sketch")) ``` ### How was this patch tested? 1. **Scala Unit Tests** (`DataFrameAggregateSuite`): - `kll_sketch_agg_{bigint,float,double}` with default and explicit k values - `kll_sketch_to_string` functions for all data types - `kll_sketch_get_n` functions for all data types - `kll_sketch_merge` operations - `kll_sketch_get_quantile` with single rank and array of ranks - `kll_sketch_get_rank` operations - Null value handling tests 2. **Python Unit Tests** (`test_functions.py`): - Comprehensive tests mirroring Scala tests - Tests for Column object and string column name overloads - Tests for optional k parameter - Array input tests for quantile/rank functions - Null handling validation - Type checking (bytes/bytearray for sketches, str for to_string, int/float for values) ### Was this patch authored or co-authored using generative AI tooling? Yes, IDE assistance used `claude-4.5-sonnet` with manual validation and integration. Closes#52900 from dtenedor/dataframe-api-kll-functions. Authored-by: Daniel Tenedorio <daniel.tenedorio@databricks.com> Signed-off-by: Daniel Tenedorio <daniel.tenedorio@databricks.com>
cloud-fan
commented
Nov 11, 2025
I've backported both #52800 and this PR to 4.1 |
dtenedor
commented
Nov 12, 2025
Thank you @cloud-fan for the review + backporting 👍 |
dtenedor
commented
Nov 12, 2025
@dongjoon-hyun my mistake on the review protocols, I did not know that part as a new committer; I will get another committer(s) review for my own future work |
…etch functions ### What changes were proposed in this pull request? This PR adds DataFrame API support for the KLL quantile sketch functions that were previously added to Spark SQL in apache#52800. This lets users leverage KLL sketches through both Scala and Python DataFrame APIs in addition to the existing SQL interface. **Key additions:** 1. **Scala DataFrame API** (`sql/api/src/main/scala/org/apache/spark/sql/functions.scala`): - 18 new functions covering aggregate, merge, quantile, and rank operations - Multiple overloads for each function supporting: - `Column` parameters for computed values - `String` parameters for column names - `Int` parameters for literal k values - Optional k parameters with sensible defaults - Functions for all three data type variants: bigint, float, double 2. **Python DataFrame API** (`python/pyspark/sql/functions/builtin.py`): - 18 corresponding Python functions with: - Comprehensive docstrings with usage examples - Proper type hints (`ColumnOrName`, `Optional[Union[int, Column]]`) - Support for both column objects and column name strings - Added to PySpark documentation reference 3. **Python Spark Connect Support** (`python/pyspark/sql/connect/functions/builtin.py`): - Full compatibility with Spark Connect architecture - All 18 functions properly registered ### Why are the changes needed? While the SQL API for KLL sketches was previously added, DataFrame API support is essential for full usability. Without DataFrame API support, users would be forced to use SQL expressions via `expr()` or `selectExpr()`, which is less ergonomic and type-safe. ### Does this PR introduce any user-facing change? Yes, this PR adds DataFrame API support for the 18 KLL sketch functions: **Scala DataFrame API Example:** ```scala import org.apache.spark.sql.functions._ // Create sketch with default k val df = Seq(1, 2, 3, 4, 5).toDF("value") val sketch = df.agg(kll_sketch_agg_bigint($"value")) // Create sketch with custom k value val sketch2 = df.agg(kll_sketch_agg_bigint("value", 400)) // Get median (0.5 quantile) val sketchDf = df.agg(kll_sketch_agg_bigint($"value").alias("sketch")) val median = sketchDf.select(kll_sketch_get_quantile_bigint($"sketch", lit(0.5))) // Get multiple quantiles val quantiles = sketchDf.select( kll_sketch_get_quantile_bigint($"sketch", array(lit(0.25), lit(0.5), lit(0.75))) ) // Merge sketches val merged = sketchDf.select( kll_sketch_merge_bigint($"sketch", $"sketch").alias("merged") ) // Get count of items val count = sketchDf.select(kll_sketch_get_n_bigint($"sketch")) ``` **Python DataFrame API Example:** ```python from pyspark.sql import functions as sf # Create sketch with default k df = spark.createDataFrame([1, 2, 3, 4, 5], "INT") sketch = df.agg(sf.kll_sketch_agg_bigint("value")) # Create sketch with custom k value sketch2 = df.agg(sf.kll_sketch_agg_bigint("value", 400)) # Get median (0.5 quantile) sketch_df = df.agg(sf.kll_sketch_agg_bigint("value").alias("sketch")) median = sketch_df.select(sf.kll_sketch_get_quantile_bigint("sketch", sf.lit(0.5))) # Get multiple quantiles quantiles = sketch_df.select( sf.kll_sketch_get_quantile_bigint("sketch", sf.array(sf.lit(0.25), sf.lit(0.5), sf.lit(0.75))) ) # Merge sketches merged = sketch_df.select( sf.kll_sketch_merge_bigint("sketch", "sketch").alias("merged") ) # Get count of items count = sketch_df.select(sf.kll_sketch_get_n_bigint("sketch")) ``` ### How was this patch tested? 1. **Scala Unit Tests** (`DataFrameAggregateSuite`): - `kll_sketch_agg_{bigint,float,double}` with default and explicit k values - `kll_sketch_to_string` functions for all data types - `kll_sketch_get_n` functions for all data types - `kll_sketch_merge` operations - `kll_sketch_get_quantile` with single rank and array of ranks - `kll_sketch_get_rank` operations - Null value handling tests 2. **Python Unit Tests** (`test_functions.py`): - Comprehensive tests mirroring Scala tests - Tests for Column object and string column name overloads - Tests for optional k parameter - Array input tests for quantile/rank functions - Null handling validation - Type checking (bytes/bytearray for sketches, str for to_string, int/float for values) ### Was this patch authored or co-authored using generative AI tooling? Yes, IDE assistance used `claude-4.5-sonnet` with manual validation and integration. Closesapache#52900 from dtenedor/dataframe-api-kll-functions. Authored-by: Daniel Tenedorio <daniel.tenedorio@databricks.com> Signed-off-by: Daniel Tenedorio <daniel.tenedorio@databricks.com>
…etch functions ### What changes were proposed in this pull request? This PR adds DataFrame API support for the KLL quantile sketch functions that were previously added to Spark SQL in apache#52800. This lets users leverage KLL sketches through both Scala and Python DataFrame APIs in addition to the existing SQL interface. **Key additions:** 1. **Scala DataFrame API** (`sql/api/src/main/scala/org/apache/spark/sql/functions.scala`): - 18 new functions covering aggregate, merge, quantile, and rank operations - Multiple overloads for each function supporting: - `Column` parameters for computed values - `String` parameters for column names - `Int` parameters for literal k values - Optional k parameters with sensible defaults - Functions for all three data type variants: bigint, float, double 2. **Python DataFrame API** (`python/pyspark/sql/functions/builtin.py`): - 18 corresponding Python functions with: - Comprehensive docstrings with usage examples - Proper type hints (`ColumnOrName`, `Optional[Union[int, Column]]`) - Support for both column objects and column name strings - Added to PySpark documentation reference 3. **Python Spark Connect Support** (`python/pyspark/sql/connect/functions/builtin.py`): - Full compatibility with Spark Connect architecture - All 18 functions properly registered ### Why are the changes needed? While the SQL API for KLL sketches was previously added, DataFrame API support is essential for full usability. Without DataFrame API support, users would be forced to use SQL expressions via `expr()` or `selectExpr()`, which is less ergonomic and type-safe. ### Does this PR introduce any user-facing change? Yes, this PR adds DataFrame API support for the 18 KLL sketch functions: **Scala DataFrame API Example:** ```scala import org.apache.spark.sql.functions._ // Create sketch with default k val df = Seq(1, 2, 3, 4, 5).toDF("value") val sketch = df.agg(kll_sketch_agg_bigint($"value")) // Create sketch with custom k value val sketch2 = df.agg(kll_sketch_agg_bigint("value", 400)) // Get median (0.5 quantile) val sketchDf = df.agg(kll_sketch_agg_bigint($"value").alias("sketch")) val median = sketchDf.select(kll_sketch_get_quantile_bigint($"sketch", lit(0.5))) // Get multiple quantiles val quantiles = sketchDf.select( kll_sketch_get_quantile_bigint($"sketch", array(lit(0.25), lit(0.5), lit(0.75))) ) // Merge sketches val merged = sketchDf.select( kll_sketch_merge_bigint($"sketch", $"sketch").alias("merged") ) // Get count of items val count = sketchDf.select(kll_sketch_get_n_bigint($"sketch")) ``` **Python DataFrame API Example:** ```python from pyspark.sql import functions as sf # Create sketch with default k df = spark.createDataFrame([1, 2, 3, 4, 5], "INT") sketch = df.agg(sf.kll_sketch_agg_bigint("value")) # Create sketch with custom k value sketch2 = df.agg(sf.kll_sketch_agg_bigint("value", 400)) # Get median (0.5 quantile) sketch_df = df.agg(sf.kll_sketch_agg_bigint("value").alias("sketch")) median = sketch_df.select(sf.kll_sketch_get_quantile_bigint("sketch", sf.lit(0.5))) # Get multiple quantiles quantiles = sketch_df.select( sf.kll_sketch_get_quantile_bigint("sketch", sf.array(sf.lit(0.25), sf.lit(0.5), sf.lit(0.75))) ) # Merge sketches merged = sketch_df.select( sf.kll_sketch_merge_bigint("sketch", "sketch").alias("merged") ) # Get count of items count = sketch_df.select(sf.kll_sketch_get_n_bigint("sketch")) ``` ### How was this patch tested? 1. **Scala Unit Tests** (`DataFrameAggregateSuite`): - `kll_sketch_agg_{bigint,float,double}` with default and explicit k values - `kll_sketch_to_string` functions for all data types - `kll_sketch_get_n` functions for all data types - `kll_sketch_merge` operations - `kll_sketch_get_quantile` with single rank and array of ranks - `kll_sketch_get_rank` operations - Null value handling tests 2. **Python Unit Tests** (`test_functions.py`): - Comprehensive tests mirroring Scala tests - Tests for Column object and string column name overloads - Tests for optional k parameter - Array input tests for quantile/rank functions - Null handling validation - Type checking (bytes/bytearray for sketches, str for to_string, int/float for values) ### Was this patch authored or co-authored using generative AI tooling? Yes, IDE assistance used `claude-4.5-sonnet` with manual validation and integration. Closesapache#52900 from dtenedor/dataframe-api-kll-functions. Authored-by: Daniel Tenedorio <daniel.tenedorio@databricks.com> Signed-off-by: Daniel Tenedorio <daniel.tenedorio@databricks.com>
### What changes were proposed in this pull request? This PR adds SQL aggregate functions with their tests for the KLL merge aggregate functions: - `kll_merge_agg_bigint` - `kll_merge_agg_float` - `kll_merge_agg_double` These aggregate functions merge multiple binary KLL sketch representations. Initial PRs: - #52900 - #52800 ### Why are the changes needed? The existing scalar `kll_sketch_merge_*` functions can only merge two sketches at a time. In distributed computing scenarios where sketches are pre-computed across multiple partitions, time windows, or datasets, users need to merge many sketches together. ### Does this PR introduce _any_ user-facing change? Yes, this PR adds 3 new aggregate functions. ### How was this patch tested? New SQL tests were added to `sql/core/src/test/resources/sql-tests/inputs/kllquantiles.sql`: **Positive tests:** - Merging bigint/float/double sketches from multiple rows - Merging with custom k parameters (400, 300, 500) - NULL value handling **Negative tests:** - Type mismatches (passing non-binary types) - Invalid binary data - k parameter validation (too small, too large, NULL, non-constant) ### Was this patch authored or co-authored using generative AI tooling? claude-4.5-sonnet and manual changes. Closes#53548 from cboumalh/cboumalh-kll-enhancement. Lead-authored-by: Chris Boumalhab <cboumalh@amazon.com> Co-authored-by: Chris Boumalhab <84485659+cboumalh@users.noreply.github.com> Signed-off-by: Daniel Tenedorio <daniel.tenedorio@databricks.com>
This PR adds SQL aggregate functions with their tests for the KLL merge aggregate functions: - `kll_merge_agg_bigint` - `kll_merge_agg_float` - `kll_merge_agg_double` These aggregate functions merge multiple binary KLL sketch representations. Initial PRs: - #52900 - #52800 The existing scalar `kll_sketch_merge_*` functions can only merge two sketches at a time. In distributed computing scenarios where sketches are pre-computed across multiple partitions, time windows, or datasets, users need to merge many sketches together. Yes, this PR adds 3 new aggregate functions. New SQL tests were added to `sql/core/src/test/resources/sql-tests/inputs/kllquantiles.sql`: **Positive tests:** - Merging bigint/float/double sketches from multiple rows - Merging with custom k parameters (400, 300, 500) - NULL value handling **Negative tests:** - Type mismatches (passing non-binary types) - Invalid binary data - k parameter validation (too small, too large, NULL, non-constant) claude-4.5-sonnet and manual changes. Closes#53548 from cboumalh/cboumalh-kll-enhancement. Lead-authored-by: Chris Boumalhab <cboumalh@amazon.com> Co-authored-by: Chris Boumalhab <84485659+cboumalh@users.noreply.github.com> Signed-off-by: Daniel Tenedorio <daniel.tenedorio@databricks.com> (cherry picked from commit fc15f72) Signed-off-by: Daniel Tenedorio <daniel.tenedorio@databricks.com>

What changes were proposed in this pull request?
This PR adds DataFrame API support for the KLL quantile sketch functions that were previously added to Spark SQL in #52800. This lets users leverage KLL sketches through both Scala and Python DataFrame APIs in addition to the existing SQL interface.
Key additions:
Scala DataFrame API (
sql/api/src/main/scala/org/apache/spark/sql/functions.scala):Columnparameters for computed valuesStringparameters for column namesIntparameters for literal k valuesPython DataFrame API (
python/pyspark/sql/functions/builtin.py):ColumnOrName,Optional[Union[int, Column]])Python Spark Connect Support (
python/pyspark/sql/connect/functions/builtin.py):Why are the changes needed?
While the SQL API for KLL sketches was previously added, DataFrame API support is essential for full usability. Without DataFrame API support, users would be forced to use SQL expressions via
expr()orselectExpr(), which is less ergonomic and type-safe.Does this PR introduce any user-facing change?
Yes, this PR adds DataFrame API support for the 18 KLL sketch functions:
Scala DataFrame API Example:
Python DataFrame API Example:
How was this patch tested?
Scala Unit Tests (
DataFrameAggregateSuite):kll_sketch_agg_{bigint,float,double}with default and explicit k valueskll_sketch_to_stringfunctions for all data typeskll_sketch_get_nfunctions for all data typeskll_sketch_mergeoperationskll_sketch_get_quantilewith single rank and array of rankskll_sketch_get_rankoperationsPython Unit Tests (
test_functions.py):Was this patch authored or co-authored using generative AI tooling?
Yes, IDE assistance used
claude-4.5-sonnetwith manual validation and integration.