Skip to content

Commit 578e4e9

Browse files
committed
new correlation aggregation supported
1 parent b621631 commit 578e4e9

3 files changed

Lines changed: 135 additions & 0 deletions

File tree

‎crates/frontend-sql/src/sql/mod.rs‎

Lines changed: 39 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1602,6 +1602,13 @@ fn lower_agg_intent(expr: &Expr) -> Result<AggIntent<ColumnRef>, LoweringError>
16021602
if let Some(intent) = lower_arg_selector(&name, &agg_fn.args)? {
16031603
return Ok(intent);
16041604
}
1605+
// `corr` IS a native DataFusion aggregate — what it lacks is an
1606+
// `AggSemantic`, since no core `AggIntent` reads two columns.
1607+
// Handled here for the same reason as the arg selectors above:
1608+
// the `lookup_native` lookup below has nothing to return for it.
1609+
if let Some(intent) = lower_corr(&name, &agg_fn.args)? {
1610+
return Ok(intent);
1611+
}
16051612
let semantic = asap_sql_function_catalog::lookup_native(&name)
16061613
.ok_or_else(|| LoweringError::UnsupportedAggregate(name.clone()))?;
16071614
// The canonical intent algebra has no DISTINCT modifier for the
@@ -1776,6 +1783,38 @@ fn lower_arg_selector(
17761783
}))
17771784
}
17781785

1786+
/// Pearson correlation `corr(x, y)` — a two-column statistic.
1787+
///
1788+
/// Every core `AggIntent` reducer folds *one* column (`col: Option<C>`);
1789+
/// correlation reads two and would be the first binary core variant. Per
1790+
/// `AggIntent::Extension`'s own "core only grows for intents ≥2 deployment
1791+
/// models actually use" bar, and a search that found no second model wanting
1792+
/// it — PromQL has no correlation — this lowers to `Extension`, the same
1793+
/// judgement `lower_arg_selector` records for `argMax`/`argMin`.
1794+
///
1795+
/// Both columns are kept as validated bare-column `ColumnRef`s in `payload`
1796+
/// (`reducer_col`'s "no expression arguments" rule, issue #115); core carries
1797+
/// them without resolving them, since `Extension` has no typed column field.
1798+
/// The order is the order written: a later rewrite into co-moments has to tell
1799+
/// x from y.
1800+
fn lower_corr(name: &str, args: &[Expr]) -> Result<Option<AggIntent<ColumnRef>>, LoweringError> {
1801+
if name != "corr" {
1802+
return Ok(None);
1803+
}
1804+
let [x, y] = args else {
1805+
unreachable!(
1806+
"corr's DataFusion signature fixes its arity at 2 -- the planner already rejected any other argument count before lower_agg_intent runs"
1807+
);
1808+
};
1809+
Ok(Some(AggIntent::Extension {
1810+
ext_kind: "corr".to_string(),
1811+
payload: serde_json::json!({
1812+
"x_col": reducer_col(name, std::slice::from_ref(x))?,
1813+
"y_col": reducer_col(name, std::slice::from_ref(y))?,
1814+
}),
1815+
}))
1816+
}
1817+
17791818
/// The name an `IN (subquery)`'s key column is projected under, so the join
17801819
/// predicate cannot bind it to a same-named column of the outer relation.
17811820
const IN_SUBQUERY_KEY: &str = "__asap_in_key";

‎crates/frontend-sql/tests/sql_lowering.rs‎

