diff --git a/src/promql/src/functions/aggr_over_time.rs b/src/promql/src/functions/aggr_over_time.rs index 63f045f134..3eeae47e61 100644 --- a/src/promql/src/functions/aggr_over_time.rs +++ b/src/promql/src/functions/aggr_over_time.rs @@ -156,7 +156,10 @@ fn evaluate_presence( display_name = prom_avg_over_time )] pub fn avg_over_time(_: &TimestampMillisecondArray, values: &Float64Array) -> Option { - compute::sum(values).map(|result| result / values.len() as f64) + // `sum` already skips null slots and yields `None` for an all-null window, so only the + // divisor needs to count samples instead of slots. + let sample_count = values.len() - values.null_count(); + compute::sum(values).map(|result| result / sample_count as f64) } /// The minimum value of all points in the specified interval. @@ -259,26 +262,21 @@ pub fn present_over_time(_: &TimestampMillisecondArray, values: &Float64Array) - display_name = prom_stdvar_over_time )] pub fn stdvar_over_time(_: &TimestampMillisecondArray, values: &Float64Array) -> Option { - if values.is_empty() { - None - } else { - let mut count = 0; - let mut mean: f64 = 0.0; - let mut result: f64 = 0.0; - for value in values { - let value = value.unwrap(); - let new_count = count + 1; - let delta1 = value - mean; - let new_mean = delta1 / new_count as f64 + mean; - let delta2 = value - new_mean; - let new_result = result + delta1 * delta2; + let mut count = 0; + let mut mean: f64 = 0.0; + let mut result: f64 = 0.0; + for value in values.iter().flatten() { + let new_count = count + 1; + let delta1 = value - mean; + let new_mean = delta1 / new_count as f64 + mean; + let delta2 = value - new_mean; + let new_result = result + delta1 * delta2; - count += 1; - mean = new_mean; - result = new_result; - } - Some(result / count as f64) + count = new_count; + mean = new_mean; + result = new_result; } + (count > 0).then(|| result / count as f64) } /// the population standard deviation of the values in the specified interval. @@ -289,35 +287,32 @@ pub fn stdvar_over_time(_: &TimestampMillisecondArray, values: &Float64Array) -> display_name = prom_stddev_over_time )] pub fn stddev_over_time(_: &TimestampMillisecondArray, values: &Float64Array) -> Option { - if values.is_empty() { - None - } else { - let mut count = 0.0; - let mut mean = 0.0; - let mut comp_mean = 0.0; - let mut deviations_sum_sq = 0.0; - let mut comp_deviations_sum_sq = 0.0; - for v in values { - count += 1.0; - let current_value = v.unwrap(); - let delta = current_value - (mean + comp_mean); - let (new_mean, new_comp_mean) = compensated_sum_inc(delta / count, mean, comp_mean); - mean = new_mean; - comp_mean = new_comp_mean; - let (new_deviations_sum_sq, new_comp_deviations_sum_sq) = compensated_sum_inc( - delta * (current_value - (mean + comp_mean)), - deviations_sum_sq, - comp_deviations_sum_sq, - ); - deviations_sum_sq = new_deviations_sum_sq; - comp_deviations_sum_sq = new_comp_deviations_sum_sq; - } - Some(((deviations_sum_sq + comp_deviations_sum_sq) / count).sqrt()) + let mut count = 0.0; + let mut mean = 0.0; + let mut comp_mean = 0.0; + let mut deviations_sum_sq = 0.0; + let mut comp_deviations_sum_sq = 0.0; + for current_value in values.iter().flatten() { + count += 1.0; + let delta = current_value - (mean + comp_mean); + let (new_mean, new_comp_mean) = compensated_sum_inc(delta / count, mean, comp_mean); + mean = new_mean; + comp_mean = new_comp_mean; + let (new_deviations_sum_sq, new_comp_deviations_sum_sq) = compensated_sum_inc( + delta * (current_value - (mean + comp_mean)), + deviations_sum_sq, + comp_deviations_sum_sq, + ); + deviations_sum_sq = new_deviations_sum_sq; + comp_deviations_sum_sq = new_comp_deviations_sum_sq; } + (count > 0.0).then(|| ((deviations_sum_sq + comp_deviations_sum_sq) / count).sqrt()) } #[cfg(test)] mod test { + use datafusion::arrow::buffer::NullBuffer; + use super::*; use crate::functions::test_util::simple_range_udf_runner; @@ -996,4 +991,57 @@ mod test { vec![Some(0.0), Some(3.249615361854384)], ); } + + /// Timestamps and value ranges shared by the null-sample assertions below. + fn null_sample_range_arrays() -> (RangeArray, RangeArray) { + let ts_array = Arc::new(TimestampMillisecondArray::from_iter_values([ + 0i64, 1000, 2000, 3000, + ])); + // Samples are 2.0@0 and 8.0@3000; the null slots keep a payload that would skew every + // aggregate if it were read. + let values_array = Arc::new(Float64Array::new( + vec![2.0, 1000.0, -1000.0, 8.0].into(), + Some(NullBuffer::from_iter([true, false, false, true])), + )); + // The second window holds no sample at all. + let ranges = [(0, 4), (1, 2)]; + + ( + RangeArray::from_ranges(ts_array, ranges).unwrap(), + RangeArray::from_ranges(values_array, ranges).unwrap(), + ) + } + + #[test] + fn avg_over_time_divides_by_sample_count() { + let (ts_array, value_array) = null_sample_range_arrays(); + simple_range_udf_runner( + AvgOverTime::scalar_udf(), + ts_array, + value_array, + vec![], + vec![Some(5.0), None], + ); + } + + #[test] + fn stdvar_and_stddev_over_time_skip_null_samples() { + let (ts_array, value_array) = null_sample_range_arrays(); + simple_range_udf_runner( + StdvarOverTime::scalar_udf(), + ts_array, + value_array, + vec![], + vec![Some(9.0), None], + ); + + let (ts_array, value_array) = null_sample_range_arrays(); + simple_range_udf_runner( + StddevOverTime::scalar_udf(), + ts_array, + value_array, + vec![], + vec![Some(3.0), None], + ); + } } diff --git a/src/promql/src/functions/changes.rs b/src/promql/src/functions/changes.rs index c143d29bef..8630c914a8 100644 --- a/src/promql/src/functions/changes.rs +++ b/src/promql/src/functions/changes.rs @@ -43,8 +43,8 @@ mod test { use super::*; use crate::functions::test_util::{ - self, STALE_NAN, TinyPrng, assert_execution_error, build_test_range_arrays, - invoke_range_udf, simple_range_udf_runner, + self, STALE_NAN, assert_execution_error, build_test_range_arrays, invoke_range_udf, + simple_range_udf_runner, }; use crate::range_array::RangeArray; @@ -151,43 +151,17 @@ mod test { changes_oracle, Changes::scalar_udf(), ); - assert_eq!(expected, vec![Some(2.0), Some(0.0)]); + assert_eq!(expected, vec![Some(0.0), None]); } #[test] fn changes_range_array_seeded_differential() { - let mut prng = TinyPrng(0x2f6e_2b1d_834a_90c5); - let raw_values = (0..48) - .map(|_| match prng.next_index(12) { - 0 => -0.0, - 1 => 0.0, - 2 => -2.0, - 3 => -1.0, - 4 => 1.0, - 5 => 2.0, - 6 => f64::INFINITY, - 7 => f64::NEG_INFINITY, - 8 | 9 => f64::NAN, - _ => STALE_NAN, - }) - .collect::>(); - let values = raw_values.iter().copied().map(Some).collect(); - let mut timestamp_ranges = Vec::new(); - let mut value_ranges = Vec::new(); - for _ in 0..32 { - let length = prng.next_index(13) as u32; - timestamp_ranges.push((prng.next_index(65 - length as usize) as u32, length)); - value_ranges.push((prng.next_index(49 - length as usize) as u32, length)); - } + test_util::run_seeded_differential(changes_oracle, Changes::scalar_udf(), false); + } - test_util::run_oracle_ranges( - values, - raw_values, - timestamp_ranges, - value_ranges, - changes_oracle, - Changes::scalar_udf(), - ); + #[test] + fn changes_range_array_seeded_differential_with_nulls() { + test_util::run_seeded_differential(changes_oracle, Changes::scalar_udf(), true); } #[test] @@ -263,6 +237,8 @@ mod test { vec![Some(0.0), Some(0.0), Some(0.0)], ); + // The raw payload under the null slot would look like two changes; the sample sequence + // is 10 -> 10, which is none. The last window holds no sample at all. let values = Arc::new(Float64Array::new( vec![10.0, 7.0, 10.0].into(), Some(NullBuffer::from(vec![true, false, true])), @@ -276,10 +252,10 @@ mod test { ])); simple_range_udf_runner( Changes::scalar_udf(), - RangeArray::from_ranges(timestamps, [(0, 3)]).unwrap(), - RangeArray::from_ranges(values, [(0, 3)]).unwrap(), + RangeArray::from_ranges(timestamps, [(0, 3), (1, 1)]).unwrap(), + RangeArray::from_ranges(values, [(0, 3), (1, 1)]).unwrap(), vec![], - vec![Some(2.0)], + vec![Some(0.0), None], ); } } diff --git a/src/promql/src/functions/edge_count.rs b/src/promql/src/functions/edge_count.rs index 9109c74207..7cd5907a5e 100644 --- a/src/promql/src/functions/edge_count.rs +++ b/src/promql/src/functions/edge_count.rs @@ -81,7 +81,11 @@ fn calc( .unwrap(); let requested_edges = validate_windows(×tamp_ranges, &value_ranges, name)?; let raw_values = values.values(); - let direct = should_scan_direct(requested_edges, raw_values.len()); + // A NULL field value means the series has no sample at that timestamp. The prefix sums + // encode edges between physically adjacent slots, which no longer holds once nulls are + // skipped, so a null-bearing input falls back to scanning each window. + let has_nulls = values.null_count() > 0; + let direct = has_nulls || should_scan_direct(requested_edges, raw_values.len()); let prefix = (!direct).then(|| build_prefix(raw_values.as_ref(), kind)); let mut result = Vec::with_capacity(value_ranges.len()); @@ -90,6 +94,7 @@ fn calc( let end = checked_end(offset, len, index, name)?; let count = match len { 0 => None, + _ if has_nulls => count_edges_skipping_nulls(values, offset, end, kind), 1 => Some(0), _ if direct => Some(count_edges(raw_values.as_ref(), offset, end, kind)), _ => { @@ -172,6 +177,29 @@ fn count_edges(values: &[f64], offset: usize, end: usize, kind: EdgeKind) -> u64 count } +/// Counts edges between consecutive samples in `[offset, end)`, treating null slots as +/// absent. Returns `None` when the window holds no sample. +fn count_edges_skipping_nulls( + values: &Float64Array, + offset: usize, + end: usize, + kind: EdgeKind, +) -> Option { + let raw_values = values.values(); + let mut count = 0; + let mut previous = None; + for index in offset..end { + if values.is_null(index) { + continue; + } + let current = raw_values[index]; + if let Some(previous) = previous.replace(current) { + count += u64::from(is_edge(previous, current, kind)); + } + } + previous.is_some().then_some(count) +} + fn is_edge(previous: f64, current: f64, kind: EdgeKind) -> bool { match kind { EdgeKind::Changes => previous != current && !(previous.is_nan() && current.is_nan()), diff --git a/src/promql/src/functions/extrapolate_rate.rs b/src/promql/src/functions/extrapolate_rate.rs index f072400d3e..263ae6337b 100644 --- a/src/promql/src/functions/extrapolate_rate.rs +++ b/src/promql/src/functions/extrapolate_rate.rs @@ -181,12 +181,16 @@ impl ExtrapolatedRate() .expect("validated by extract_range_dict") .values(); - let all_values = value_dict + let value_array = value_dict .values() .as_any() .downcast_ref::() - .expect("validated by extract_range_dict") - .values(); + .expect("validated by extract_range_dict"); + // A NULL field value means the series has no sample at that timestamp, so the padding + // under a null slot must not be read. Skip the per-window null scan when the whole + // backing array is null-free, which is the common case. + let has_nulls = value_array.null_count() > 0; + let all_values = value_array.values(); let eval_ts = eval_ts_array.values(); let mut result_builder = Float64Builder::with_capacity(num_windows); @@ -198,7 +202,7 @@ impl ExtrapolatedRate ExtrapolatedRate bounds, + None => { + result_builder.append_null(); + continue; + } + } + } else { + (offset, end.saturating_sub(1), length) + }; + + if sample_count < 2 { result_builder.append_null(); continue; } - let end = offset + length; - let first_value = all_values[offset]; - let last_value = all_values[end - 1]; + let first_value = all_values[first_index]; + let last_value = all_values[last_index]; let mut result_value = last_value - first_value; if IS_COUNTER { - result_value = match &mut reset_index { - Some(reset_index) => reset_index.add_resets(result_value, offset, end), - None => add_counter_resets(result_value, &all_values[offset..end]), + result_value = if has_nulls { + add_counter_resets_between_samples( + result_value, + value_array, + first_index, + last_index, + ) + } else { + match &mut reset_index { + Some(reset_index) => reset_index.add_resets(result_value, offset, end), + None => add_counter_resets(result_value, &all_values[offset..end]), + } }; } - let first_ts = all_timestamps[offset]; - let last_ts = all_timestamps[end - 1]; + let first_ts = all_timestamps[first_index]; + let last_ts = all_timestamps[last_index]; let range_end = eval_ts[index]; let range_start = range_end - range_length; let sampled_interval_ms = (last_ts - first_ts) as f64; - let average_interval_ms = sampled_interval_ms / (length - 1) as f64; + let average_interval_ms = sampled_interval_ms / (sample_count - 1) as f64; let mut duration_to_start_ms = (first_ts - range_start) as f64; - let duration_to_end_ms = (range_end - last_ts) as f64; + let mut duration_to_end_ms = (range_end - last_ts) as f64; + let extrapolation_threshold = average_interval_ms * 1.1; - // Counters cannot be negative, so Prometheus allows the extrapolation window to snap - // back to the inferred zero point instead of extending into negative values. + // Mirror Prometheus extrapolation: extend to the real range boundary when a sample is + // close enough, otherwise only half an average sampling interval, which is the guess + // for where the series actually starts or ends. + if duration_to_start_ms >= 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; @@ -381,6 +407,52 @@ impl<'a> CounterResetIndex<'a> { } } +/// Same additions [`add_counter_resets`] performs, over the samples in `[first, last]` instead +/// of over every slot. +fn add_counter_resets_between_samples( + result: f64, + values: &Float64Array, + first: usize, + last: usize, +) -> f64 { + let raw_values = values.values(); + let mut result = result; + let mut previous = raw_values[first]; + for index in first + 1..=last { + if values.is_null(index) { + continue; + } + let current = raw_values[index]; + if current < previous { + result += previous; + } + previous = current; + } + result +} + +/// Locates the samples inside `[offset, offset + length)`, returning the first and last +/// non-null index together with the number of non-null slots. Returns `None` when the +/// window holds no sample. +fn valid_window_bounds( + values: &Float64Array, + offset: usize, + length: usize, +) -> Option<(usize, usize, usize)> { + let mut first = None; + let mut last = 0; + let mut count = 0; + for index in offset..offset + length { + if values.is_null(index) { + continue; + } + first.get_or_insert(index); + last = index; + count += 1; + } + first.map(|first| (first, last, count)) +} + fn extract_eval_timestamps( columnar_value: &ColumnarValue, func_name: &str, @@ -453,9 +525,11 @@ impl Display for ExtrapolatedRate { mod test { use datafusion::arrow::array::ArrayRef; + use datafusion::arrow::buffer::NullBuffer; use datafusion_common::ScalarValue; use super::*; + use crate::functions::test_util::TinyPrng; /// Range length is fixed to 5 fn extrapolated_rate_runner( @@ -607,6 +681,263 @@ mod test { assert_counter_windows_match_single(&values, &ranges); } + /// Builds a value array whose null slots keep a distinguishable raw payload, so a + /// function that reads the padding instead of the samples produces a different result. + fn values_with_nulls(values: Vec>, padding: f64) -> Arc { + let raw = values + .iter() + .map(|value| value.unwrap_or(padding)) + .collect::>(); + Arc::new(Float64Array::new( + raw.into(), + Some(NullBuffer::from_iter( + values.iter().map(|value| value.is_some()), + )), + )) + } + + fn nullable_rate_runner( + timestamps: Vec, + values: Arc, + ranges: Vec<(u32, u32)>, + eval_timestamps: Vec, + range_length: i64, + ) -> Vec> { + let ts_array = Arc::new(TimestampMillisecondArray::from_iter_values(timestamps)); + let ts_range = RangeArray::from_ranges(ts_array, ranges.clone()).unwrap(); + let value_range = RangeArray::from_ranges(values, ranges).unwrap(); + let input = vec![ + ColumnarValue::Array(Arc::new(ts_range.into_dict())), + ColumnarValue::Array(Arc::new(value_range.into_dict())), + ColumnarValue::Array(Arc::new(TimestampMillisecondArray::from_iter_values( + eval_timestamps, + ))), + ColumnarValue::Array(Arc::new(Int64Array::from(vec![range_length]))), + ]; + let output = extract_array( + &ExtrapolatedRate::::new(range_length) + .calc(&input) + .unwrap(), + ) + .unwrap(); + let output = output.as_any().downcast_ref::().unwrap(); + output.iter().collect() + } + + #[test] + fn rate_uses_samples_not_null_padding() { + // Samples are 1.0@0 and 4.0@3000; the padding under the null slots would add two more. + let output = nullable_rate_runner::( + vec![0, 1000, 2000, 3000], + values_with_nulls(vec![Some(1.0), None, None, Some(4.0)], 99.0), + vec![(0, 4)], + vec![3000], + 4000, + ); + + assert_eq!(output, vec![Some(1.0)]); + } + + #[test] + fn rate_returns_null_for_windows_without_enough_samples() { + let output = nullable_rate_runner::( + vec![0, 1000, 2000], + values_with_nulls(vec![None, Some(2.0), None], 7.0), + vec![(0, 3), (0, 2), (2, 1)], + vec![2000, 2000, 2000], + 4000, + ); + + assert_eq!(output, vec![None, None, None]); + } + + #[test] + fn increase_corrects_counter_reset_between_samples() { + // The sample sequence is 5.0 -> 3.0, one reset. Reading the padding would see + // 5.0 -> 100.0 -> 3.0 and charge the correction against the wrong value. + let output = nullable_rate_runner::( + vec![0, 1000, 2000], + values_with_nulls(vec![Some(5.0), None, Some(3.0)], 100.0), + vec![(0, 3)], + vec![2000], + 2000, + ); + + assert_eq!(output, vec![Some(3.0)]); + } + + #[test] + fn delta_extrapolates_from_sample_timestamps() { + // The window spans (-1000, 3000] but its samples only cover 1000..2000, so the + // extrapolation adds half an average interval on the leading side. + let output = nullable_rate_runner::( + vec![0, 1000, 2000, 3000], + values_with_nulls(vec![None, Some(2.0), Some(5.0), None], 42.0), + vec![(0, 4)], + vec![3000], + 4000, + ); + + 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) @@ -684,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], ); } @@ -715,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], ); } @@ -789,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], ); } @@ -809,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], ); } @@ -984,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], ); } @@ -1015,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/src/promql/src/functions/idelta.rs b/src/promql/src/functions/idelta.rs index 0772b0bf1e..35ce44ba27 100644 --- a/src/promql/src/functions/idelta.rs +++ b/src/promql/src/functions/idelta.rs @@ -101,12 +101,12 @@ impl IDelta { .unwrap() .values(); - let value_values = value_range.values(); - let value_values = value_values - .as_any() - .downcast_ref::() - .unwrap() - .values(); + let value_array = value_range.values(); + let value_array = value_array.as_any().downcast_ref::().unwrap(); + // A NULL field value means the series has no sample at that timestamp, so the last two + // samples are not necessarily the last two slots. + let has_nulls = value_array.null_count() > 0; + let value_values = value_array.values(); let mut result_builder = Float64Builder::with_capacity(ts_range.len()); @@ -122,20 +122,29 @@ impl IDelta { value_len )), )?; - if len < 2 { - result_builder.append_null(); - continue; - } + let (last_position, prev_position) = if has_nulls { + match last_two_samples(value_array, value_offset, len) { + Some(positions) => positions, + None => { + result_builder.append_null(); + continue; + } + } + } else { + if len < 2 { + result_builder.append_null(); + continue; + } + (len - 1, len - 2) + }; - let last_offset = ts_offset + len - 1; - let prev_offset = last_offset - 1; + let last_offset = ts_offset + last_position; + let prev_offset = ts_offset + prev_position; let sampled_interval = (ts_values[last_offset] - ts_values[prev_offset]) as f64 / 1000.0; - let last_value_offset = value_offset + len - 1; - let prev_value_offset = last_value_offset - 1; - let last_value = value_values[last_value_offset]; - let prev_value = value_values[prev_value_offset]; + let last_value = value_values[value_offset + last_position]; + let prev_value = value_values[value_offset + prev_position]; if !IS_RATE { result_builder.append_value(last_value - prev_value); @@ -157,6 +166,22 @@ impl IDelta { } } +/// Locates the last two samples inside `[offset, offset + len)`, returning their positions +/// relative to `offset`. Returns `None` when the window holds fewer than two samples. +fn last_two_samples(values: &Float64Array, offset: usize, len: usize) -> Option<(usize, usize)> { + let mut last = None; + for position in (0..len).rev() { + if values.is_null(offset + position) { + continue; + } + match last { + None => last = Some(position), + Some(last) => return Some((last, position)), + } + } + None +} + impl Display for IDelta { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { write!(f, "PromQL Idelta Function (is_rate: {IS_RATE})",) @@ -166,6 +191,8 @@ impl Display for IDelta { #[cfg(test)] mod test { + use datafusion::arrow::buffer::NullBuffer; + use super::*; use crate::functions::test_util::simple_range_udf_runner; @@ -207,4 +234,33 @@ mod test { vec![Some(0.5), Some(0.0), None, Some(3.0), None, None], ); } + + #[test] + fn idelta_uses_last_two_samples_not_last_two_slots() { + let ts_array = Arc::new(TimestampMillisecondArray::from_iter_values([ + 0i64, 1000, 2000, 3000, + ])); + // Samples are 1.0@0 and 4.0@1000; the trailing slots only carry padding. + let values_array = Arc::new(Float64Array::new( + vec![1.0, 4.0, 100.0, 200.0].into(), + Some(NullBuffer::from_iter([true, true, false, false])), + )); + let ranges = [(0, 4), (2, 2), (3, 1)]; + + simple_range_udf_runner( + IDelta::::scalar_udf(), + RangeArray::from_ranges(ts_array.clone(), ranges).unwrap(), + RangeArray::from_ranges(values_array.clone(), ranges).unwrap(), + vec![], + vec![Some(3.0), None, None], + ); + + simple_range_udf_runner( + IDelta::::scalar_udf(), + RangeArray::from_ranges(ts_array, ranges).unwrap(), + RangeArray::from_ranges(values_array, ranges).unwrap(), + vec![], + vec![Some(3.0), None, None], + ); + } } diff --git a/src/promql/src/functions/quantile.rs b/src/promql/src/functions/quantile.rs index 35f0c1aa81..3794057da5 100644 --- a/src/promql/src/functions/quantile.rs +++ b/src/promql/src/functions/quantile.rs @@ -93,14 +93,14 @@ impl QuantileOverTime { )), )?; - let all_values = value_range - .values() - .as_any() - .downcast_ref::() - .unwrap() - .values(); + let value_array = value_range.values(); + let value_array = value_array.as_any().downcast_ref::().unwrap(); + // A NULL field value means the series has no sample at that timestamp, so a window's + // samples are not simply its slots. + let has_nulls = value_array.null_count() > 0; let mut result_builder = Float64Builder::with_capacity(ts_range.len()); let mut scratch = Vec::new(); + let mut samples = Vec::new(); match quantile_col { ColumnarValue::Scalar(quantile_scalar) => { @@ -125,11 +125,14 @@ impl QuantileOverTime { )), )?; - match quantile_with_scratch( - &all_values[value_offset..value_offset + value_len], - quantile, - &mut scratch, - ) { + let window = window_samples( + value_array, + has_nulls, + value_offset, + value_len, + &mut samples, + ); + match window_quantile(window, quantile, &mut scratch) { Some(value) => result_builder.append_value(value), None => result_builder.append_null(), } @@ -173,11 +176,14 @@ impl QuantileOverTime { } else { quantile_array.value(index) }; - match quantile_with_scratch( - &all_values[value_offset..value_offset + value_len], - quantile, - &mut scratch, - ) { + let window = window_samples( + value_array, + has_nulls, + value_offset, + value_len, + &mut samples, + ); + match window_quantile(window, quantile, &mut scratch) { Some(value) => result_builder.append_value(value), None => result_builder.append_null(), } @@ -190,6 +196,40 @@ impl QuantileOverTime { } } +/// Returns the samples of the window `[offset, offset + len)`, collecting the non-null ones +/// into `samples` when the backing array has nulls and borrowing the slice otherwise. +fn window_samples<'a>( + values: &'a Float64Array, + has_nulls: bool, + offset: usize, + len: usize, + samples: &'a mut Vec, +) -> &'a [f64] { + let raw_values = values.values(); + if !has_nulls { + return &raw_values[offset..offset + len]; + } + samples.clear(); + samples.extend( + (offset..offset + len) + .filter(|index| values.is_valid(*index)) + .map(|index| raw_values[index]), + ); + samples +} + +/// Quantile of one range window, or `None` when the window holds no sample. +/// +/// Prometheus returns an empty vector for a range without float samples rather than the NaN +/// that [`quantile_impl`] yields for an empty slice, so the emptiness check belongs here and +/// not in the shared kernel. +fn window_quantile(values: &[f64], quantile: f64, scratch: &mut Vec) -> Option { + if values.is_empty() { + return None; + } + quantile_with_scratch(values, quantile, scratch) +} + /// Refer to pub(crate) fn quantile_impl(values: &[f64], quantile: f64) -> Option { let mut scratch = Vec::new(); @@ -226,6 +266,9 @@ fn quantile_with_scratch(values: &[f64], quantile: f64, scratch: &mut Vec) #[cfg(test)] mod tests { + use datafusion::arrow::array::TimestampMillisecondArray; + use datafusion::arrow::buffer::NullBuffer; + use super::*; #[test] @@ -276,4 +319,64 @@ mod tests { let q = 0.25; assert_eq!(quantile_impl(values, q).unwrap(), 2.0); } + + #[test] + fn quantile_over_time_ranks_samples_only() { + let ts_array = Arc::new(TimestampMillisecondArray::from_iter_values([ + 0i64, 1000, 2000, + ])); + // Samples are 1.0 and 4.0; ranking the padding too would pull the median down. + let values_array = Arc::new(Float64Array::new( + vec![1.0, -100.0, 4.0].into(), + Some(NullBuffer::from_iter([true, false, true])), + )); + // The second window holds no sample, the third holds no slot at all. + let ranges = [(0, 3), (1, 1), (3, 0)]; + + let input = vec![ + ColumnarValue::Array(Arc::new( + RangeArray::from_ranges(ts_array, ranges) + .unwrap() + .into_dict(), + )), + ColumnarValue::Array(Arc::new( + RangeArray::from_ranges(values_array, ranges) + .unwrap() + .into_dict(), + )), + ColumnarValue::Scalar(ScalarValue::Float64(Some(0.5))), + ]; + let output = extract_array(&QuantileOverTime::quantile_over_time(&input).unwrap()).unwrap(); + let output = output.as_any().downcast_ref::().unwrap(); + + assert_eq!( + output.iter().collect::>(), + vec![Some(2.5), None, None] + ); + } + + #[test] + fn quantile_over_time_keeps_nan_for_an_invalid_quantile() { + let ts_array = Arc::new(TimestampMillisecondArray::from_iter_values([0i64, 1000])); + let values_array = Arc::new(Float64Array::from_iter_values([1.0, 4.0])); + let ranges = [(0, 2)]; + + let input = vec![ + ColumnarValue::Array(Arc::new( + RangeArray::from_ranges(ts_array, ranges) + .unwrap() + .into_dict(), + )), + ColumnarValue::Array(Arc::new( + RangeArray::from_ranges(values_array, ranges) + .unwrap() + .into_dict(), + )), + ColumnarValue::Scalar(ScalarValue::Float64(None)), + ]; + let output = extract_array(&QuantileOverTime::quantile_over_time(&input).unwrap()).unwrap(); + let output = output.as_any().downcast_ref::().unwrap(); + + assert!(output.value(0).is_nan()); + } } diff --git a/src/promql/src/functions/resets.rs b/src/promql/src/functions/resets.rs index 66e83cde31..7a5c0db71c 100644 --- a/src/promql/src/functions/resets.rs +++ b/src/promql/src/functions/resets.rs @@ -43,8 +43,8 @@ mod test { use super::*; use crate::functions::test_util::{ - self, STALE_NAN, TinyPrng, assert_execution_error, build_test_range_arrays, - invoke_range_udf, simple_range_udf_runner, + self, STALE_NAN, assert_execution_error, build_test_range_arrays, invoke_range_udf, + simple_range_udf_runner, }; use crate::range_array::RangeArray; @@ -151,43 +151,17 @@ mod test { resets_oracle, Resets::scalar_udf(), ); - assert_eq!(expected, vec![Some(1.0), Some(0.0)]); + assert_eq!(expected, vec![Some(0.0), None]); } #[test] fn resets_range_array_seeded_differential() { - let mut prng = TinyPrng(0x2f6e_2b1d_834a_90c5); - let raw_values = (0..48) - .map(|_| match prng.next_index(12) { - 0 => -0.0, - 1 => 0.0, - 2 => -2.0, - 3 => -1.0, - 4 => 1.0, - 5 => 2.0, - 6 => f64::INFINITY, - 7 => f64::NEG_INFINITY, - 8 | 9 => f64::NAN, - _ => STALE_NAN, - }) - .collect::>(); - let values = raw_values.iter().copied().map(Some).collect(); - let mut timestamp_ranges = Vec::new(); - let mut value_ranges = Vec::new(); - for _ in 0..32 { - let length = prng.next_index(13) as u32; - timestamp_ranges.push((prng.next_index(65 - length as usize) as u32, length)); - value_ranges.push((prng.next_index(49 - length as usize) as u32, length)); - } + test_util::run_seeded_differential(resets_oracle, Resets::scalar_udf(), false); + } - test_util::run_oracle_ranges( - values, - raw_values, - timestamp_ranges, - value_ranges, - resets_oracle, - Resets::scalar_udf(), - ); + #[test] + fn resets_range_array_seeded_differential_with_nulls() { + test_util::run_seeded_differential(resets_oracle, Resets::scalar_udf(), true); } #[test] @@ -263,6 +237,8 @@ mod test { vec![Some(0.0), Some(0.0), Some(0.0)], ); + // The raw payload under the null slot would look like a reset; the sample sequence is + // 10 -> 10, which is none. The last window holds no sample at all. let values = Arc::new(Float64Array::new( vec![10.0, 7.0, 10.0].into(), Some(NullBuffer::from(vec![true, false, true])), @@ -276,10 +252,10 @@ mod test { ])); simple_range_udf_runner( Resets::scalar_udf(), - RangeArray::from_ranges(timestamps, [(0, 3)]).unwrap(), - RangeArray::from_ranges(values, [(0, 3)]).unwrap(), + RangeArray::from_ranges(timestamps, [(0, 3), (1, 1)]).unwrap(), + RangeArray::from_ranges(values, [(0, 3), (1, 1)]).unwrap(), vec![], - vec![Some(1.0)], + vec![Some(0.0), None], ); } } diff --git a/src/promql/src/functions/test_util.rs b/src/promql/src/functions/test_util.rs index e52f52e3bb..6148e06c39 100644 --- a/src/promql/src/functions/test_util.rs +++ b/src/promql/src/functions/test_util.rs @@ -15,6 +15,7 @@ use std::sync::Arc; use datafusion::arrow::array::{Float64Array, TimestampMillisecondArray}; +use datafusion::arrow::buffer::NullBuffer; use datafusion::common::DataFusionError; use datafusion::logical_expr::ScalarUDF; use datafusion::physical_plan::ColumnarValue; @@ -144,12 +145,54 @@ impl TinyPrng { } } -/// Run the oracle-based differential test: build range arrays, verify raw bits and -/// null validity, invoke the UDF via [`simple_range_udf_runner`], and return the -/// expected values for further assertions. +/// Run [`run_oracle_ranges`] over a seeded mix of signed zeros, infinities, NaNs and stale +/// markers spread across random windows. `nullable` marks roughly a quarter of the slots as +/// null while keeping their raw payload, which moves the UDF off its null-free fast path. +pub fn run_seeded_differential(oracle: fn(&[f64]) -> Option, udf: ScalarUDF, nullable: bool) { + let mut prng = TinyPrng(0x2f6e_2b1d_834a_90c5); + let raw_values = (0..48) + .map(|_| match prng.next_index(12) { + 0 => -0.0, + 1 => 0.0, + 2 => -2.0, + 3 => -1.0, + 4 => 1.0, + 5 => 2.0, + 6 => f64::INFINITY, + 7 => f64::NEG_INFINITY, + 8 | 9 => f64::NAN, + _ => STALE_NAN, + }) + .collect::>(); + let values = raw_values + .iter() + .map(|value| (!nullable || prng.next_index(4) != 0).then_some(*value)) + .collect::>(); + let mut timestamp_ranges = Vec::new(); + let mut value_ranges = Vec::new(); + for _ in 0..32 { + let length = prng.next_index(13) as u32; + timestamp_ranges.push((prng.next_index(65 - length as usize) as u32, length)); + value_ranges.push((prng.next_index(49 - length as usize) as u32, length)); + } + + run_oracle_ranges( + values, + raw_values, + timestamp_ranges, + value_ranges, + oracle, + udf, + ); +} + +/// Run the oracle-based differential test: build range arrays, invoke the UDF via +/// [`simple_range_udf_runner`], and return the expected values for further assertions. /// /// `oracle` is the behavior-specific function (e.g. `changes_oracle` or `resets_oracle`) -/// that computes the expected count for a slice of raw f64 values. +/// that computes the expected count for a window's samples. Null slots are dropped before +/// the oracle runs: a NULL field value means the series has no sample at that timestamp, +/// so the raw payload underneath it is not a sample. pub fn run_oracle_ranges( values: Vec>, raw_values: Vec, @@ -166,29 +209,37 @@ pub fn run_oracle_ranges( .all(|((_, timestamp_length), (_, value_length))| timestamp_length == value_length) ); assert_eq!(values.len(), raw_values.len()); + // `values` marks the null slots, `raw_values` carries the payload of every slot including + // the null ones. The two must agree wherever a sample exists. + for (index, (value, raw)) in values.iter().zip(&raw_values).enumerate() { + assert!( + value.is_none_or(|value| value.to_bits() == raw.to_bits()), + "values[{index}] and raw_values[{index}] disagree on the sample" + ); + } let nulls = values.iter().map(Option::is_none).collect::>(); let expected = value_ranges .iter() - .map(|(offset, length)| oracle(&raw_values[*offset as usize..(*offset + *length) as usize])) + .map(|(offset, length)| { + let samples = (*offset as usize..(*offset + *length) as usize) + .filter(|index| !nulls[*index]) + .map(|index| raw_values[index]) + .collect::>(); + oracle(&samples) + }) .collect::>(); let timestamp_values = (0..64) .map(|value| Some(i64::from(value) * 1_000)) .collect::>(); let timestamp_array = Arc::new(TimestampMillisecondArray::from_iter(timestamp_values)); - let value_array = Arc::new(Float64Array::from_iter(values)); - for (index, ((actual, expected), is_null)) in value_array - .values() - .iter() - .zip(&raw_values) - .zip(nulls) - .enumerate() - { - assert_eq!(actual.to_bits(), expected.to_bits()); - assert_eq!(value_array.is_null(index), is_null); - } - + // Build from the raw payload plus a null bitmap instead of `from_iter`, which would zero + // the null slots and hide a function that reads them. + let value_array = Arc::new(Float64Array::new( + raw_values.clone().into(), + Some(NullBuffer::from_iter(nulls.iter().map(|null| !null))), + )); let timestamp_ranges = RangeArray::from_ranges(timestamp_array, timestamp_ranges).unwrap(); let value_ranges = RangeArray::from_ranges(value_array, value_ranges).unwrap(); simple_range_udf_runner( diff --git a/src/query/src/promql/planner.rs b/src/query/src/promql/planner.rs index 8e37c55dae..6ca04ffca7 100644 --- a/src/query/src/promql/planner.rs +++ b/src/query/src/promql/planner.rs @@ -2195,8 +2195,6 @@ impl PromPlanner { ), }) }; - let preserve_any_value = - Self::field_columns_are_alternative_samples(input.schema(), &self.ctx.field_columns); let (mut func_exprs, new_tags) = self.create_function_expr( func, args.literals.clone(), @@ -2211,10 +2209,15 @@ impl PromPlanner { func_exprs.push(tsid_col); } + // A row survives as long as one field column produced a sample, and the fields without + // one stay NULL, which is the shape a selector already emits. Requiring every field to + // be non-NULL would drop one field's samples because another field has none in the same + // window — the reason alternative float/histogram columns already needed this form. A + // single field column reduces to the same predicate either way. let builder = LogicalPlanBuilder::from(input) .project(func_exprs) .context(DataFusionPlanningSnafu)? - .filter(self.create_empty_values_filter_expr(preserve_any_value)?) + .filter(self.create_empty_values_filter_expr(true)?) .context(DataFusionPlanningSnafu)?; let builder = match func.name { diff --git a/tests/cases/standalone/common/promql/null_samples.result b/tests/cases/standalone/common/promql/null_samples.result new file mode 100644 index 0000000000..8e0f159fe1 --- /dev/null +++ b/tests/cases/standalone/common/promql/null_samples.result @@ -0,0 +1,252 @@ +-- NULL fields are missing samples, not zero-valued samples. +CREATE TABLE null_samples ( + ts TIMESTAMP(3) TIME INDEX, + host STRING PRIMARY KEY, + val DOUBLE, +); + +Affected Rows: 0 + +INSERT INTO null_samples VALUES + (0, 'a', 1.0), + (1000, 'a', NULL), + (2000, 'a', NULL), + (3000, 'a', 4.0), + (0, 'b', NULL), + (1000, 'b', NULL), + (2000, 'b', NULL), + (3000, 'b', NULL); + +Affected Rows: 8 + +-- At t=2 the trailing NULLs must not hide 1; at t=3 count must be 2. +-- At t=7 the left-open window is empty. Valid results must disappear. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (2, 7, '1s') count_over_time(null_samples{host="a"}[4s]); + ++---------------------+------------------------------------+------+ +| ts | prom_count_over_time(ts_range,val) | host | ++---------------------+------------------------------------+------+ +| 1970-01-01T00:00:02 | 1.0 | a | +| 1970-01-01T00:00:03 | 2.0 | a | +| 1970-01-01T00:00:04 | 1.0 | a | +| 1970-01-01T00:00:05 | 1.0 | a | +| 1970-01-01T00:00:06 | 1.0 | a | ++---------------------+------------------------------------+------+ + +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (2, 7, '1s') last_over_time(null_samples{host="a"}[4s]); + ++---------------------+-----------------------------------+------+ +| ts | prom_last_over_time(ts_range,val) | host | ++---------------------+-----------------------------------+------+ +| 1970-01-01T00:00:02 | 1.0 | a | +| 1970-01-01T00:00:03 | 4.0 | a | +| 1970-01-01T00:00:04 | 4.0 | a | +| 1970-01-01T00:00:05 | 4.0 | a | +| 1970-01-01T00:00:06 | 4.0 | a | ++---------------------+-----------------------------------+------+ + +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (2, 7, '1s') present_over_time(null_samples{host="a"}[4s]); + ++---------------------+--------------------------------------+------+ +| ts | prom_present_over_time(ts_range,val) | host | ++---------------------+--------------------------------------+------+ +| 1970-01-01T00:00:02 | 1.0 | a | +| 1970-01-01T00:00:03 | 1.0 | a | +| 1970-01-01T00:00:04 | 1.0 | a | +| 1970-01-01T00:00:05 | 1.0 | a | +| 1970-01-01T00:00:06 | 1.0 | a | ++---------------------+--------------------------------------+------+ + +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (2, 7, '1s') absent_over_time(null_samples{host="a"}[4s]); + ++---------------------+-------------------------------------+------+ +| ts | prom_absent_over_time(ts_range,val) | host | ++---------------------+-------------------------------------+------+ +| 1970-01-01T00:00:07 | 1.0 | a | ++---------------------+-------------------------------------+------+ + +-- All-NULL windows have no samples: only absent_over_time returns 1. +TQL EVAL (3, 3, '1s') count_over_time(null_samples{host="b"}[4s]); + +++ +++ + +TQL EVAL (3, 3, '1s') last_over_time(null_samples{host="b"}[4s]); + +++ +++ + +TQL EVAL (3, 3, '1s') present_over_time(null_samples{host="b"}[4s]); + +++ +++ + +TQL EVAL (3, 3, '1s') absent_over_time(null_samples{host="b"}[4s]); + ++---------------------+-------------------------------------+------+ +| ts | prom_absent_over_time(ts_range,val) | host | ++---------------------+-------------------------------------+------+ +| 1970-01-01T00:00:03 | 1.0 | b | ++---------------------+-------------------------------------+------+ + +-- Every function below sees the same two samples, 1.0 at 0s and 4.0 at 3s. +TQL EVAL (3, 3, '1s') rate(null_samples{host="a"}[4s]); + ++---------------------+----------------------------------------+------+ +| ts | prom_rate(ts_range,val,ts,Int64(4000)) | host | ++---------------------+----------------------------------------+------+ +| 1970-01-01T00:00:03 | 1.0 | a | ++---------------------+----------------------------------------+------+ + +TQL EVAL (3, 3, '1s') increase(null_samples{host="a"}[4s]); + ++---------------------+--------------------------------------------+------+ +| ts | prom_increase(ts_range,val,ts,Int64(4000)) | host | ++---------------------+--------------------------------------------+------+ +| 1970-01-01T00:00:03 | 4.0 | a | ++---------------------+--------------------------------------------+------+ + +TQL EVAL (3, 3, '1s') delta(null_samples{host="a"}[4s]); + ++---------------------+-----------------------------------------+------+ +| ts | prom_delta(ts_range,val,ts,Int64(4000)) | host | ++---------------------+-----------------------------------------+------+ +| 1970-01-01T00:00:03 | 4.0 | a | ++---------------------+-----------------------------------------+------+ + +TQL EVAL (3, 3, '1s') idelta(null_samples{host="a"}[4s]); + ++---------------------+---------------------------+------+ +| ts | prom_idelta(ts_range,val) | host | ++---------------------+---------------------------+------+ +| 1970-01-01T00:00:03 | 3.0 | a | ++---------------------+---------------------------+------+ + +TQL EVAL (3, 3, '1s') irate(null_samples{host="a"}[4s]); + ++---------------------+--------------------------+------+ +| ts | prom_irate(ts_range,val) | host | ++---------------------+--------------------------+------+ +| 1970-01-01T00:00:03 | 1.0 | a | ++---------------------+--------------------------+------+ + +TQL EVAL (3, 3, '1s') changes(null_samples{host="a"}[4s]); + ++---------------------+----------------------------+------+ +| ts | prom_changes(ts_range,val) | host | ++---------------------+----------------------------+------+ +| 1970-01-01T00:00:03 | 1.0 | a | ++---------------------+----------------------------+------+ + +TQL EVAL (3, 3, '1s') resets(null_samples{host="a"}[4s]); + ++---------------------+---------------------------+------+ +| ts | prom_resets(ts_range,val) | host | ++---------------------+---------------------------+------+ +| 1970-01-01T00:00:03 | 0.0 | a | ++---------------------+---------------------------+------+ + +TQL EVAL (3, 3, '1s') avg_over_time(null_samples{host="a"}[4s]); + ++---------------------+----------------------------------+------+ +| ts | prom_avg_over_time(ts_range,val) | host | ++---------------------+----------------------------------+------+ +| 1970-01-01T00:00:03 | 2.5 | a | ++---------------------+----------------------------------+------+ + +TQL EVAL (3, 3, '1s') stddev_over_time(null_samples{host="a"}[4s]); + ++---------------------+-------------------------------------+------+ +| ts | prom_stddev_over_time(ts_range,val) | host | ++---------------------+-------------------------------------+------+ +| 1970-01-01T00:00:03 | 1.5 | a | ++---------------------+-------------------------------------+------+ + +TQL EVAL (3, 3, '1s') stdvar_over_time(null_samples{host="a"}[4s]); + ++---------------------+-------------------------------------+------+ +| ts | prom_stdvar_over_time(ts_range,val) | host | ++---------------------+-------------------------------------+------+ +| 1970-01-01T00:00:03 | 2.25 | a | ++---------------------+-------------------------------------+------+ + +TQL EVAL (3, 3, '1s') quantile_over_time(0.5, null_samples{host="a"}[4s]); + ++---------------------+----------------------------------------------------+------+ +| ts | prom_quantile_over_time(ts_range,val,Float64(0.5)) | host | ++---------------------+----------------------------------------------------+------+ +| 1970-01-01T00:00:03 | 2.5 | a | ++---------------------+----------------------------------------------------+------+ + +-- A window whose only slot is NULL holds no sample, like a window with no row. +TQL EVAL (1, 1, '1s') rate(null_samples{host="a"}[1s]); + +++ +++ + +-- Prometheus returns an empty vector for a range without samples, not NaN. +TQL EVAL (1, 1, '1s') quantile_over_time(0.5, null_samples{host="a"}[1s]); + +++ +++ + +-- `b` never has a sample, so it must not reach the aggregation. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 15, '1s') avg by (host) (rate(null_samples[4s])); + ++------+---------------------+---------------------------------------------+ +| host | ts | avg(prom_rate(ts_range,val,ts,Int64(4000))) | ++------+---------------------+---------------------------------------------+ +| a | 1970-01-01T00:00:03 | 1.0 | ++------+---------------------+---------------------------------------------+ + +DROP TABLE null_samples; + +Affected Rows: 0 + +CREATE TABLE multi_field ( + ts TIMESTAMP(3) TIME INDEX, + host STRING PRIMARY KEY, + f1 DOUBLE, + f2 DOUBLE +); + +Affected Rows: 0 + +INSERT INTO multi_field VALUES + (0, 'a', 1.0, 10.0), + (1000, 'a', 2.0, NULL), + (2000, 'a', 3.0, NULL), + (3000, 'a', 4.0, 40.0); + +Affected Rows: 4 + +-- f1 has two samples in this window and f2 only one, so f1 keeps its result while f2 is NULL. +-- Dropping the whole row would take f1's samples with it. +-- SQLNESS SORT_RESULT 3 1 +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.375 | | a | ++---------------------+---------------------------------------+---------------------------------------+------+ + +-- A selector emits the same shape: the row stays, the field without a sample is NULL. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (1, 1, '1s') multi_field; + ++---------------------+------+-----+----+ +| ts | host | f1 | f2 | ++---------------------+------+-----+----+ +| 1970-01-01T00:00:01 | a | 2.0 | | ++---------------------+------+-----+----+ + +DROP TABLE multi_field; + +Affected Rows: 0 + diff --git a/tests/cases/standalone/common/promql/null_samples.sql b/tests/cases/standalone/common/promql/null_samples.sql new file mode 100644 index 0000000000..d0a02158a4 --- /dev/null +++ b/tests/cases/standalone/common/promql/null_samples.sql @@ -0,0 +1,92 @@ +-- NULL fields are missing samples, not zero-valued samples. +CREATE TABLE null_samples ( + ts TIMESTAMP(3) TIME INDEX, + host STRING PRIMARY KEY, + val DOUBLE, +); + +INSERT INTO null_samples VALUES + (0, 'a', 1.0), + (1000, 'a', NULL), + (2000, 'a', NULL), + (3000, 'a', 4.0), + (0, 'b', NULL), + (1000, 'b', NULL), + (2000, 'b', NULL), + (3000, 'b', NULL); + +-- At t=2 the trailing NULLs must not hide 1; at t=3 count must be 2. +-- At t=7 the left-open window is empty. Valid results must disappear. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (2, 7, '1s') count_over_time(null_samples{host="a"}[4s]); +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (2, 7, '1s') last_over_time(null_samples{host="a"}[4s]); +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (2, 7, '1s') present_over_time(null_samples{host="a"}[4s]); +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (2, 7, '1s') absent_over_time(null_samples{host="a"}[4s]); + +-- All-NULL windows have no samples: only absent_over_time returns 1. +TQL EVAL (3, 3, '1s') count_over_time(null_samples{host="b"}[4s]); +TQL EVAL (3, 3, '1s') last_over_time(null_samples{host="b"}[4s]); +TQL EVAL (3, 3, '1s') present_over_time(null_samples{host="b"}[4s]); +TQL EVAL (3, 3, '1s') absent_over_time(null_samples{host="b"}[4s]); + +-- Every function below sees the same two samples, 1.0 at 0s and 4.0 at 3s. +TQL EVAL (3, 3, '1s') rate(null_samples{host="a"}[4s]); + +TQL EVAL (3, 3, '1s') increase(null_samples{host="a"}[4s]); + +TQL EVAL (3, 3, '1s') delta(null_samples{host="a"}[4s]); + +TQL EVAL (3, 3, '1s') idelta(null_samples{host="a"}[4s]); + +TQL EVAL (3, 3, '1s') irate(null_samples{host="a"}[4s]); + +TQL EVAL (3, 3, '1s') changes(null_samples{host="a"}[4s]); + +TQL EVAL (3, 3, '1s') resets(null_samples{host="a"}[4s]); + +TQL EVAL (3, 3, '1s') avg_over_time(null_samples{host="a"}[4s]); + +TQL EVAL (3, 3, '1s') stddev_over_time(null_samples{host="a"}[4s]); + +TQL EVAL (3, 3, '1s') stdvar_over_time(null_samples{host="a"}[4s]); + +TQL EVAL (3, 3, '1s') quantile_over_time(0.5, null_samples{host="a"}[4s]); + +-- A window whose only slot is NULL holds no sample, like a window with no row. +TQL EVAL (1, 1, '1s') rate(null_samples{host="a"}[1s]); + +-- Prometheus returns an empty vector for a range without samples, not NaN. +TQL EVAL (1, 1, '1s') quantile_over_time(0.5, null_samples{host="a"}[1s]); + +-- `b` never has a sample, so it must not reach the aggregation. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 15, '1s') avg by (host) (rate(null_samples[4s])); + +DROP TABLE null_samples; + +CREATE TABLE multi_field ( + ts TIMESTAMP(3) TIME INDEX, + host STRING PRIMARY KEY, + f1 DOUBLE, + f2 DOUBLE +); + +INSERT INTO multi_field VALUES + (0, 'a', 1.0, 10.0), + (1000, 'a', 2.0, NULL), + (2000, 'a', 3.0, NULL), + (3000, 'a', 4.0, 40.0); + +-- f1 has two samples in this window and f2 only one, so f1 keeps its result while f2 is NULL. +-- Dropping the whole row would take f1's samples with it. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (1, 1, '1s') rate(multi_field[4s]); + +-- A selector emits the same shape: the row stays, the field without a sample is NULL. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (1, 1, '1s') multi_field; + +DROP TABLE multi_field; diff --git a/tests/cases/standalone/common/promql/range_presence_null.result b/tests/cases/standalone/common/promql/range_presence_null.result deleted file mode 100644 index d0f28ad9fb..0000000000 --- a/tests/cases/standalone/common/promql/range_presence_null.result +++ /dev/null @@ -1,99 +0,0 @@ --- NULL fields are missing samples, not zero-valued samples. -CREATE TABLE range_presence_null ( - ts TIMESTAMP(3) TIME INDEX, - host STRING PRIMARY KEY, - val DOUBLE, -); - -Affected Rows: 0 - -INSERT INTO range_presence_null VALUES - (0, 'a', 1.0), - (1000, 'a', NULL), - (2000, 'a', NULL), - (3000, 'a', 4.0), - (0, 'b', NULL), - (1000, 'b', NULL), - (2000, 'b', NULL), - (3000, 'b', NULL); - -Affected Rows: 8 - --- At t=2 the trailing NULLs must not hide 1; at t=3 count must be 2. --- At t=7 the left-open window is empty. Valid results must disappear. --- SQLNESS SORT_RESULT 3 1 -TQL EVAL (2, 7, '1s') count_over_time(range_presence_null{host="a"}[4s]); - -+---------------------+------------------------------------+------+ -| ts | prom_count_over_time(ts_range,val) | host | -+---------------------+------------------------------------+------+ -| 1970-01-01T00:00:02 | 1.0 | a | -| 1970-01-01T00:00:03 | 2.0 | a | -| 1970-01-01T00:00:04 | 1.0 | a | -| 1970-01-01T00:00:05 | 1.0 | a | -| 1970-01-01T00:00:06 | 1.0 | a | -+---------------------+------------------------------------+------+ - --- SQLNESS SORT_RESULT 3 1 -TQL EVAL (2, 7, '1s') last_over_time(range_presence_null{host="a"}[4s]); - -+---------------------+-----------------------------------+------+ -| ts | prom_last_over_time(ts_range,val) | host | -+---------------------+-----------------------------------+------+ -| 1970-01-01T00:00:02 | 1.0 | a | -| 1970-01-01T00:00:03 | 4.0 | a | -| 1970-01-01T00:00:04 | 4.0 | a | -| 1970-01-01T00:00:05 | 4.0 | a | -| 1970-01-01T00:00:06 | 4.0 | a | -+---------------------+-----------------------------------+------+ - --- SQLNESS SORT_RESULT 3 1 -TQL EVAL (2, 7, '1s') present_over_time(range_presence_null{host="a"}[4s]); - -+---------------------+--------------------------------------+------+ -| ts | prom_present_over_time(ts_range,val) | host | -+---------------------+--------------------------------------+------+ -| 1970-01-01T00:00:02 | 1.0 | a | -| 1970-01-01T00:00:03 | 1.0 | a | -| 1970-01-01T00:00:04 | 1.0 | a | -| 1970-01-01T00:00:05 | 1.0 | a | -| 1970-01-01T00:00:06 | 1.0 | a | -+---------------------+--------------------------------------+------+ - --- SQLNESS SORT_RESULT 3 1 -TQL EVAL (2, 7, '1s') absent_over_time(range_presence_null{host="a"}[4s]); - -+---------------------+-------------------------------------+------+ -| ts | prom_absent_over_time(ts_range,val) | host | -+---------------------+-------------------------------------+------+ -| 1970-01-01T00:00:07 | 1.0 | a | -+---------------------+-------------------------------------+------+ - --- All-NULL windows have no samples: only absent_over_time returns 1. -TQL EVAL (3, 3, '1s') count_over_time(range_presence_null{host="b"}[4s]); - -++ -++ - -TQL EVAL (3, 3, '1s') last_over_time(range_presence_null{host="b"}[4s]); - -++ -++ - -TQL EVAL (3, 3, '1s') present_over_time(range_presence_null{host="b"}[4s]); - -++ -++ - -TQL EVAL (3, 3, '1s') absent_over_time(range_presence_null{host="b"}[4s]); - -+---------------------+-------------------------------------+------+ -| ts | prom_absent_over_time(ts_range,val) | host | -+---------------------+-------------------------------------+------+ -| 1970-01-01T00:00:03 | 1.0 | b | -+---------------------+-------------------------------------+------+ - -DROP TABLE range_presence_null; - -Affected Rows: 0 - diff --git a/tests/cases/standalone/common/promql/range_presence_null.sql b/tests/cases/standalone/common/promql/range_presence_null.sql deleted file mode 100644 index e442679b2b..0000000000 --- a/tests/cases/standalone/common/promql/range_presence_null.sql +++ /dev/null @@ -1,35 +0,0 @@ --- NULL fields are missing samples, not zero-valued samples. -CREATE TABLE range_presence_null ( - ts TIMESTAMP(3) TIME INDEX, - host STRING PRIMARY KEY, - val DOUBLE, -); - -INSERT INTO range_presence_null VALUES - (0, 'a', 1.0), - (1000, 'a', NULL), - (2000, 'a', NULL), - (3000, 'a', 4.0), - (0, 'b', NULL), - (1000, 'b', NULL), - (2000, 'b', NULL), - (3000, 'b', NULL); - --- At t=2 the trailing NULLs must not hide 1; at t=3 count must be 2. --- At t=7 the left-open window is empty. Valid results must disappear. --- SQLNESS SORT_RESULT 3 1 -TQL EVAL (2, 7, '1s') count_over_time(range_presence_null{host="a"}[4s]); --- SQLNESS SORT_RESULT 3 1 -TQL EVAL (2, 7, '1s') last_over_time(range_presence_null{host="a"}[4s]); --- SQLNESS SORT_RESULT 3 1 -TQL EVAL (2, 7, '1s') present_over_time(range_presence_null{host="a"}[4s]); --- SQLNESS SORT_RESULT 3 1 -TQL EVAL (2, 7, '1s') absent_over_time(range_presence_null{host="a"}[4s]); - --- All-NULL windows have no samples: only absent_over_time returns 1. -TQL EVAL (3, 3, '1s') count_over_time(range_presence_null{host="b"}[4s]); -TQL EVAL (3, 3, '1s') last_over_time(range_presence_null{host="b"}[4s]); -TQL EVAL (3, 3, '1s') present_over_time(range_presence_null{host="b"}[4s]); -TQL EVAL (3, 3, '1s') absent_over_time(range_presence_null{host="b"}[4s]); - -DROP TABLE range_presence_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;