diff --git a/crates/asap-aware-mapping/src/replacement.rs b/crates/asap-aware-mapping/src/replacement.rs index b69df8a33..386bd028d 100644 --- a/crates/asap-aware-mapping/src/replacement.rs +++ b/crates/asap-aware-mapping/src/replacement.rs @@ -358,6 +358,7 @@ use asap_types::post_asap::{ }; use asap_types::post_asap::{AccuracyError, CompositionOperator, GuaranteeSource, ResultGuarantee}; use asap_types::pre_asap::agg_intent::{agg_is_mergeable, AggIntent}; +use asap_types::pre_asap::column_resolution::resolve_column_ref; use asap_types::pre_asap::cse::{share_common_subtrees, structural_hash, HashCache}; use asap_types::pre_asap::expr_ir::{ArithmeticOpKind, ColumnRef}; use asap_types::pre_asap::query_expr::any_measure_filtered; @@ -2867,6 +2868,15 @@ fn construct_summary_agg( { // State identity is independent of which statistic reads it. field.name = "univmon".into(); + } else if let (AggIntent::Quantile { .. }, SummaryInputExpr::Column(col)) = + (intent, &summary_input.weight) + { + // The quantile is a readout parameter: name the state after the + // column it summarizes, not after the query's output column. + let child_schema = input.child.output_schema()?; + if let Ok(i) = resolve_column_ref(col, &child_schema) { + field.name = child_schema.columns[i].name.clone(); + } } } @@ -9791,9 +9801,10 @@ mod tests { ); assert_eq!(input, &SummaryUpdate::column(ColumnRef::SampleValue)); assert_eq!(reduction, &ReductionTy::by(vec![2])); - // SummaryAgg edge: the state column carries the committed family. + // SummaryAgg edge: the state column, named after its input, carries + // the committed family. assert_eq!( - field(&summary_input.schema, "quantile_0_99").dtype, + field(&summary_input.schema, "value").dtype, SummaryFamilyType::Sketch( SketchKind::new(SketchAlgorithm::Kll, SketchParams::Kll { k: 269 }), GroupingStrategy::default() diff --git a/crates/integration-tests/tests/promql_to_post_asap.rs b/crates/integration-tests/tests/promql_to_post_asap.rs index 18ebd9c76..1a8d364e0 100644 --- a/crates/integration-tests/tests/promql_to_post_asap.rs +++ b/crates/integration-tests/tests/promql_to_post_asap.rs @@ -887,7 +887,7 @@ fn planner_heap_topk_reference_execution_matches_ground_truth() { /// /// ```text /// SummaryEstimate { query: Quantile{0.99} } → {quantile_0_99: Float64} -/// └─ SummaryAgg { Kll{k:269}, input: SampleValue } → {quantile_0_99: Sketch(Kll, {k:269})} +/// └─ SummaryAgg { Kll{k:269}, input: SampleValue } → {value: Sketch(Kll, {k:269})} /// └─ SummaryAgg { Rate, input: SampleValue } → {ts, value: ExactAggregate(Rate), …} /// └─ KeepPreAsap(TimeRange{5m} → Scan) → {ts, value} /// ``` @@ -948,7 +948,7 @@ fn promql_quantile_of_rate_binds_kll_over_rate_accumulator() { "global quantile — no group keys, full reduction" ); assert_eq!( - dtype(&summary_input.schema, "quantile_0_99"), + dtype(&summary_input.schema, "value"), &SummaryFamilyType::Sketch( SketchKind::new(SketchAlgorithm::Kll, SketchParams::Kll { k: 269 }), GroupingStrategy::default() diff --git a/crates/planner/tests/summary_sharing.rs b/crates/planner/tests/summary_sharing.rs index 2346537e7..c344296bf 100644 --- a/crates/planner/tests/summary_sharing.rs +++ b/crates/planner/tests/summary_sharing.rs @@ -267,6 +267,19 @@ async fn different_producers_are_not_shared() { } } +/// Cross-series quantiles name their KLL state after the input column, not the +/// quantile, so p50 and p99 over one selector share it. +#[tokio::test] +async fn cross_series_p50_and_p99_share_one_producer() { + let output = plan_promql( + &[("quantile(0.5, lat)", 0.01), ("quantile(0.99, lat)", 0.01)], + &CHEAP_SUMMARY, + ) + .await; + assert!(same_states(&states(&output))); + assert_eq!(unique_deployments(&output), 1); +} + /// An ungrouped aggregate has no unique key, so pre-ASAP CSE keeps the two /// copies apart; their identical producers (rate, then sum) are shared here. #[tokio::test] @@ -278,8 +291,7 @@ async fn identical_ungrouped_queries_share_their_producers() { } /// The SQL frontend reaches the same sharing for two copies of one filtered -/// percentile. (Its state schema names the query's output column, so p50 and -/// p99 over one SQL sketch are not yet structurally identical.) +/// percentile. #[tokio::test] async fn identical_sql_percentiles_share_one_producer() { let query = @@ -289,6 +301,55 @@ async fn identical_sql_percentiles_share_one_producer() { assert_eq!(unique_deployments(&output), 1); } +/// The quantile is a readout parameter: SQL p50 and p99 over one filtered +/// column build one KLL, named after its input, while each query keeps its +/// own output column. +#[tokio::test] +async fn sql_p50_and_p99_share_one_producer() { + let p50 = + "SELECT approx_percentile_cont(l_extendedprice, 0.5) FROM lineitem WHERE l_orderkey > 10"; + let p99 = + "SELECT approx_percentile_cont(l_extendedprice, 0.99) FROM lineitem WHERE l_orderkey > 10"; + let output = plan_sql(&[p50, p99], &CHEAP_SUMMARY).await; + assert!(same_states(&states(&output))); + assert_eq!(unique_deployments(&output), 1); + let names: Vec<_> = output + .plans + .iter() + .map(|plan| { + plan.plan + .root + .schema + .fields + .iter() + .map(|field| field.name.clone()) + .collect::>() + }) + .collect(); + assert_eq!( + names, + [ + ["approx_percentile_cont(lineitem.l_extendedprice,Float64(0.5))"], + ["approx_percentile_cont(lineitem.l_extendedprice,Float64(0.99))"], + ] + ); + + for queries in [ + [ + p50, + "SELECT approx_percentile_cont(l_extendedprice, 0.99) FROM lineitem WHERE l_orderkey > 20", + ], + [ + p50, + "SELECT approx_percentile_cont(l_orderkey, 0.99) FROM lineitem WHERE l_orderkey > 10", + ], + ] { + let output = plan_sql(&queries, &CHEAP_SUMMARY).await; + assert!(!same_states(&states(&output)), "{queries:?}"); + assert_eq!(unique_deployments(&output), 2, "{queries:?}"); + } +} + /// A state costs 100 and recomputing a query costs 60 over its six reads: /// alone, the query recomputes raw. Shared by p50 and p99, the state costs 50 /// per query, so both keep it.