From 60899d19ee9d7cc74e28d4d71eb6b81c12de3697 Mon Sep 17 00:00:00 2001 From: dennis zhuang Date: Thu, 10 Sep 2026 04:20:57 +0000 Subject: [PATCH] perf(servers): defer Prometheus sample value formatting to serialization (#9091) Building a matrix response allocated one `String` per sample while the record batches were scanned, then dropped it after the JSON body was written. Keep the `f64` in `PromSampleValue::Number` instead and format it with ryu while serializing, so no per-sample string is allocated. `PromSampleValue::Text` keeps values parsed from a JSON body, so deserializing and re-serializing a response is unchanged. Vector and scalar results still expose `String`, since they hold a single sample. Signed-off-by: Dennis Zhuang --- Cargo.lock | 1 + src/servers/Cargo.toml | 7 + src/servers/benches/prometheus_response.rs | 231 ++++++++++++++++++ src/servers/src/http/prometheus.rs | 4 +- .../src/http/result/prometheus_resp.rs | 163 ++++++++++-- tests-integration/tests/grpc.rs | 16 +- 6 files changed, 401 insertions(+), 21 deletions(-) create mode 100644 src/servers/benches/prometheus_response.rs diff --git a/Cargo.lock b/Cargo.lock index 00e3d85442..ea51491c88 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -13584,6 +13584,7 @@ dependencies = [ "table", "tempfile", "tikv-jemalloc-ctl", + "tikv-jemallocator", "tokio", "tokio-postgres", "tokio-postgres-rustls", diff --git a/src/servers/Cargo.toml b/src/servers/Cargo.toml index 509699d8ba..549a64e253 100644 --- a/src/servers/Cargo.toml +++ b/src/servers/Cargo.toml @@ -174,6 +174,9 @@ tokio-postgres-rustls = "0.14" [target.'cfg(unix)'.dev-dependencies] pprof = { version = "0.14", features = ["criterion", "flamegraph"] } +[target.'cfg(not(windows))'.dev-dependencies] +tikv-jemallocator = "0.6" + [build-dependencies] common-version.workspace = true @@ -190,6 +193,10 @@ required-features = ["testing"] name = "to_http_output" harness = false +[[bench]] +name = "prometheus_response" +harness = false + [[bench]] name = "loki_labels" harness = false diff --git a/src/servers/benches/prometheus_response.rs b/src/servers/benches/prometheus_response.rs new file mode 100644 index 0000000000..5d9325253a --- /dev/null +++ b/src/servers/benches/prometheus_response.rs @@ -0,0 +1,231 @@ +// Copyright 2023 Greptime Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use std::hint::black_box; +use std::sync::Arc; +use std::time::Duration; + +use axum::response::IntoResponse; +use common_query::Output; +use common_recordbatch::{RecordBatch, RecordBatches}; +use criterion::{Criterion, Throughput, criterion_group, criterion_main}; +use datatypes::data_type::ConcreteDataType; +use datatypes::schema::{ColumnSchema, Schema, SchemaRef}; +use datatypes::vectors::{Float64Vector, StringVector, TimestampMillisecondVector, VectorRef}; +use promql_parser::parser::value::ValueType; +use servers::http::prometheus::{ + PromQueryResult, PromSampleValue, PrometheusJsonResponse, PrometheusResponse, +}; + +const SERIES: usize = 64; +const POINTS: usize = 2048; +/// Label count of the one-sample-per-series shape an instant query returns. +const LABELS: usize = 4; + +// Response building is dominated by allocation, so match the allocator the +// server binary uses instead of the platform default. +#[cfg(not(windows))] +#[global_allocator] +static ALLOC: tikv_jemallocator::Jemalloc = tikv_jemallocator::Jemalloc; + +/// Builds a query result of `series_count` series with `points` samples each. +/// +/// `run_length` is the number of consecutive rows that belong to the same +/// series, and `permuted` shuffles timestamps so the response has to sort them. +fn input( + label_count: usize, + run_length: usize, + series_count: usize, + points: usize, + permuted: bool, +) -> (SchemaRef, Vec) { + let mut columns = vec![ + ColumnSchema::new( + "timestamp", + ConcreteDataType::timestamp_millisecond_datatype(), + false, + ), + ColumnSchema::new("value", ConcreteDataType::float64_datatype(), false), + ]; + columns.extend((0..label_count).map(|index| { + ColumnSchema::new( + format!("label_{index}"), + ConcreteDataType::string_datatype(), + false, + ) + })); + let schema = Arc::new(Schema::new(columns)); + let point = |row: usize| { + let point = row / (series_count * run_length) * run_length + row % run_length; + if permuted { + (point * 109 + 17) % points + } else { + point + } + }; + let mut vectors: Vec = vec![ + Arc::new(TimestampMillisecondVector::from_vec( + (0..series_count * points) + .map(|row| point(row) as i64 * 300_000) + .collect(), + )), + Arc::new(Float64Vector::from( + (0..series_count * points) + .map(|row| Some((point(row) % 1000) as f64 * 0.25)) + .collect::>(), + )), + ]; + for label in 0..label_count { + let values: Vec<_> = (0..series_count) + .map(|series| format!("label-{label}-series-{series:04}")) + .collect(); + vectors.push(Arc::new(StringVector::from( + (0..series_count * points) + .map(|row| Some(values[(row / run_length) % series_count].as_str())) + .collect::>(), + ))); + } + let batch = RecordBatch::new(schema.clone(), vectors).unwrap(); + let batches = (0..batch.num_rows()) + .step_by(1024) + .map(|offset| batch.slice(offset, 1024).unwrap()) + .collect(); + (schema, batches) +} + +fn bench_prometheus_response(c: &mut Criterion) { + let runtime = tokio::runtime::Builder::new_current_thread() + .build() + .unwrap(); + let mut group = c.benchmark_group("prometheus_response_complete"); + group.sample_size(20); + group.warm_up_time(Duration::from_millis(500)); + group.measurement_time(Duration::from_secs(2)); + group.throughput(Throughput::Elements((SERIES * POINTS) as u64)); + for labels in [1, 4] { + for run in [1, 32, POINTS] { + for permuted in [false, true] { + let (schema, batches) = input(labels, run, SERIES, POINTS, permuted); + let convert = || { + runtime.block_on(PrometheusJsonResponse::from_query_result( + Ok(Output::new_with_record_batches( + RecordBatches::try_new(schema.clone(), batches.clone()).unwrap(), + )), + Some("metric".to_string()), + ValueType::Matrix, + None, + )) + }; + let response = convert(); + assert_eq!(response.status, "success"); + let PrometheusResponse::PromData(data) = response.data else { + panic!("expected Prometheus data"); + }; + let PromQueryResult::Matrix(series) = data.result else { + panic!("expected matrix"); + }; + assert_eq!(series.len(), SERIES); + for (index, series) in series.iter().enumerate() { + assert_eq!(series.metric.len(), labels + 1); + assert_eq!(series.metric["__name__"], "metric"); + for label in 0..labels { + assert_eq!( + series.metric[&format!("label_{label}")], + format!("label-{label}-series-{index:04}") + ); + } + assert!(series.histograms.is_empty()); + assert_eq!(series.values.len(), POINTS); + for (point, (timestamp, value)) in series.values.iter().enumerate() { + assert_eq!(*timestamp, point as f64 * 300.0); + assert!( + matches!(value, PromSampleValue::Number(value) if *value == (point % 1000) as f64 * 0.25) + ); + } + } + let order = if permuted { "permuted" } else { "ordered" }; + group.bench_function(format!("labels{labels}_run{run}_{order}"), |b| { + b.iter(|| black_box(convert().into_response())); + }); + } + } + } + group.finish(); +} + +fn bench_prometheus_single_point(c: &mut Criterion) { + let runtime = tokio::runtime::Builder::new_current_thread() + .build() + .unwrap(); + let mut group = c.benchmark_group("prometheus_single_point_complete"); + group.sample_size(20); + group.warm_up_time(Duration::from_millis(500)); + group.measurement_time(Duration::from_secs(2)); + for series_count in [4096, 65536] { + let (schema, batches) = input(LABELS, 1, series_count, 1, false); + group.throughput(Throughput::Elements(series_count as u64)); + for result_type in [ValueType::Matrix, ValueType::Vector] { + let convert = || { + runtime.block_on(PrometheusJsonResponse::from_query_result( + Ok(Output::new_with_record_batches( + RecordBatches::try_new(schema.clone(), batches.clone()).unwrap(), + )), + Some("metric".to_string()), + result_type, + None, + )) + }; + let response = convert(); + assert_eq!(response.status, "success"); + let PrometheusResponse::PromData(data) = response.data else { + panic!("expected Prometheus data"); + }; + match data.result { + PromQueryResult::Matrix(series) => { + assert_eq!(series.len(), series_count); + assert!(series.iter().all(|series| { + series.metric.len() == LABELS + 1 + && series.values.len() == 1 + && series.values[0].0 == 0.0 + && matches!(series.values[0].1, PromSampleValue::Number(0.0)) + && series.histograms.is_empty() + })); + } + PromQueryResult::Vector(series) => { + assert_eq!(series.len(), series_count); + assert!(series.iter().all(|series| { + series.metric.len() == LABELS + 1 + && series.value.as_ref().is_some_and(|(timestamp, value)| { + *timestamp == 0.0 && value == "0.0" + }) + && series.histogram.is_none() + })); + } + _ => panic!("expected matrix or vector"), + } + group.bench_function( + format!("{result_type}_labels{LABELS}_series{series_count}"), + |b| b.iter(|| black_box(convert().into_response())), + ); + } + } + group.finish(); +} + +criterion_group!( + benches, + bench_prometheus_response, + bench_prometheus_single_point +); +criterion_main!(benches); diff --git a/src/servers/src/http/prometheus.rs b/src/servers/src/http/prometheus.rs index e1147c92b6..2c9d8267ee 100644 --- a/src/servers/src/http/prometheus.rs +++ b/src/servers/src/http/prometheus.rs @@ -71,7 +71,7 @@ use table::requests::{ SEMANTIC_METRIC_UNIT, SEMANTIC_VALUE_MIXED, }; -pub use super::result::prometheus_resp::PrometheusJsonResponse; +pub use super::result::prometheus_resp::{PromSampleValue, PrometheusJsonResponse}; use crate::error::{ CollectRecordbatchSnafu, ConvertScalarValueSnafu, DataFusionSnafu, Error, InvalidQuerySnafu, NotSupportedSnafu, Result, TableNotFoundSnafu, UnexpectedResultSnafu, @@ -100,7 +100,7 @@ pub struct PromSeriesVector { pub struct PromSeriesMatrix { pub metric: BTreeMap, #[serde(skip_serializing_if = "Vec::is_empty", default)] - pub values: Vec<(f64, String)>, + pub values: Vec<(f64, PromSampleValue)>, #[serde(skip_serializing_if = "Vec::is_empty", default)] pub histograms: Vec<(f64, PromNativeHistogram)>, } diff --git a/src/servers/src/http/result/prometheus_resp.rs b/src/servers/src/http/result/prometheus_resp.rs index bc1d084903..7ec15a8c1e 100644 --- a/src/servers/src/http/result/prometheus_resp.rs +++ b/src/servers/src/http/result/prometheus_resp.rs @@ -39,7 +39,7 @@ use indexmap::IndexMap; use promql_parser::label::METRIC_NAME; use promql_parser::parser::value::ValueType; use ryu::Buffer; -use serde::{Deserialize, Serialize}; +use serde::{Deserialize, Serialize, Serializer}; use serde_json::Value; use snafu::{OptionExt, ResultExt}; @@ -55,10 +55,45 @@ use crate::http::prometheus::{ #[derive(Default)] struct PromSeriesSamples { - values: Vec<(f64, String)>, + values: Vec<(f64, PromSampleValue)>, histograms: Vec<(f64, PromNativeHistogram)>, } +/// A sample value of the Prometheus HTTP API JSON format. +/// +/// Samples read out of a query result are kept as `f64` and formatted while the +/// response is serialized, which avoids one `String` per sample. Samples parsed +/// from a JSON body keep their original spelling, so a response that is +/// deserialized and serialized again is unchanged. +#[derive(Debug, Clone, Deserialize, PartialEq)] +#[serde(untagged)] +pub enum PromSampleValue { + #[serde(skip_deserializing)] + Number(f64), + Text(String), +} + +impl Serialize for PromSampleValue { + fn serialize(&self, serializer: S) -> std::result::Result { + match self { + Self::Number(value) if value.is_finite() => { + serializer.serialize_str(Buffer::new().format_finite(*value)) + } + Self::Number(value) => serializer.collect_str(value), + Self::Text(value) => serializer.serialize_str(value), + } + } +} + +impl PromSampleValue { + fn into_string(self) -> String { + match self { + Self::Number(value) => format_prometheus_sample_value(value), + Self::Text(value) => value, + } + } +} + fn prometheus_native_histogram(histogram: &NativeHistogram) -> Result { Ok(PromNativeHistogram { count: format_prometheus_float(histogram.count), @@ -427,7 +462,7 @@ impl PrometheusJsonResponse { } else if let Some((timestamp_millis, value)) = value { samples.values.push(( timestamp_millis as f64 / 1000.0, - format_prometheus_sample_value(value), + PromSampleValue::Number(value), )); } } @@ -451,7 +486,10 @@ impl PrometheusJsonResponse { PromQueryResult::Vector(ref mut v) => { let histogram = samples.histograms.pop(); let value = if histogram.is_none() { - samples.values.pop() + samples + .values + .pop() + .map(|(timestamp, value)| (timestamp, value.into_string())) } else { None }; @@ -481,7 +519,10 @@ impl PrometheusJsonResponse { }); } PromQueryResult::Scalar(ref mut v) => { - *v = samples.values.pop(); + *v = samples + .values + .pop() + .map(|(timestamp, value)| (timestamp, value.into_string())); } PromQueryResult::String(ref mut _v) => { // TODO(ruihang): Not supported yet @@ -665,6 +706,101 @@ mod tests { ); } + #[test] + fn sample_value_serialization_matches_eager_formatting() { + let mut values = vec![ + 0.0, + -0.0, + f64::MAX, + f64::MIN, + f64::MIN_POSITIVE, + f64::from_bits(1), + 1e-7, + 1e21, + f64::NAN, + f64::INFINITY, + f64::NEG_INFINITY, + ]; + let mut bits = 0x1234_5678_9876_5432_u64; + for _ in 0..1000 { + bits ^= bits << 13; + bits ^= bits >> 7; + bits ^= bits << 17; + values.push(f64::from_bits(bits)); + } + for value in values { + assert_eq!( + serde_json::to_string(&PromSampleValue::Number(value)).unwrap(), + serde_json::to_string(&format_prometheus_sample_value(value)).unwrap() + ); + } + for value in ["1.00", "+Inf", "-0", "NaN", "not-a-number", "", "\"\\\n"] { + let json = serde_json::to_string(value).unwrap(); + let parsed: PromSampleValue = serde_json::from_str(&json).unwrap(); + assert_eq!(serde_json::to_string(&parsed).unwrap(), json); + } + for value in ["1", "null", "true", "[]", "{}"] { + assert!(serde_json::from_str::(value).is_err()); + } + } + + #[tokio::test] + async fn matrix_response_body_matches_eagerly_formatted_json() { + let schema = Arc::new(Schema::new(vec![ + ColumnSchema::new( + "timestamp", + ConcreteDataType::timestamp_millisecond_datatype(), + false, + ), + ColumnSchema::new("host", ConcreteDataType::string_datatype(), false), + ColumnSchema::new("value", ConcreteDataType::float64_datatype(), false), + ])); + let batches = RecordBatches::try_new( + schema.clone(), + vec![ + RecordBatch::new( + schema, + vec![ + Arc::new(TimestampMillisecondVector::from_values([ + 1000, 2000, 3000, 4000, 5000, + ])) as _, + Arc::new(StringVector::from(vec![Some("a"); 5])) as _, + Arc::new(Float64Vector::from_values([ + -0.0, + f64::NAN, + 1e-7, + f64::INFINITY, + f64::NEG_INFINITY, + ])) as _, + ], + ) + .unwrap(), + ], + ) + .unwrap(); + let actual = PrometheusJsonResponse::from_query_result( + Ok(Output::new_with_record_batches(batches)), + None, + ValueType::Matrix, + None, + ) + .await; + // Deserializing the expectation yields `PromSampleValue::Text`, so this + // compares the deferred numeric encoding against eagerly built strings. + let expected: PrometheusJsonResponse = serde_json::from_value(serde_json::json!({ + "status": "success", + "data": {"resultType": "matrix", "result": [{ + "metric": {"host": "a"}, + "values": [[1.0, "-0.0"], [2.0, "NaN"], [3.0, "1e-7"], [4.0, "inf"], [5.0, "-inf"]] + }]} + })) + .unwrap(); + assert_eq!( + serde_json::to_string(&actual).unwrap(), + serde_json::to_string(&expected).unwrap() + ); + } + #[test] fn matrix_response_preserves_ordinary_nan_and_filters_stale_markers() { let schema = Arc::new(Schema::new(vec![ @@ -704,8 +840,8 @@ mod tests { assert_eq!(series.len(), 1); assert_eq!( - series[0].values, - vec![(1.0, "1.0".to_string()), (2.0, "NaN".to_string())] + serde_json::to_value(&series[0].values).unwrap(), + serde_json::json!([[1.0, "1.0"], [2.0, "NaN"]]) ); } @@ -766,7 +902,10 @@ mod tests { ((index + 1) as f64, expected_value) }) .collect::>(); - assert_eq!(series[0].values, expected); + assert_eq!( + serde_json::to_value(&series[0].values).unwrap(), + serde_json::to_value(expected).unwrap() + ); } #[test] @@ -809,12 +948,8 @@ mod tests { assert_eq!(series.len(), 1); assert_eq!( - series[0].values, - vec![ - (1.0, "inf".to_string()), - (2.0, "-inf".to_string()), - (3.0, "NaN".to_string()), - ] + serde_json::to_value(&series[0].values).unwrap(), + serde_json::json!([[1.0, "inf"], [2.0, "-inf"], [3.0, "NaN"]]) ); } diff --git a/tests-integration/tests/grpc.rs b/tests-integration/tests/grpc.rs index 242161654b..1d9b123d1b 100644 --- a/tests-integration/tests/grpc.rs +++ b/tests-integration/tests/grpc.rs @@ -55,8 +55,8 @@ use otel_arrow_rust::schema::consts as arrow_consts; use servers::grpc::GrpcServerConfig; use servers::grpc::builder::GrpcServerBuilder; use servers::http::prometheus::{ - PromData, PromQueryResult, PromSeriesMatrix, PromSeriesVector, PrometheusJsonResponse, - PrometheusResponse, + PromData, PromQueryResult, PromSampleValue, PromSeriesMatrix, PromSeriesVector, + PrometheusJsonResponse, PrometheusResponse, }; use servers::request_memory_limiter::ServerMemoryLimiter; use servers::server::Server; @@ -1538,7 +1538,7 @@ pub async fn test_prom_gateway_query(store_type: StorageType) { panic!("unexpected result type") }; - mat.sort_unstable_by_key(|v| v.values[0].1.clone()); + mat.sort_unstable_by_key(|v| serde_json::to_string(&v.values[0].1).unwrap()); assert_eq!( mat, @@ -1550,7 +1550,10 @@ pub async fn test_prom_gateway_query(store_type: StorageType) { ] .into_iter() .collect(), - values: vec![(5.0, "1.0".to_string()), (10.0, "1.0".to_string())], + values: vec![ + (5.0, PromSampleValue::Text("1.0".to_string())), + (10.0, PromSampleValue::Text("1.0".to_string())), + ], ..Default::default() }, PromSeriesMatrix { @@ -1560,7 +1563,10 @@ pub async fn test_prom_gateway_query(store_type: StorageType) { ] .into_iter() .collect(), - values: vec![(5.0, "2.0".to_string()), (10.0, "2.0".to_string())], + values: vec![ + (5.0, PromSampleValue::Text("2.0".to_string())), + (10.0, PromSampleValue::Text("2.0".to_string())), + ], ..Default::default() }, ]