Skip to content
Merged
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
15 changes: 13 additions & 2 deletions crates/asap-aware-mapping/src/replacement.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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();
}
}
}

Expand Down Expand Up @@ -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()
Expand Down
4 changes: 2 additions & 2 deletions crates/integration-tests/tests/promql_to_post_asap.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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}
/// ```
Expand Down Expand Up @@ -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()
Expand Down
65 changes: 63 additions & 2 deletions crates/planner/tests/summary_sharing.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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]
Expand All @@ -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 =
Expand All @@ -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::<Vec<_>>()
})
.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.
Expand Down
Loading