Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
16 changes: 16 additions & 0 deletions benchmarks/bench.sh
Original file line number Diff line number Diff line change
Expand Up @@ -177,6 +177,7 @@ hj: Benchmark for simple hash joins, testing various join sc
smj: Benchmark for simple sort merge joins, testing various join scenarios
dict: Benchmark for dictionary-encoded group-by scenarios
array_agg_distinct: 1000K-group, two-row-per-group array_agg(DISTINCT) benchmark
grouped_count_distinct: Grouped count(DISTINCT) by key type, group count and distinct cardinality
compile_profile: Compile and execute TPC-H across selected Cargo profiles, reporting timing and binary size


Expand Down Expand Up @@ -295,6 +296,9 @@ main() {
# Data is generated inline by the suite's load SQL.
echo "projection_subquery: no external data to generate"
;;
grouped_count_distinct)
echo "grouped_count_distinct: no external data to generate"
;;
asof_join)
data_asof_join
;;
Expand Down Expand Up @@ -697,6 +701,9 @@ main() {
array_agg_distinct)
run_array_agg_distinct
;;
grouped_count_distinct)
run_grouped_count_distinct
;;
compile_profile)
run_compile_profile "${PROFILE_ARGS[@]}"
;;
Expand Down Expand Up @@ -1831,6 +1838,15 @@ run_dict() {
debug_run $CARGO_COMMAND --bin dfbench -- dict --iterations 5 -o "${RESULTS_FILE}" ${QUERY_ARG} ${LATENCY_ARG}
}

# Runs grouped DISTINCT counts over inline data with varying key types and cardinalities.
run_grouped_count_distinct() {
debug_run env BENCH_NAME=grouped_count_distinct \
BENCH_RESULTS_FILE="$(sql_results_file grouped_count_distinct)" \
GROUPED_DISTINCT_ROWS="${GROUPED_DISTINCT_ROWS:-4000000}" \
${QUERY:+BENCH_QUERY="${QUERY}"} \
bash -c "$SQL_CARGO_COMMAND"
}