Lines changed: 88 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2454,3 +2454,91 @@ async fn clickhouse_tuple_element_preserves_declared_field_metadata() {
24542454
);
24552455
}
24562456
}
2457+
2458+
// ── `corr(x, y)`: a two-column statistic with no core `AggIntent` shape.
2459+
// Every core reducer folds ONE column (`col: Option<C>`); correlation reads
2460+
// two and would be the first binary core variant. A repo-wide search turned up
2461+
// no second deployment model wanting it — PromQL has no correlation — so per
2462+
// `AggIntent::Extension`'s own "core only grows for intents ≥2 deployment
2463+
// models actually use" bar it lowers to an `Extension`, exactly as `argMax`
2464+
// does. Unlike `argMax` it IS a native DataFusion aggregate; what is missing is
2465+
// an `AggSemantic` for it, which is why `lookup_native` rejected it. ─────────
2466+
2467+
#[tokio::test]
2468+
async fn corr_lowers_to_an_extension_intent() {
2469+
let qe = lower("SELECT corr(latency, bytes) AS r FROM metrics").await;
2470+
let (by, measures) = find_aggregate(&qe).expect("expected an Aggregate");
2471+
assert!(by.is_empty());
2472+
assert!(
2473+
matches!(
2474+
measures.as_slice(),
2475+
[AggIntent::Extension { ext_kind, .. }] if ext_kind == "corr"
2476+
),
2477+
"expected Extension {{ ext_kind: \"corr\", .. }}, got {measures:?}"
2478+
);
2479+
}
2480+
2481+
#[tokio::test]
2482+
async fn corr_payload_preserves_both_column_names() {
2483+
// Core never resolves an `Extension`'s payload, so both columns stay as
2484+
// validated bare-column `ColumnRef`s — and in the order written, since
2485+
// a later rewrite into co-moments needs to tell the two apart.
2486+
let qe = lower("SELECT corr(latency, bytes) AS r FROM metrics").await;
2487+
let (_, measures) = find_aggregate(&qe).expect("expected an Aggregate");
2488+
let AggIntent::Extension { payload, .. } = &measures[0] else {
2489+
panic!("expected an Extension intent, got {:?}", measures[0]);
2490+
};
2491+
let named = |key: &str| {
2492+
payload
2493+
.get(key)
2494+
.and_then(|c| c.get("Named"))
2495+
.and_then(|n| n.as_str())
2496+
.map(str::to_string)
2497+
};
2498+
assert_eq!(named("x_col"), Some("latency".to_string()));
2499+
assert_eq!(named("y_col"), Some("bytes".to_string()));
2500+
}
2501+
2502+
#[tokio::test]
2503+
async fn corr_over_an_expression_binds_the_derived_column() {
2504+
// `reducer_col`'s "bare column only" rule (issue #115) exists so an
2505+
// expression argument is never silently DROPPED. Here nothing is dropped:
2506+
// `corr` is a native DataFusion aggregate, so the planner materializes
2507+
// `latency * 2` into the Project below the Aggregate and hands the
2508+
// aggregate a column reference to it. The payload names that derived
2509+
// column, which resolves in the child's own output schema — so the
2510+
// expression is carried, not lost.
2511+
//
2512+
// `argMax` behaves differently only because it reaches `lower_agg_intent`
2513+
// as a stub ClickHouse builtin, before that rewrite applies.
2514+
let qe = lower("SELECT corr(latency * 2, bytes) AS r FROM metrics").await;
2515+
let (_, measures) = find_aggregate(&qe).expect("expected an Aggregate");
2516+
let AggIntent::Extension { payload, .. } = &measures[0] else {
2517+
panic!("expected an Extension intent, got {:?}", measures[0]);
2518+
};
2519+
let x = payload
2520+
.get("x_col")
2521+
.and_then(|c| c.get("Named"))
2522+
.and_then(|n| n.as_str())
2523+
.expect("x_col is a Named ColumnRef");
2524+
assert!(
2525+
x.contains('*'),
2526+
"expected the derived column carrying `latency * 2`, got {x:?}"
2527+
);
2528+
}
2529+
2530+
#[tokio::test]
2531+
async fn corr_output_column_is_a_nullable_float() {
2532+
// `Extension`'s generic guess is `Utf8` — correct only for a shape core
2533+
// knows nothing about. Correlation always yields a float, and NULL when
2534+
// a variance is zero or fewer than two rows contributed, so the schema
2535+
// says so rather than carrying the placeholder downstream.
2536+
let qe = lower("SELECT corr(latency, bytes) AS r FROM metrics").await;
2537+
let schema = qe.output_schema().expect("output schema");
2538+
let out = schema.columns.last().expect("at least one output column");
2539+
assert_eq!(out.dtype, DataType::Float64, "got {:?}", schema.columns);
2540+
assert!(
2541+
out.nullable,
2542+
"corr is NULL on zero variance / fewer than 2 rows"
2543+
);
2544+
}

‎crates/types/src/pre_asap/agg_intent.rs‎

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -568,6 +568,14 @@ impl<C: Clone> AggIntent<C> {
568568
// the column after `kind` and leave the type unconstrained.
569569
// The owning deployment model is expected to re-derive the
570570
// real schema itself rather than rely on this generic guess.
571+
// Correlation always yields a float, and NULL when a variance is
572+
// zero or fewer than two rows contributed — unlike the arg
573+
// selectors, whose output type follows the selected column and is
574+
// patched in during aggregate schema derivation, this one needs no
575+
// schema and so is settled here.
576+
AggIntent::Extension { ext_kind, .. } if ext_kind == "corr" => {
577+
col("corr", DataType::Float64, true)
578+
}
571579
AggIntent::Extension { ext_kind, .. } => col(ext_kind, DataType::Utf8, true),
572580
}
573581
}

0 commit comments

Comments
 (0)