From a164c4a66d2e9efec6eb7caca9ff2e1d96a5d945 Mon Sep 17 00:00:00 2001 From: Dennis Zhuang Date: Fri, 11 Sep 2026 14:59:33 +0800 Subject: [PATCH] fix(promql): clamp extrapolation before snapping a counter to zero Prometheus clamps `durationToStart` to half an average interval once the first sample is past the extrapolation threshold, and only then lets the counter zero-snap shorten it further, so the snap can never lengthen the leading extrapolation. Running the snap first let it rescue a duration the clamp should have cut, and `rate` and `increase` over-extrapolated to the left. For samples 1@0s and 2@1s in a 4s window at 1s, upstream gives 0.375 and this returned 0.5. `extrapolation_matches_prometheus_on_seeded_windows` carries a line-by-line port of upstream `extrapolatedRate` as an oracle and diffs it against the UDF over seeded windows, so the order stays pinned. `factor` also picked up upstream's guard against a zero sampled interval, which previously divided by zero. Two sqlness results move, both verified against the upstream algorithm. Signed-off-by: Dennis Zhuang --- src/promql/src/functions/extrapolate_rate.rs | 210 +++++++++++++++--- .../common/promql/null_samples.result | 6 +- .../common/promql/set_operation.result | 8 +- 3 files changed, 191 insertions(+), 33 deletions(-) diff --git a/src/promql/src/functions/extrapolate_rate.rs b/src/promql/src/functions/extrapolate_rate.rs index ef3106e108..263ae6337b 100644 --- a/src/promql/src/functions/extrapolate_rate.rs +++ b/src/promql/src/functions/extrapolate_rate.rs @@ -266,34 +266,35 @@ impl ExtrapolatedRate= extrapolation_threshold { + duration_to_start_ms = average_interval_ms / 2.0; + } + // Counters cannot be negative, so the extrapolation can snap back to the inferred + // zero point instead of extending into negative values. Prometheus applies this + // after the threshold clamp, so it can only shorten the leading extrapolation. if IS_COUNTER && result_value > 0.0 && first_value >= 0.0 { let duration_to_zero = sampled_interval_ms * (first_value / result_value); if duration_to_zero < duration_to_start_ms { duration_to_start_ms = duration_to_zero; } } - - let extrapolation_threshold = average_interval_ms * 1.1; - let mut extrapolated_interval_ms = sampled_interval_ms; - - // Mirror Prometheus extrapolation: extend to the real range boundary when a sample is - // close enough, otherwise add half an average sampling interval on that side. - if duration_to_start_ms < extrapolation_threshold { - extrapolated_interval_ms += duration_to_start_ms; - } else { - extrapolated_interval_ms += average_interval_ms / 2.0; - } - if duration_to_end_ms < extrapolation_threshold { - extrapolated_interval_ms += duration_to_end_ms; - } else { - extrapolated_interval_ms += average_interval_ms / 2.0; + if duration_to_end_ms >= extrapolation_threshold { + duration_to_end_ms = average_interval_ms / 2.0; } - let mut factor = extrapolated_interval_ms / sampled_interval_ms; + // Samples sharing one timestamp leave nothing to extrapolate over. + let mut factor = if sampled_interval_ms == 0.0 { + 1.0 + } else { + (sampled_interval_ms + duration_to_start_ms + duration_to_end_ms) + / sampled_interval_ms + }; if IS_RATE { factor /= range_length_secs; @@ -528,6 +529,7 @@ mod test { use datafusion_common::ScalarValue; use super::*; + use crate::functions::test_util::TinyPrng; /// Range length is fixed to 5 fn extrapolated_rate_runner( @@ -779,6 +781,163 @@ mod test { assert_eq!(output, vec![Some(7.5)]); } + /// Line-by-line port of Prometheus `extrapolatedRate` (promql/functions.go), float path + /// without start timestamps. Kept as a second implementation so that the order of the + /// threshold clamp and the counter zero-snap stays pinned to the upstream one. + fn prometheus_extrapolated_rate( + timestamps: &[i64], + values: &[f64], + eval_ts: i64, + range_ms: i64, + is_counter: bool, + is_rate: bool, + ) -> Option { + if values.len() < 2 { + return None; + } + let num_samples_minus_one = values.len() - 1; + let first_t = timestamps[0]; + let last_t = timestamps[num_samples_minus_one]; + let mut result = values[num_samples_minus_one] - values[0]; + if is_counter { + for index in 1..values.len() { + if values[index] < values[index - 1] { + result += values[index - 1]; + } + } + } + + let range_start = eval_ts - range_ms; + let mut duration_to_start = (first_t - range_start) as f64 / 1000.0; + let mut duration_to_end = (eval_ts - last_t) as f64 / 1000.0; + let sampled_interval = (last_t - first_t) as f64 / 1000.0; + let average_duration_between_samples = sampled_interval / num_samples_minus_one as f64; + let extrapolation_threshold = average_duration_between_samples * 1.1; + + if duration_to_start >= extrapolation_threshold { + duration_to_start = average_duration_between_samples / 2.0; + } + if is_counter { + let mut duration_to_zero = duration_to_start; + if result > 0.0 && values[0] >= 0.0 { + duration_to_zero = sampled_interval * (values[0] / result); + } + if duration_to_zero < duration_to_start { + duration_to_start = duration_to_zero; + } + } + if duration_to_end >= extrapolation_threshold { + duration_to_end = average_duration_between_samples / 2.0; + } + + let mut factor = 1.0; + if sampled_interval != 0.0 { + factor = (sampled_interval + duration_to_start + duration_to_end) / sampled_interval; + } + if is_rate { + factor /= range_ms as f64 / 1000.0; + } + Some(result * factor) + } + + fn assert_matches_prometheus( + timestamps: &[i64], + values: &[f64], + ranges: &[(u32, u32)], + eval_timestamps: &[i64], + range_ms: i64, + ) { + let actual = nullable_rate_runner::( + timestamps.to_vec(), + Arc::new(Float64Array::from(values.to_vec())), + ranges.to_vec(), + eval_timestamps.to_vec(), + range_ms, + ); + + for (index, ((offset, length), eval_ts)) in ranges.iter().zip(eval_timestamps).enumerate() { + let window = *offset as usize..(*offset + *length) as usize; + let expected = prometheus_extrapolated_rate( + ×tamps[window.clone()], + &values[window], + *eval_ts, + range_ms, + IS_COUNTER, + IS_RATE, + ); + match (actual[index], expected) { + (None, None) => {} + (Some(actual), Some(expected)) => assert!( + (actual - expected).abs() <= expected.abs() * 1e-9, + "window {index} {:?}: got {actual}, Prometheus gives {expected}", + ranges[index] + ), + (actual, expected) => { + panic!("window {index}: got {actual:?}, Prometheus gives {expected:?}") + } + } + } + } + + #[test] + fn extrapolation_matches_prometheus_on_seeded_windows() { + let mut prng = TinyPrng(0x51ed_270b_8f26_1a37); + // Uneven spacing so the average interval, and with it the extrapolation threshold, + // differs from window to window. + let timestamps = (0..48) + .scan(0i64, |clock, _| { + *clock += 1_000 + prng.next_index(4) as i64 * 500; + Some(*clock) + }) + .collect::>(); + // A counter that resets a few times, so the zero-snap branch is reached with both + // small and large leading values. + let values = (0..48) + .scan(0.0f64, |counter, _| { + *counter = match prng.next_index(8) { + 0 => 0.0, + 1 => *counter / 2.0, + _ => *counter + prng.next_index(50) as f64, + }; + Some(*counter) + }) + .collect::>(); + + let mut ranges = Vec::new(); + let mut eval_timestamps = Vec::new(); + for _ in 0..32 { + let length = 2 + prng.next_index(10) as u32; + let offset = prng.next_index(48 - length as usize) as u32; + ranges.push((offset, length)); + // Land the range boundary at varying distances from the samples, so the clamp + // fires on neither, one, or both sides. + let last = timestamps[(offset + length - 1) as usize]; + eval_timestamps.push(last + prng.next_index(5) as i64 * 500); + } + + assert_matches_prometheus::( + ×tamps, + &values, + &ranges, + &eval_timestamps, + 20_000, + ); + assert_matches_prometheus::( + ×tamps, + &values, + &ranges, + &eval_timestamps, + 20_000, + ); + assert_matches_prometheus::( + ×tamps, + &values, + &ranges, + &eval_timestamps, + 20_000, + ); + } + #[test] fn rate_rejects_wrong_input_arity() { let err = ExtrapolatedRate::::new(5) @@ -856,7 +1015,7 @@ mod test { ts_range, value_range, timestamps, - vec![2.0, 5.0, 0.0, 2.5, 0.0, 0.0], + vec![1.5, 5.0, 0.0, 2.5, 0.0, 0.0], ); } @@ -887,8 +1046,7 @@ mod test { ts_range, value_range, timestamps, - // `2.0` is because that `duration_to_zero` less than `extrapolation_threshold` - vec![2.0, 1.5, 1.5, 1.5, 1.5, 1.5, 1.5, 1.5], + vec![1.5, 1.5, 1.5, 1.5, 1.5, 1.5, 1.5, 1.5], ); } @@ -961,7 +1119,7 @@ mod test { // that two `2.0` is because `duration_to_start` are shrunk to // `duration_to_zero`, and causes `duration_to_zero` less than // `extrapolation_threshold`. - vec![2.0, 1.5, 1.5, 1.5, 2.0, 1.5, 1.5, 1.5], + vec![1.5, 1.5, 1.5, 1.5, 1.5, 1.5, 1.5, 1.5], ); } @@ -981,7 +1139,7 @@ mod test { ts_range, value_range, timestamps, - vec![4.0, 3.5, 3.5, 4.0], + vec![3.5, 3.5, 3.5, 3.5], ); } @@ -1156,7 +1314,7 @@ mod test { ts_range, value_range, timestamps, - vec![400.0, 300.0, 300.0, 300.0, 400.0, 300.0, 300.0, 300.0], + vec![300.0, 300.0, 300.0, 300.0, 300.0, 300.0, 300.0, 300.0], ); } @@ -1187,7 +1345,7 @@ mod test { ts_range, value_range, timestamps, - vec![400.0, 300.0, 300.0, 300.0, 300.0, 300.0, 300.0, 300.0], + vec![300.0, 300.0, 300.0, 300.0, 300.0, 300.0, 300.0, 300.0], ); } diff --git a/tests/cases/standalone/common/promql/null_samples.result b/tests/cases/standalone/common/promql/null_samples.result index 190815103f..e55fe887b9 100644 --- a/tests/cases/standalone/common/promql/null_samples.result +++ b/tests/cases/standalone/common/promql/null_samples.result @@ -24,8 +24,8 @@ TQL EVAL (0, 15, '1s') avg by (region) (rate(null_samples[4s])); +--------+---------------------+---------------------------------------------+ | region | ts | avg(prom_rate(ts_range,val,ts,Int64(4000))) | +--------+---------------------+---------------------------------------------+ -| ok | 1970-01-01T00:00:01 | 0.5 | -| ok | 1970-01-01T00:00:02 | 0.75 | +| ok | 1970-01-01T00:00:01 | 0.375 | +| ok | 1970-01-01T00:00:02 | 0.625 | | ok | 1970-01-01T00:00:03 | 1.0 | | ok | 1970-01-01T00:00:04 | 0.625 | +--------+---------------------+---------------------------------------------+ @@ -195,7 +195,7 @@ TQL EVAL (1, 1, '1s') rate(multi_field[4s]); +---------------------+---------------------------------------+---------------------------------------+------+ | ts | prom_rate(ts_range,f1,ts,Int64(4000)) | prom_rate(ts_range,f2,ts,Int64(4000)) | host | +---------------------+---------------------------------------+---------------------------------------+------+ -| 1970-01-01T00:00:01 | 0.5 | | a | +| 1970-01-01T00:00:01 | 0.375 | | a | +---------------------+---------------------------------------+---------------------------------------+------+ -- A selector emits the same shape: the row stays, the field without a sample is NULL. diff --git a/tests/cases/standalone/common/promql/set_operation.result b/tests/cases/standalone/common/promql/set_operation.result index 50a0b6cf3b..2c56567f08 100644 --- a/tests/cases/standalone/common/promql/set_operation.result +++ b/tests/cases/standalone/common/promql/set_operation.result @@ -816,12 +816,12 @@ tql eval(1000, 2000, '300s') sum by (src, src_pod, src_namespace, src_node, dest +---------------------+-----------+---------+----------+---------+----------+---------------+----------+---------+----------------------------------------------------------------------------------------------+ | greptime_timestamp | az | cloud | dest | region | src | src_namespace | src_node | src_pod | sum(prom_increase(greptime_timestamp_range,greptime_value,greptime_timestamp,Int64(900000))) | +---------------------+-----------+---------+----------+---------+----------+---------------+----------+---------+----------------------------------------------------------------------------------------------+ -| 1970-01-01T00:21:40 | us-west-6 | cloud-1 | 10.0.0.2 | us-west | 10.0.0.1 | namespace-1 | node-1 | pod-1 | 2500.0 | +| 1970-01-01T00:21:40 | us-west-6 | cloud-1 | 10.0.0.2 | us-west | 10.0.0.1 | namespace-1 | node-1 | pod-1 | 2000.0 | | 1970-01-01T00:21:40 | us-west-6 | cloud-1 | 10.0.0.3 | us-west | 10.0.0.1 | namespace-1 | node-2 | pod-2 | 2000.0 | -| 1970-01-01T00:21:40 | us-west-6 | cloud-2 | 10.0.0.5 | us-west | 10.0.0.4 | namespace-2 | node-3 | pod-3 | 2300.0 | -| 1970-01-01T00:26:40 | us-west-6 | cloud-1 | 10.0.0.2 | us-west | 10.0.0.1 | namespace-1 | node-1 | pod-1 | 2500.0 | +| 1970-01-01T00:21:40 | us-west-6 | cloud-2 | 10.0.0.5 | us-west | 10.0.0.4 | namespace-2 | node-3 | pod-3 | 2000.0 | +| 1970-01-01T00:26:40 | us-west-6 | cloud-1 | 10.0.0.2 | us-west | 10.0.0.1 | namespace-1 | node-1 | pod-1 | 2000.0 | | 1970-01-01T00:26:40 | us-west-6 | cloud-1 | 10.0.0.3 | us-west | 10.0.0.1 | namespace-1 | node-2 | pod-2 | 2000.0 | -| 1970-01-01T00:26:40 | us-west-6 | cloud-2 | 10.0.0.5 | us-west | 10.0.0.4 | namespace-2 | node-3 | pod-3 | 2300.0 | +| 1970-01-01T00:26:40 | us-west-6 | cloud-2 | 10.0.0.5 | us-west | 10.0.0.4 | namespace-2 | node-3 | pod-3 | 2000.0 | +---------------------+-----------+---------+----------+---------+----------+---------------+----------+---------+----------------------------------------------------------------------------------------------+ DROP TABLE node_network_transmit_bytes_total;