# Runs the data-free high-cardinality array_agg(DISTINCT) SQL benchmark.
run_array_agg_distinct() {
echo "Running array_agg_distinct benchmark..."
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,25 @@
name Q01
group grouped_count_distinct

load sql_benchmarks/grouped_count_distinct/init/load.sql

assert I
WITH actual AS (
SELECT g AS g, count(DISTINCT x) AS n FROM grouped_distinct GROUP BY g
), expected AS (
SELECT g, count(x) AS n
FROM (SELECT DISTINCT g AS g, x AS x FROM grouped_distinct)
GROUP BY g
)
SELECT NOT EXISTS (SELECT * FROM actual EXCEPT SELECT * FROM expected)
AND NOT EXISTS (SELECT * FROM expected EXCEPT SELECT * FROM actual);
----
true

expect_plan AggregateExec

run
-- Int64, 2,000 groups, high distinct cardinality (issue #25936).
SELECT g, count(DISTINCT x) FROM grouped_distinct GROUP BY g;

cleanup sql_benchmarks/grouped_count_distinct/init/cleanup.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1,25 @@
name Q02
group grouped_count_distinct

load sql_benchmarks/grouped_count_distinct/init/load.sql

assert I
WITH actual AS (
SELECT g AS g, count(DISTINCT low_x) AS n FROM grouped_distinct GROUP BY g
), expected AS (
SELECT g, count(x) AS n
FROM (SELECT DISTINCT g AS g, low_x AS x FROM grouped_distinct)
GROUP BY g
)
SELECT NOT EXISTS (SELECT * FROM actual EXCEPT SELECT * FROM expected)
AND NOT EXISTS (SELECT * FROM expected EXCEPT SELECT * FROM actual);
----
true

expect_plan AggregateExec

run
-- Int64, 2,000 groups, at most eight distinct values per group.
SELECT g, count(DISTINCT low_x) FROM grouped_distinct GROUP BY g;

cleanup sql_benchmarks/grouped_count_distinct/init/cleanup.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1,25 @@
name Q03
group grouped_count_distinct

load sql_benchmarks/grouped_count_distinct/init/load.sql

assert I
WITH actual AS (
SELECT few_g AS g, count(DISTINCT x) AS n FROM grouped_distinct GROUP BY few_g
), expected AS (
SELECT g, count(x) AS n
FROM (SELECT DISTINCT few_g AS g, x AS x FROM grouped_distinct)
GROUP BY g
)
SELECT NOT EXISTS (SELECT * FROM actual EXCEPT SELECT * FROM expected)
AND NOT EXISTS (SELECT * FROM expected EXCEPT SELECT * FROM actual);
----
true

expect_plan AggregateExec

run
-- Int64, eight groups, high distinct cardinality.
SELECT few_g, count(DISTINCT x) FROM grouped_distinct GROUP BY few_g;

cleanup sql_benchmarks/grouped_count_distinct/init/cleanup.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1,25 @@
name Q04
group grouped_count_distinct

load sql_benchmarks/grouped_count_distinct/init/load.sql

assert I
WITH actual AS (
SELECT g AS g, count(DISTINCT unsigned_x) AS n FROM grouped_distinct GROUP BY g
), expected AS (
SELECT g, count(x) AS n
FROM (SELECT DISTINCT g AS g, unsigned_x AS x FROM grouped_distinct)
GROUP BY g
)
SELECT NOT EXISTS (SELECT * FROM actual EXCEPT SELECT * FROM expected)
AND NOT EXISTS (SELECT * FROM expected EXCEPT SELECT * FROM actual);
----
true

expect_plan AggregateExec

run
-- UInt64, 2,000 groups, high distinct cardinality.
SELECT g, count(DISTINCT unsigned_x) FROM grouped_distinct GROUP BY g;

cleanup sql_benchmarks/grouped_count_distinct/init/cleanup.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1,25 @@
name Q05
group grouped_count_distinct

load sql_benchmarks/grouped_count_distinct/init/load.sql

assert I
WITH actual AS (
SELECT g AS g, count(DISTINCT string_x) AS n FROM grouped_distinct GROUP BY g
), expected AS (
SELECT g, count(x) AS n
FROM (SELECT DISTINCT g AS g, string_x AS x FROM grouped_distinct)
GROUP BY g
)
SELECT NOT EXISTS (SELECT * FROM actual EXCEPT SELECT * FROM expected)
AND NOT EXISTS (SELECT * FROM expected EXCEPT SELECT * FROM actual);
----
true

expect_plan AggregateExec

run
-- VARCHAR control: keep the existing two-phase rewrite.
SELECT g, count(DISTINCT string_x) FROM grouped_distinct GROUP BY g;

cleanup sql_benchmarks/grouped_count_distinct/init/cleanup.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1,25 @@
name Q06
group grouped_count_distinct

load sql_benchmarks/grouped_count_distinct/init/load.sql

assert I
WITH actual AS (
SELECT g AS g, count(DISTINCT float_x) AS n FROM grouped_distinct GROUP BY g
), expected AS (
SELECT g, count(x) AS n
FROM (SELECT DISTINCT g AS g, float_x AS x FROM grouped_distinct)
GROUP BY g
)
SELECT NOT EXISTS (SELECT * FROM actual EXCEPT SELECT * FROM expected)
AND NOT EXISTS (SELECT * FROM expected EXCEPT SELECT * FROM actual);
----
true

expect_plan AggregateExec

run
-- Float64 control: no specialized DISTINCT GroupsAccumulator.
SELECT g, count(DISTINCT float_x) FROM grouped_distinct GROUP BY g;

cleanup sql_benchmarks/grouped_count_distinct/init/cleanup.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1,11 @@
description = "Grouped COUNT(DISTINCT) across integer, string and floating-point keys, group counts and distinct cardinalities"

[[options]]
name = "rows"
env = "GROUPED_DISTINCT_ROWS"
default = "4000000"
help = "Input rows generated outside the timed query."

[[examples]]
command = "cargo run --release --bin benchmark_runner -- grouped_count_distinct --iterations 10 --partitions 1"
description = "Compare grouped DISTINCT counts at the issue's four-million-row scale."
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
DROP TABLE grouped_distinct;
Original file line number Diff line number Diff line change
@@ -0,0 +1,9 @@
CREATE TABLE grouped_distinct AS
SELECT value % 2000 AS g,
value % 8 AS few_g,
(value * 48271) % 999983 AS x,
(value / 2000) % 8 AS low_x,
arrow_cast((value * 48271) % 999983, 'UInt64') AS unsigned_x,
CAST((value * 48271) % 999983 AS VARCHAR) AS string_x,
CAST((value * 48271) % 999983 AS DOUBLE) AS float_x
FROM range(${GROUPED_DISTINCT_ROWS:-4000000});
47 changes: 37 additions & 10 deletions datafusion/optimizer/src/single_distinct_to_groupby.rs
Original file line number Diff line number Diff line change
Expand Up @@ -125,6 +125,7 @@ fn is_single_distinct_agg(
aggr_expr: &[Expr],
input_schema: &DFSchema,
count_rollup: Option<&CountRollup>,
has_group_by: bool,
) -> Result<bool> {
let mut fields_set = HashSet::new();
let mut aggregate_count = 0;
Expand Down Expand Up @@ -167,12 +168,39 @@ fn is_single_distinct_agg(
if aggregate_count != aggr_expr.len() || fields_set.len() != 1 {
return Ok(false);
}
// Grouped DISTINCT aggregates with native GroupsAccumulators do not need
// the adapter this rewrite avoids. Keep them on the direct path instead of
// building an extra aggregate with a row per (group, distinct value) pair.
// Leave mixed aggregates and global DISTINCT aggregation unchanged.
if has_group_by
&& distinct_aggs.len() == aggregate_count
&& all_have_groups_accumulators(&distinct_aggs, input_schema)?
{
return Ok(false);
}
if has_count_rollup && !rewrite_pays_for_count(&distinct_aggs, input_schema)? {
return Ok(false);
}
Ok(true)
}

fn all_have_groups_accumulators(
distinct_aggs: &[(&Arc<AggregateUDF>, &[Expr])],
input_schema: &DFSchema,
) -> Result<bool> {
for (func, args) in distinct_aggs {
let arg_types = args
.iter()
.map(|arg| arg.get_type(input_schema))
.collect::<Result<Vec<_>>>()?;
// An unknown answer must preserve the existing rewrite.
if func.groups_accumulator_supported_for_types(&arg_types, true) != Some(true) {
return Ok(false);
}
}
Ok(true)
}

/// Whether the rewrite is worth extending to a plan that only qualifies because
/// of the non-distinct `count`.
///
Expand Down Expand Up @@ -248,6 +276,7 @@ impl OptimizerRule for SingleDistinctToGroupBy {
&aggr_expr,
input.schema(),
count_rollup.as_ref(),
!group_expr.is_empty(),
)? && !contains_grouping_set(&group_expr) =>
{
let group_size = group_expr.len();
Expand Down Expand Up @@ -737,14 +766,12 @@ mod tests {
.aggregate(vec![col("a")], vec![count_distinct(col("b"))])?
.build()?;

// Should work
// The integer DISTINCT count already has a GroupsAccumulator.
assert_optimized_plan_equal!(
plan,
@r"
Projection: test.a, count(alias1) AS count(DISTINCT test.b) [a:UInt32, count(DISTINCT test.b):Int64]
Aggregate: groupBy=[[test.a]], aggr=[[count(alias1)]] [a:UInt32, count(alias1):Int64]
Aggregate: groupBy=[[test.a, test.b AS alias1]], aggr=[[]] [a:UInt32, alias1:UInt32]
TableScan: test [a:UInt32, b:UInt32, c:UInt32]
Aggregate: groupBy=[[test.a]], aggr=[[count(DISTINCT test.b)]] [a:UInt32, count(DISTINCT test.b):Int64]
TableScan: test [a:UInt32, b:UInt32, c:UInt32]
"
)
}
Expand Down Expand Up @@ -914,20 +941,20 @@ mod tests {

#[test]
fn group_by_with_expr() -> Result<()> {
let table_scan = test_table_scan().unwrap();
let table_scan = test_table_scan_utf8_b()?;

let plan = LogicalPlanBuilder::from(table_scan)
.aggregate(vec![col("a") + lit(1)], vec![count_distinct(col("c"))])?
.aggregate(vec![col("a") + lit(1)], vec![count_distinct(col("b"))])?
.build()?;

// Should work
assert_optimized_plan_equal!(
plan,
@r"
Projection: group_alias_0 AS test.a + Int32(1), count(alias1) AS count(DISTINCT test.c) [test.a + Int32(1):Int64, count(DISTINCT test.c):Int64]
Projection: group_alias_0 AS test.a + Int32(1), count(alias1) AS count(DISTINCT test.b) [test.a + Int32(1):Int64, count(DISTINCT test.b):Int64]
Aggregate: groupBy=[[group_alias_0]], aggr=[[count(alias1)]] [group_alias_0:Int64, count(alias1):Int64]
Aggregate: groupBy=[[test.a + Int32(1) AS group_alias_0, test.c AS alias1]], aggr=[[]] [group_alias_0:Int64, alias1:UInt32]
TableScan: test [a:UInt32, b:UInt32, c:UInt32]
Aggregate: groupBy=[[test.a + Int32(1) AS group_alias_0, test.b AS alias1]], aggr=[[]] [group_alias_0:Int64, alias1:Utf8]
TableScan: test [a:UInt32, b:Utf8, c:UInt32]
"
)
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -113,7 +113,7 @@ FROM (
)
----
<slt:ignore>
04)------AggregateExec: mode=Single, gby=[group_alias_0@0 as group_alias_0], aggr=[count(alias1)], metrics=[<slt:ignore>spill_count=7,<slt:ignore>]
04)------AggregateExec: mode=Single, gby=[v@0 * 7 % 100000 as t.v * Int64(7) % Int64(100000)], aggr=[count(DISTINCT t.v)], metrics=[<slt:ignore>spill_count=7,<slt:ignore>]
<slt:ignore>

# --- Case D: multiple aggregates (sum/min/max) under memory limit ---
Expand Down
22 changes: 8 additions & 14 deletions datafusion/sqllogictest/test_files/aggregates_simplify.slt
Original file line number Diff line number Diff line change
Expand Up @@ -543,23 +543,17 @@ query TT
EXPLAIN SELECT g, count(DISTINCT v) FROM distinct_simplify_t GROUP BY g;
----
logical_plan
01)Projection: distinct_simplify_t.g, count(alias1) AS count(DISTINCT distinct_simplify_t.v)
02)--Aggregate: groupBy=[[distinct_simplify_t.g]], aggr=[[count(alias1)]]
03)----Aggregate: groupBy=[[distinct_simplify_t.g, distinct_simplify_t.v AS alias1]], aggr=[[]]
04)------TableScan: distinct_simplify_t projection=[g, v]
01)Aggregate: groupBy=[[distinct_simplify_t.g]], aggr=[[count(DISTINCT distinct_simplify_t.v)]]
02)--TableScan: distinct_simplify_t projection=[g, v]
physical_plan
01)ProjectionExec: expr=[g@0 as g, count(alias1)@1 as count(DISTINCT distinct_simplify_t.v)]
02)--AggregateExec: mode=FinalPartitioned, gby=[g@0 as g], aggr=[count(alias1)]
03)----RepartitionExec: partitioning=Hash([g@0], 4), input_partitions=4
04)------AggregateExec: mode=Partial, gby=[g@0 as g], aggr=[count(alias1)]
05)--------AggregateExec: mode=FinalPartitioned, gby=[g@0 as g, alias1@1 as alias1], aggr=[]
06)----------RepartitionExec: partitioning=Hash([g@0, alias1@1], 4), input_partitions=1
07)------------AggregateExec: mode=Partial, gby=[g@0 as g, v@1 as alias1], aggr=[]
08)--------------DataSourceExec: partitions=1, partition_sizes=[1]
01)AggregateExec: mode=FinalPartitioned, gby=[g@0 as g], aggr=[count(DISTINCT distinct_simplify_t.v)]
02)--RepartitionExec: partitioning=Hash([g@0], 4), input_partitions=1
03)----AggregateExec: mode=Partial, gby=[g@0 as g], aggr=[count(DISTINCT distinct_simplify_t.v)]
04)------DataSourceExec: partitions=1, partition_sizes=[1]

# Mixed: a node that still needs one DISTINCT keeps all of them, so this rule
# cannot change which plans SingleDistinctToGroupBy rewrites. `count` needs the
# rewrite, and the plan below shows that it still happens.
# leaves min(DISTINCT) in place. The mixed min/count aggregate still gets the
# existing SingleDistinctToGroupBy rewrite.
query III
SELECT g, min(DISTINCT v), count(DISTINCT v) FROM distinct_simplify_t GROUP BY g ORDER BY g;
----
Expand Down
Loading
Loading