From e9beb62eef520520171e083b03ec2e769635b0ca Mon Sep 17 00:00:00 2001 From: discord9 Date: Thu, 17 Sep 2026 10:08:54 +0000 Subject: [PATCH] perf(promql): avoid concatenating constant series tags (#9108) * perf(promql): experiment with constant-tag series concat Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com> * test(promql): verify logical constant-tag concat equivalence Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com> * test(perf): qualify constant-tag series concat Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com> * test(perf): cover fragmented millisecond series concat Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com> * fix(promql): compact constant dictionary tags at construction Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com> * perf(promql): construct constant string dictionaries directly Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com> * fix(perf): benchmark ordinary TQL queries for constant tags Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com> * fix(promql): address constant-tag review and cardinality coverage Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com> * fix(promql): scope concat optimization to string dictionaries Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com> * test(perf): move high-cardinality constant-tag cases to the heavy set The 10k and 100k constant-tag direct-SST cases repeatedly kill the self-hosted query-regression runner (lost communication during the run), while the default-cardinality case passes. Move them out of the default 'all' set into the heavy set so they only run on demand (case=heavy or the heavy-regression label), and qualify them locally instead. Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com> --------- Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com> --- .github/scripts/query-regression-run.py | 5 + .../src/extension_plan/series_divide.rs | 640 +++++++++++++++++- .../promql_constant_tag_concat_ms/case.toml | 139 ++++ .../case.toml | 138 ++++ .../case.toml | 138 ++++ .../test_query_regression_case_selection.py | 8 +- 6 files changed, 1059 insertions(+), 9 deletions(-) create mode 100644 tests/perf/query_cases/promql_constant_tag_concat_ms/case.toml create mode 100644 tests/perf/query_cases/promql_constant_tag_concat_ms_100k/case.toml create mode 100644 tests/perf/query_cases/promql_constant_tag_concat_ms_10k/case.toml diff --git a/.github/scripts/query-regression-run.py b/.github/scripts/query-regression-run.py index 6da7c896d6..d488cee159 100644 --- a/.github/scripts/query-regression-run.py +++ b/.github/scripts/query-regression-run.py @@ -45,6 +45,11 @@ DEFAULT_CASES = [ HEAVY_CASES = [ "tests/perf/query_cases/prom_remote_write_7913/case.toml", + # High-cardinality direct-SST cases exceed the default query-regression + # runner's resources (the self-hosted runner is repeatedly lost during the + # run). Keep them in the heavy set so they only run on demand. + "tests/perf/query_cases/promql_constant_tag_concat_ms_10k/case.toml", + "tests/perf/query_cases/promql_constant_tag_concat_ms_100k/case.toml", ] CASE_GROUPS = { diff --git a/src/promql/src/extension_plan/series_divide.rs b/src/promql/src/extension_plan/series_divide.rs index b850d9b2dd..bf62f2c333 100644 --- a/src/promql/src/extension_plan/series_divide.rs +++ b/src/promql/src/extension_plan/series_divide.rs @@ -16,8 +16,13 @@ use std::pin::Pin; use std::sync::Arc; use std::task::{Context, Poll}; -use datafusion::arrow::array::{Array, ArrayRef, UInt64Array}; -use datafusion::arrow::datatypes::{DataType, SchemaRef}; +use datafusion::arrow::array::{ + Array, ArrayRef, DictionaryArray, LargeStringArray, PrimitiveArray, StringArray, + StringViewArray, UInt64Array, +}; +use datafusion::arrow::buffer::NullBuffer; +use datafusion::arrow::datatypes::{ArrowDictionaryKeyType, DataType, SchemaRef}; +use datafusion::arrow::downcast_dictionary_array; use datafusion::arrow::record_batch::RecordBatch; use datafusion::common::tree_node::TreeNodeRecursion; use datafusion::common::{DFSchema, DFSchemaRef}; @@ -493,6 +498,87 @@ impl RecordBatchStream for SeriesDivideStream { } } +fn constant_string_dictionary( + dictionary: &DictionaryArray, + total_rows: usize, +) -> DataFusionResult { + let key = dictionary.key(0); + let value = key.and_then(|key| string_array_value_at_index(dictionary.values(), key)); + let values: ArrayRef = match dictionary.values().data_type() { + DataType::Utf8 => Arc::new(StringArray::from(vec![value])), + DataType::LargeUtf8 => Arc::new(LargeStringArray::from(vec![value])), + DataType::Utf8View => Arc::new(StringViewArray::from(vec![value])), + _ => unreachable!("dictionary values must be strings"), + }; + let keys = PrimitiveArray::::new( + vec![K::Native::default(); total_rows].into(), + key.is_none().then(|| NullBuffer::new_null(total_rows)), + ); + + Ok(Arc::new(DictionaryArray::try_new(keys, values)?)) +} + +/// Concatenates batches from one isolated series; every designated tag is constant. +fn concat_series_batches( + schema: &SchemaRef, + batches: &[RecordBatch], + tag_indices: &[usize], +) -> DataFusionResult { + if batches.len() <= 1 || tag_indices.is_empty() { + return Ok(compute::concat_batches(schema, batches)?); + } + + let Some(first_batch) = batches.iter().find(|batch| batch.num_rows() > 0) else { + return Ok(compute::concat_batches(schema, batches)?); + }; + + // This endpoint-only sanity check does not validate tags in interior rows. + #[cfg(debug_assertions)] + { + let last_batch = batches + .iter() + .rfind(|batch| batch.num_rows() > 0) + .expect("first non-empty batch implies a last non-empty batch"); + if let (Ok(first_tags), Ok(last_tags)) = ( + TagIdentifier::try_new(first_batch, tag_indices), + TagIdentifier::try_new(last_batch, tag_indices), + ) { + debug_assert!( + first_tags.equal_at(0, &last_tags, last_batch.num_rows() - 1)?, + "series batch tag endpoints must match" + ); + } + } + + let total_rows: usize = batches.iter().map(RecordBatch::num_rows).sum(); + let columns = schema + .fields() + .iter() + .enumerate() + .map(|(index, field)| -> DataFusionResult { + if tag_indices.contains(&index) + && matches!(field.data_type(), DataType::Dictionary(_, value_type) if value_type.is_string()) + { + let array = first_batch.column(index); + downcast_dictionary_array! { + array => constant_string_dictionary(array, total_rows), + _ => unreachable!("dictionary keys must be integers"), + } + } else { + compute::concat( + &batches + .iter() + .map(|batch| batch.column(index).as_ref()) + .collect::>(), + ) + .map_err(Into::into) + } + }) + .collect::>>()?; + + RecordBatch::try_new(schema.clone(), columns).map_err(Into::into) +} + impl Stream for SeriesDivideStream { type Item = DataFusionResult; @@ -522,7 +608,8 @@ impl Stream for SeriesDivideStream { } else { self.buffer.remove(0); } - let result_batch = compute::concat_batches(&self.schema, &result_batches)?; + let result_batch = + concat_series_batches(&self.schema, &result_batches, &self.tag_indices)?; self.inspect_start = 0; self.num_series.add(1); @@ -540,7 +627,8 @@ impl Stream for SeriesDivideStream { continue; } else { // input stream is ended - let result = compute::concat_batches(&self.schema, &self.buffer)?; + let result = + concat_series_batches(&self.schema, &self.buffer, &self.tag_indices)?; self.buffer.clear(); self.inspect_start = 0; self.num_series.add(1); @@ -628,11 +716,16 @@ impl SeriesDivideStream { #[cfg(test)] mod test { + use std::collections::HashMap; + use datafusion::arrow::array::{ - DictionaryArray, Int32Array, Int64Array, LargeStringArray, StringArray, StringViewArray, - UInt32Array, + DictionaryArray, Int32Array, Int64Array, LargeStringArray, PrimitiveArray, StringArray, + StringViewArray, UInt32Array, UInt64Array, + }; + use datafusion::arrow::datatypes::{ + ArrowDictionaryKeyType, DataType, Field, Int8Type, Int16Type, Int32Type, Int64Type, Schema, + UInt8Type, UInt16Type, UInt32Type, UInt64Type, }; - use datafusion::arrow::datatypes::{DataType, Field, Int32Type, Schema, UInt32Type}; use datafusion::common::ToDFSchema; use datafusion::datasource::memory::MemorySourceConfig; use datafusion::datasource::source::DataSourceExec; @@ -641,6 +734,446 @@ mod test { use super::*; + fn assert_concat_matches_reference( + schema: SchemaRef, + batches: Vec, + tags: &[usize], + ) { + let expected = compute::concat_batches(&schema, &batches).unwrap(); + let actual = concat_series_batches(&schema, &batches, tags).unwrap(); + assert_eq!(actual.schema(), schema); + assert_eq!(actual.num_rows(), expected.num_rows()); + + for (index, (actual, expected)) in + actual.columns().iter().zip(expected.columns()).enumerate() + { + assert_eq!(actual.data_type(), expected.data_type(), "column {index}"); + let (actual, expected) = match actual.data_type() { + DataType::Dictionary(_, value_type) => ( + compute::cast(actual.as_ref(), value_type).unwrap(), + compute::cast(expected.as_ref(), value_type).unwrap(), + ), + _ => (actual.clone(), expected.clone()), + }; + assert_eq!(actual.to_data(), expected.to_data(), "column {index}"); + } + } + + #[test] + fn test_concat_series_batches() { + let schema = Arc::new(Schema::new(vec![ + Field::new("tag", DataType::Utf8, true), + Field::new("value", DataType::Int64, true), + ])); + let batches = vec![ + RecordBatch::try_new( + schema.clone(), + vec![ + Arc::new(StringArray::from(vec![Some("tag"), Some("tag")])), + Arc::new(Int64Array::from(vec![Some(1), None])), + ], + ) + .unwrap(), + RecordBatch::try_new( + schema.clone(), + vec![ + Arc::new(StringArray::from(vec![Some("tag"), Some("tag")])), + Arc::new(Int64Array::from(vec![Some(3), Some(4)])), + ], + ) + .unwrap(), + ]; + assert_concat_matches_reference(schema, batches, &[0]); + } + + #[cfg(debug_assertions)] + #[test] + #[should_panic(expected = "series batch tag endpoints must match")] + fn test_concat_series_batches_mismatched_tag_endpoints_panics() { + let schema = Arc::new(Schema::new(vec![ + Field::new("tag", DataType::Utf8, true), + Field::new("value", DataType::Int64, true), + ])); + let batches = vec![ + RecordBatch::try_new( + schema.clone(), + vec![ + Arc::new(StringArray::from(vec!["first"])), + Arc::new(Int64Array::from(vec![1])), + ], + ) + .unwrap(), + RecordBatch::try_new( + schema.clone(), + vec![ + Arc::new(StringArray::from(vec!["last"])), + Arc::new(Int64Array::from(vec![2])), + ], + ) + .unwrap(), + ]; + + concat_series_batches(&schema, &batches, &[0]).unwrap(); + } + + #[test] + fn test_concat_series_batches_string_tags() { + for (data_type, batches) in [ + ( + DataType::LargeUtf8, + vec![ + Arc::new(LargeStringArray::from(vec![ + Some("long tag"), + Some("long tag"), + ])) as ArrayRef, + Arc::new(LargeStringArray::from(vec![Some("long tag")])) as ArrayRef, + ], + ), + ( + DataType::Utf8View, + vec![ + Arc::new(StringViewArray::from(vec![ + Some("view tag longer than twelve bytes"), + Some("view tag longer than twelve bytes"), + ])) as ArrayRef, + Arc::new(StringViewArray::from(vec![Some( + "view tag longer than twelve bytes", + )])) as ArrayRef, + ], + ), + ( + DataType::Utf8, + vec![ + Arc::new(StringArray::from(vec![None::<&str>, None])) as ArrayRef, + Arc::new(StringArray::from(vec![None::<&str>])) as ArrayRef, + ], + ), + ( + DataType::Utf8, + vec![ + Arc::new(StringArray::from(vec![Some(""), Some("")])) as ArrayRef, + Arc::new(StringArray::from(vec![Some("")])) as ArrayRef, + ], + ), + ] { + let schema = Arc::new(Schema::new(vec![ + Field::new("tag", data_type, true), + Field::new("value", DataType::Int64, true), + ])); + let batches = batches + .into_iter() + .enumerate() + .map(|(index, tag)| { + RecordBatch::try_new( + schema.clone(), + vec![ + tag.clone(), + Arc::new(Int64Array::from_iter_values( + (0..tag.len()).map(|row| (index * 10 + row) as i64), + )), + ], + ) + .unwrap() + }) + .collect(); + assert_concat_matches_reference(schema, batches, &[0]); + } + } + + fn assert_dictionary_concat( + keys: Vec>, + values: ArrayRef, + other_keys: Vec>, + other_values: ArrayRef, + ) where + PrimitiveArray: From>>, + { + let data_type = DataType::Dictionary( + Box::new(PrimitiveArray::::from(keys.clone()).data_type().clone()), + Box::new(values.data_type().clone()), + ); + let schema = Arc::new(Schema::new(vec![ + Field::new("tag", data_type, true), + Field::new("value", DataType::Int64, true), + ])); + let first = Arc::new(DictionaryArray::::new( + PrimitiveArray::::from(keys), + values, + )) as ArrayRef; + let second = Arc::new(DictionaryArray::::new( + PrimitiveArray::::from(other_keys), + other_values, + )) as ArrayRef; + let batches = [first, second] + .into_iter() + .enumerate() + .map(|(index, tag)| { + RecordBatch::try_new( + schema.clone(), + vec![ + tag.clone(), + Arc::new(Int64Array::from_iter_values( + (0..tag.len()).map(|row| (index * 10 + row) as i64), + )), + ], + ) + .unwrap() + }) + .collect(); + assert_concat_matches_reference(schema, batches, &[0]); + } + + #[test] + fn test_concat_series_batches_dictionary_tags() { + macro_rules! dictionary_cases { + ($($key_type:ty, $key:expr),+ $(,)?) => { + $(for (values, other_values) in [ + ( + Arc::new(StringArray::from(vec!["other", "tag"])) as ArrayRef, + Arc::new(StringArray::from(vec!["tag", "other"])) as ArrayRef, + ), + ( + Arc::new(LargeStringArray::from(vec!["other", "tag"])) as ArrayRef, + Arc::new(LargeStringArray::from(vec!["tag", "other"])) as ArrayRef, + ), + ( + Arc::new(StringViewArray::from(vec![ + "other", + "view tag longer than twelve bytes", + ])) as ArrayRef, + Arc::new(StringViewArray::from(vec![ + "view tag longer than twelve bytes", + "other", + ])) as ArrayRef, + ), + ] { + assert_dictionary_concat::<$key_type>( + vec![Some(($key)(1)), Some(($key)(1))], + values, + vec![Some(($key)(0))], + other_values, + ); + })+ + }; + } + dictionary_cases!( + Int8Type, + |value| value as i8, + Int16Type, + |value| value as i16, + Int32Type, + |value: i32| value, + Int64Type, + |value| value as i64, + UInt8Type, + |value| value as u8, + UInt16Type, + |value| value as u16, + UInt32Type, + |value| value as u32, + UInt64Type, + |value| value as u64, + ); + assert_dictionary_concat::( + vec![Some(0), Some(0)], + Arc::new(LargeStringArray::from(vec![None::<&str>])), + vec![None], + Arc::new(LargeStringArray::from(vec!["unused"])), + ); + // Compare logical values rather than dictionary keys for the reverse null orientation. + assert_dictionary_concat::( + vec![None], + Arc::new(StringArray::from(vec!["unused"])), + vec![Some(0)], + Arc::new(StringArray::from(vec![None::<&str>])), + ); + } + + #[test] + fn test_concat_series_batches_dictionary_tag_drops_unused_values() { + let active = "active view tag longer than twelve bytes"; + let unused = "unused view tag ".repeat(1024); + let values = Arc::new(StringViewArray::from(vec![unused.as_str(), active])) as ArrayRef; + let schema = Arc::new(Schema::new(vec![ + Field::new( + "tag", + DataType::Dictionary(Box::new(DataType::UInt32), Box::new(DataType::Utf8View)), + false, + ), + Field::new("value", DataType::Int64, false), + ])); + let batches = [vec![1, 1], vec![1]] + .into_iter() + .enumerate() + .map(|(batch, keys)| { + let num_rows = keys.len(); + RecordBatch::try_new( + schema.clone(), + vec![ + Arc::new(DictionaryArray::::new( + UInt32Array::from(keys), + values.clone(), + )), + Arc::new(Int64Array::from_iter_values( + (0..num_rows).map(|row| (batch * 10 + row) as i64), + )), + ], + ) + .unwrap() + }) + .collect::>(); + + assert_concat_matches_reference(schema.clone(), batches.clone(), &[0]); + + let actual = concat_series_batches(&schema, &batches, &[0]).unwrap(); + let tag = actual + .column(0) + .as_any() + .downcast_ref::>() + .unwrap(); + let values = tag + .values() + .as_any() + .downcast_ref::() + .unwrap(); + assert_eq!(values.len(), 1); + assert_eq!(values.value(0), active); + assert_eq!( + values + .data_buffers() + .iter() + .map(|buffer| buffer.len()) + .sum::(), + active.len() + ); + } + + #[test] + fn test_concat_series_batches_uint64_nullable_tag() { + let schema = Arc::new(Schema::new(vec![ + Field::new("tsid", DataType::UInt64, true), + Field::new("value", DataType::Int64, true), + ])); + for tags in [vec![Some(42), Some(42)], vec![None, None]] { + let batches = [tags.clone(), tags] + .into_iter() + .enumerate() + .map(|(batch, tags)| { + RecordBatch::try_new( + schema.clone(), + vec![ + Arc::new(UInt64Array::from(tags)), + Arc::new(Int64Array::from_iter_values( + (0..2).map(|row| (batch * 10 + row) as i64), + )), + ], + ) + .unwrap() + }) + .collect(); + assert_concat_matches_reference(schema.clone(), batches, &[0]); + } + } + + #[test] + fn test_concat_series_batches_leading_empty_and_sliced_batches() { + let schema = Arc::new(Schema::new(vec![ + Field::new("tag", DataType::Utf8, true), + Field::new("value", DataType::Int64, true), + ])); + let make_sliced_batch = |values| { + RecordBatch::try_new( + schema.clone(), + vec![ + Arc::new(StringArray::from(vec!["discard", "tag", "tag"])), + Arc::new(Int64Array::from(values)), + ], + ) + .unwrap() + .slice(1, 2) + }; + let batches = vec![ + RecordBatch::new_empty(schema.clone()), + make_sliced_batch(vec![0, 1, 2]), + make_sliced_batch(vec![3, 4, 5]), + ]; + assert_concat_matches_reference(schema, batches, &[0]); + } + + #[test] + fn test_concat_series_batches_interleaved_tags_and_schema_metadata() { + let schema = Arc::new(Schema::new_with_metadata( + vec![ + Field::new("value_before", DataType::Int64, true), + Field::new("host", DataType::Utf8, true), + Field::new("value_between", DataType::UInt64, true), + Field::new("path", DataType::Utf8, true), + Field::new("value_after", DataType::Int32, true), + ], + HashMap::from([("source".to_string(), "concat test".to_string())]), + )); + let batches = [ + ( + vec![Some(1), None], + vec![Some(10), Some(11)], + vec![100, 101], + ), + ( + vec![Some(2), Some(3)], + vec![Some(12), Some(13)], + vec![102, 103], + ), + ] + .into_iter() + .map(|(before, between, after)| { + RecordBatch::try_new( + schema.clone(), + vec![ + Arc::new(Int64Array::from(before)), + Arc::new(StringArray::from(vec!["host-a", "host-a"])), + Arc::new(UInt64Array::from(between)), + Arc::new(StringArray::from(vec!["/metrics", "/metrics"])), + Arc::new(Int32Array::from(after)), + ], + ) + .unwrap() + }) + .collect(); + assert_concat_matches_reference(schema, batches, &[1, 3]); + } + + #[test] + fn test_concat_series_batches_fallbacks() { + let schema = Arc::new(Schema::new(vec![Field::new( + "value", + DataType::Int64, + true, + )])); + let batch = RecordBatch::try_new( + schema.clone(), + vec![Arc::new(Int64Array::from(vec![Some(1), None]))], + ) + .unwrap(); + assert_concat_matches_reference(schema.clone(), vec![batch.clone()], &[0]); + assert_concat_matches_reference(schema.clone(), vec![batch.clone(), batch], &[]); + + let empty_schema = Arc::new(Schema::empty()); + let empty_batch = RecordBatch::try_new_with_options( + empty_schema.clone(), + vec![], + &datafusion::arrow::record_batch::RecordBatchOptions::new().with_row_count(Some(2)), + ) + .unwrap(); + assert_concat_matches_reference( + empty_schema.clone(), + vec![empty_batch.clone(), empty_batch], + &[], + ); + + let zero_rows = RecordBatch::new_empty(schema.clone()); + assert_concat_matches_reference(schema, vec![zero_rows.clone(), zero_rows], &[0]); + } + #[test] fn test_dictionary_tag_child_null_comparison() { let dictionary: ArrayRef = Arc::new(DictionaryArray::::new( @@ -1114,6 +1647,99 @@ mod test { assert!(divide_stream.next().await.is_none()); } + #[tokio::test] + async fn test_dictionary_tags_across_batches_and_eof() { + let schema = Arc::new(Schema::new(vec![ + Field::new( + "tag", + DataType::Dictionary(Box::new(DataType::UInt32), Box::new(DataType::Utf8)), + false, + ), + Field::new("value", DataType::Int64, false), + Field::new( + "time_index", + DataType::Timestamp(datafusion::arrow::datatypes::TimeUnit::Millisecond, None), + false, + ), + ])); + let make_batch = |values: Vec<&str>, keys: Vec, payload: Vec| { + RecordBatch::try_new( + schema.clone(), + vec![ + Arc::new(DictionaryArray::::new( + UInt32Array::from(keys), + Arc::new(StringArray::from(values)), + )), + Arc::new(Int64Array::from(payload.clone())), + Arc::new(datafusion::arrow::array::TimestampMillisecondArray::from( + payload, + )), + ], + ) + .unwrap() + }; + let memory_exec: Arc = Arc::new(DataSourceExec::new(Arc::new( + MemorySourceConfig::try_new( + &[vec![ + make_batch(vec!["b", "a", "unused-a"], vec![1, 1, 0], vec![1, 2, 3]), + make_batch(vec!["c", "b", "unused-b"], vec![1, 1, 0], vec![4, 5, 6]), + make_batch(vec!["unused-c", "c"], vec![1, 1], vec![7, 8]), + ]], + schema, + None, + ) + .unwrap(), + ))); + let divide_exec = Arc::new(SeriesDivideExec { + tag_columns: vec!["tag".to_string()], + time_index_column: "time_index".to_string(), + input: memory_exec, + metric: ExecutionPlanMetricsSet::new(), + }); + let mut stream = divide_exec + .execute(0, SessionContext::default().task_ctx()) + .unwrap(); + + for (expected_tag, expected_payload, concatenated) in [ + ("a", vec![1, 2], false), + ("b", vec![3, 4, 5], true), + ("c", vec![6, 7, 8], true), + ] { + let batch = stream.next().await.unwrap().unwrap(); + assert_eq!(batch.num_rows(), expected_payload.len()); + assert_eq!( + (0..batch.num_rows()) + .map(|row| string_array_value_at_index(batch.column(0), row).unwrap()) + .collect::>(), + vec![expected_tag; expected_payload.len()] + ); + assert_eq!( + batch + .column(1) + .as_any() + .downcast_ref::() + .unwrap() + .iter() + .flatten() + .collect::>(), + expected_payload + ); + + if concatenated { + let tag = batch + .column(0) + .as_any() + .downcast_ref::>() + .unwrap(); + let values = tag.values().as_any().downcast_ref::().unwrap(); + assert_eq!(values.len(), 1, "tag {expected_tag}"); + assert_eq!(values.value(0), expected_tag); + assert!(tag.keys().iter().all(|key| key == Some(0))); + } + } + assert!(stream.next().await.is_none()); + } + #[tokio::test] async fn test_string_tag_column_types() { let schema = Arc::new(Schema::new(vec![ diff --git a/tests/perf/query_cases/promql_constant_tag_concat_ms/case.toml b/tests/perf/query_cases/promql_constant_tag_concat_ms/case.toml new file mode 100644 index 0000000000..66e432d90f --- /dev/null +++ b/tests/perf/query_cases/promql_constant_tag_concat_ms/case.toml @@ -0,0 +1,139 @@ +# Millisecond timestamps exercise PerSeries/SeriesScan without the ns-to-ms cast. +# CI performance qualification for the constant-tag concat experiment. +# +# This direct-SST fixture fixes file, row-group, and label layout independently +# of ingestion and SQL flush behavior; it does not claim to reproduce a +# SQL-flush baseline. With timestamp-major generation, every series occurs at +# every scrape timestamp and each series spans all 32 non-overlapping SSTs, +# exercising cross-file reads for the multi-evaluation queries. + +[case] +name = "promql_constant_tag_concat_ms" +description = "PromQL selector and topk/bottomk qualification on a controlled constant-tag direct-SST fixture" + +[scenario] +kind = "direct_readable_sst" +seed = 4096 + +[[scenario.tables]] +database = "public" +name = "promql_constant_tag_concat_ms" +engine = "mito" +append_mode = true +sst_format = "flat" +primary_key = ["host", "instance"] +time_index = "ts" + +[[scenario.tables.columns]] +name = "host" +type = "STRING" +semantic = "tag" +distribution = { kind = "cardinality", values = 16, prefix = "host" } + +[[scenario.tables.columns]] +name = "instance" +type = "STRING" +semantic = "tag" +distribution = { kind = "cardinality", values = 4096, prefix = "instance" } + +[[scenario.tables.columns]] +name = "reading" +type = "DOUBLE" +semantic = "field" +distribution = { kind = "deterministic_wave", min = 0.0, max = 1000.0 } + +[[scenario.tables.columns]] +name = "ts" +type = "TIMESTAMP(3)" +semantic = "timestamp" + +[scenario.layout] +regions = 1 +sst_count = 32 +rows_per_sst = 32768 +row_group_size = 8192 +series_count = 4096 +start_unix_nanos = 1_704_067_200_000_000_000 # 2024-01-01T00:00:00Z +step_nanos = 15_000_000_000 +time_range_layout = "non_overlapping_per_sst" +series_layout = "timestamp_major" + +[[scenario.queries]] +name = "selector_full" +kind = "tql" +query = "TQL EVAL (1704069000, 1704069900, '15s') promql_constant_tag_concat_ms" +warmup = 3 +iterations = 9 + +[scenario.queries.thresholds] +max_candidate_latency_regression_pct = 10 + +[[scenario.queries]] +name = "selector_host1" +kind = "tql" +query = "TQL EVAL (1704069000, 1704069900, '15s') promql_constant_tag_concat_ms{host='host1'}" +warmup = 3 +iterations = 9 + +[scenario.queries.thresholds] +max_candidate_latency_regression_pct = 10 + +[[scenario.queries]] +name = "topk_5_full" +kind = "tql" +query = "TQL EVAL (1704069000, 1704069900, '15s') topk(5, promql_constant_tag_concat_ms)" +warmup = 3 +iterations = 9 + +[scenario.queries.thresholds] +max_candidate_latency_regression_pct = 10 + +[[scenario.queries]] +name = "topk_5_host1" +kind = "tql" +query = "TQL EVAL (1704069000, 1704069900, '15s') topk(5, promql_constant_tag_concat_ms{host='host1'})" +warmup = 3 +iterations = 9 + +[scenario.queries.thresholds] +max_candidate_latency_regression_pct = 10 + +[[scenario.queries]] +name = "topk_100_full" +kind = "tql" +query = "TQL EVAL (1704069000, 1704069900, '15s') topk(100, promql_constant_tag_concat_ms)" +warmup = 3 +iterations = 9 + +[scenario.queries.thresholds] +max_candidate_latency_regression_pct = 10 + +[[scenario.queries]] +name = "bottomk_5_full" +kind = "tql" +query = "TQL EVAL (1704069000, 1704069900, '15s') bottomk(5, promql_constant_tag_concat_ms)" +warmup = 3 +iterations = 9 + +[scenario.queries.thresholds] +max_candidate_latency_regression_pct = 10 + +# One evaluation time is a likely single-batch selector control, not a guarantee: +# PromQL lookback can still require multiple batches. +[[scenario.queries]] +name = "instant_selector_single_evaluation_control" +kind = "tql" +query = "TQL EVAL (1704069000, 1704069000, '15s') promql_constant_tag_concat_ms" +warmup = 3 +iterations = 9 + +[scenario.queries.thresholds] +max_candidate_latency_regression_pct = 10 + +# A single evaluation for output-parity inspection, separate from timed controls. +[[scenario.queries]] +name = "instant_topk_5_output_parity_control" +kind = "tql" +query = "TQL EVAL (1704069000, 1704069000, '15s') topk(5, promql_constant_tag_concat_ms)" +warmup = 0 +iterations = 1 diff --git a/tests/perf/query_cases/promql_constant_tag_concat_ms_100k/case.toml b/tests/perf/query_cases/promql_constant_tag_concat_ms_100k/case.toml new file mode 100644 index 0000000000..59e45f226b --- /dev/null +++ b/tests/perf/query_cases/promql_constant_tag_concat_ms_100k/case.toml @@ -0,0 +1,138 @@ +# Millisecond timestamp high-cardinality qualification for constant-tag concat. +# +# This direct-SST fixture fixes file, row-group, and label layout independently +# of ingestion and SQL flush behavior; it does not claim to reproduce a +# SQL-flush baseline. With timestamp-major generation, every series occurs at +# every scrape timestamp and each series has eight samples in all 32 +# non-overlapping SSTs, exercising cross-file reads for the multi-evaluation queries. + +[case] +name = "promql_constant_tag_concat_ms_100k" +description = "PromQL selector and topk/bottomk qualification on a controlled constant-tag direct-SST fixture" + +[scenario] +kind = "direct_readable_sst" +seed = 4096 + +[[scenario.tables]] +database = "public" +name = "promql_constant_tag_concat_ms_100k" +engine = "mito" +append_mode = true +sst_format = "flat" +primary_key = ["host", "instance"] +time_index = "ts" + +[[scenario.tables.columns]] +name = "host" +type = "STRING" +semantic = "tag" +distribution = { kind = "cardinality", values = 16, prefix = "host" } + +[[scenario.tables.columns]] +name = "instance" +type = "STRING" +semantic = "tag" +distribution = { kind = "cardinality", values = 100_000, prefix = "instance" } + +[[scenario.tables.columns]] +name = "reading" +type = "DOUBLE" +semantic = "field" +distribution = { kind = "deterministic_wave", min = 0.0, max = 1000.0 } + +[[scenario.tables.columns]] +name = "ts" +type = "TIMESTAMP(3)" +semantic = "timestamp" + +[scenario.layout] +regions = 1 +sst_count = 32 +rows_per_sst = 800000 +row_group_size = 8192 +series_count = 100_000 +start_unix_nanos = 1_704_067_200_000_000_000 # 2024-01-01T00:00:00Z +step_nanos = 15_000_000_000 +time_range_layout = "non_overlapping_per_sst" +series_layout = "timestamp_major" + +[[scenario.queries]] +name = "selector_full" +kind = "tql" +query = "TQL EVAL (1704069000, 1704069900, '15s') promql_constant_tag_concat_ms_100k" +warmup = 3 +iterations = 9 + +[scenario.queries.thresholds] +max_candidate_latency_regression_pct = 10 + +[[scenario.queries]] +name = "selector_host1" +kind = "tql" +query = "TQL EVAL (1704069000, 1704069900, '15s') promql_constant_tag_concat_ms_100k{host='host1'}" +warmup = 3 +iterations = 9 + +[scenario.queries.thresholds] +max_candidate_latency_regression_pct = 10 + +[[scenario.queries]] +name = "topk_5_full" +kind = "tql" +query = "TQL EVAL (1704069000, 1704069900, '15s') topk(5, promql_constant_tag_concat_ms_100k)" +warmup = 3 +iterations = 9 + +[scenario.queries.thresholds] +max_candidate_latency_regression_pct = 10 + +[[scenario.queries]] +name = "topk_5_host1" +kind = "tql" +query = "TQL EVAL (1704069000, 1704069900, '15s') topk(5, promql_constant_tag_concat_ms_100k{host='host1'})" +warmup = 3 +iterations = 9 + +[scenario.queries.thresholds] +max_candidate_latency_regression_pct = 10 + +[[scenario.queries]] +name = "topk_100_full" +kind = "tql" +query = "TQL EVAL (1704069000, 1704069900, '15s') topk(100, promql_constant_tag_concat_ms_100k)" +warmup = 3 +iterations = 9 + +[scenario.queries.thresholds] +max_candidate_latency_regression_pct = 10 + +[[scenario.queries]] +name = "bottomk_5_full" +kind = "tql" +query = "TQL EVAL (1704069000, 1704069900, '15s') bottomk(5, promql_constant_tag_concat_ms_100k)" +warmup = 3 +iterations = 9 + +[scenario.queries.thresholds] +max_candidate_latency_regression_pct = 10 + +# One evaluation time is a likely single-batch selector control, not a guarantee: +# PromQL lookback can still require multiple batches. +[[scenario.queries]] +name = "instant_selector_single_evaluation_control" +kind = "tql" +query = "TQL EVAL (1704069000, 1704069000, '15s') promql_constant_tag_concat_ms_100k" +warmup = 3 +iterations = 9 + +[scenario.queries.thresholds] +max_candidate_latency_regression_pct = 10 + +# A single evaluation for output-parity inspection, separate from timed controls. +[[scenario.queries]] +name = "instant_topk_5_output_parity_control" +kind = "tql" +query = "TQL EVAL (1704069000, 1704069000, '15s') topk(5, promql_constant_tag_concat_ms_100k)" +warmup = 0 +iterations = 1 diff --git a/tests/perf/query_cases/promql_constant_tag_concat_ms_10k/case.toml b/tests/perf/query_cases/promql_constant_tag_concat_ms_10k/case.toml new file mode 100644 index 0000000000..0046fa6742 --- /dev/null +++ b/tests/perf/query_cases/promql_constant_tag_concat_ms_10k/case.toml @@ -0,0 +1,138 @@ +# Millisecond timestamp high-cardinality qualification for constant-tag concat. +# +# This direct-SST fixture fixes file, row-group, and label layout independently +# of ingestion and SQL flush behavior; it does not claim to reproduce a +# SQL-flush baseline. With timestamp-major generation, every series occurs at +# every scrape timestamp and each series has eight samples in all 32 +# non-overlapping SSTs, exercising cross-file reads for the multi-evaluation queries. + +[case] +name = "promql_constant_tag_concat_ms_10k" +description = "PromQL selector and topk/bottomk qualification on a controlled constant-tag direct-SST fixture" + +[scenario] +kind = "direct_readable_sst" +seed = 4096 + +[[scenario.tables]] +database = "public" +name = "promql_constant_tag_concat_ms_10k" +engine = "mito" +append_mode = true +sst_format = "flat" +primary_key = ["host", "instance"] +time_index = "ts" + +[[scenario.tables.columns]] +name = "host" +type = "STRING" +semantic = "tag" +distribution = { kind = "cardinality", values = 16, prefix = "host" } + +[[scenario.tables.columns]] +name = "instance" +type = "STRING" +semantic = "tag" +distribution = { kind = "cardinality", values = 10_000, prefix = "instance" } + +[[scenario.tables.columns]] +name = "reading" +type = "DOUBLE" +semantic = "field" +distribution = { kind = "deterministic_wave", min = 0.0, max = 1000.0 } + +[[scenario.tables.columns]] +name = "ts" +type = "TIMESTAMP(3)" +semantic = "timestamp" + +[scenario.layout] +regions = 1 +sst_count = 32 +rows_per_sst = 80000 +row_group_size = 8192 +series_count = 10_000 +start_unix_nanos = 1_704_067_200_000_000_000 # 2024-01-01T00:00:00Z +step_nanos = 15_000_000_000 +time_range_layout = "non_overlapping_per_sst" +series_layout = "timestamp_major" + +[[scenario.queries]] +name = "selector_full" +kind = "tql" +query = "TQL EVAL (1704069000, 1704069900, '15s') promql_constant_tag_concat_ms_10k" +warmup = 3 +iterations = 9 + +[scenario.queries.thresholds] +max_candidate_latency_regression_pct = 10 + +[[scenario.queries]] +name = "selector_host1" +kind = "tql" +query = "TQL EVAL (1704069000, 1704069900, '15s') promql_constant_tag_concat_ms_10k{host='host1'}" +warmup = 3 +iterations = 9 + +[scenario.queries.thresholds] +max_candidate_latency_regression_pct = 10 + +[[scenario.queries]] +name = "topk_5_full" +kind = "tql" +query = "TQL EVAL (1704069000, 1704069900, '15s') topk(5, promql_constant_tag_concat_ms_10k)" +warmup = 3 +iterations = 9 + +[scenario.queries.thresholds] +max_candidate_latency_regression_pct = 10 + +[[scenario.queries]] +name = "topk_5_host1" +kind = "tql" +query = "TQL EVAL (1704069000, 1704069900, '15s') topk(5, promql_constant_tag_concat_ms_10k{host='host1'})" +warmup = 3 +iterations = 9 + +[scenario.queries.thresholds] +max_candidate_latency_regression_pct = 10 + +[[scenario.queries]] +name = "topk_100_full" +kind = "tql" +query = "TQL EVAL (1704069000, 1704069900, '15s') topk(100, promql_constant_tag_concat_ms_10k)" +warmup = 3 +iterations = 9 + +[scenario.queries.thresholds] +max_candidate_latency_regression_pct = 10 + +[[scenario.queries]] +name = "bottomk_5_full" +kind = "tql" +query = "TQL EVAL (1704069000, 1704069900, '15s') bottomk(5, promql_constant_tag_concat_ms_10k)" +warmup = 3 +iterations = 9 + +[scenario.queries.thresholds] +max_candidate_latency_regression_pct = 10 + +# One evaluation time is a likely single-batch selector control, not a guarantee: +# PromQL lookback can still require multiple batches. +[[scenario.queries]] +name = "instant_selector_single_evaluation_control" +kind = "tql" +query = "TQL EVAL (1704069000, 1704069000, '15s') promql_constant_tag_concat_ms_10k" +warmup = 3 +iterations = 9 + +[scenario.queries.thresholds] +max_candidate_latency_regression_pct = 10 + +# A single evaluation for output-parity inspection, separate from timed controls. +[[scenario.queries]] +name = "instant_topk_5_output_parity_control" +kind = "tql" +query = "TQL EVAL (1704069000, 1704069000, '15s') topk(5, promql_constant_tag_concat_ms_10k)" +warmup = 0 +iterations = 1 diff --git a/tests/perf/test_query_regression_case_selection.py b/tests/perf/test_query_regression_case_selection.py index 04a542e23a..16ba726142 100644 --- a/tests/perf/test_query_regression_case_selection.py +++ b/tests/perf/test_query_regression_case_selection.py @@ -51,10 +51,14 @@ class QueryRegressionCaseSelectionTest(unittest.TestCase): self.assertEqual(len(runner.DEFAULT_CASES), 9) self.assertEqual(runner.split_cases([]), runner.DEFAULT_CASES) - def test_heavy_selects_only_remote_write_7913(self) -> None: + def test_heavy_selects_the_heavy_case_set(self) -> None: self.assertEqual( runner.split_cases(["heavy"]), - ["tests/perf/query_cases/prom_remote_write_7913/case.toml"], + [ + "tests/perf/query_cases/prom_remote_write_7913/case.toml", + "tests/perf/query_cases/promql_constant_tag_concat_ms_10k/case.toml", + "tests/perf/query_cases/promql_constant_tag_concat_ms_100k/case.toml", + ], ) def test_explicit_paths_remain_selectable(self) -> None: