Skip to main content

promql/functions/
aggr_over_time.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use std::sync::Arc;
16
17use common_macro::range_fn;
18use datafusion::arrow::array::{Float64Array, TimestampMillisecondArray};
19use datafusion::common::DataFusionError;
20use datafusion::logical_expr::{ScalarUDF, Volatility};
21use datafusion::physical_plan::ColumnarValue;
22use datatypes::arrow::array::Array;
23use datatypes::arrow::compute;
24use datatypes::arrow::datatypes::DataType;
25
26use crate::functions::{compensated_sum_inc, extract_array};
27use crate::range_array::RangeArray;
28
29/// The average value of all points in the specified interval.
30#[range_fn(
31    name = AvgOverTime,
32    ret = Float64Array,
33    display_name = prom_avg_over_time
34)]
35pub fn avg_over_time(_: &TimestampMillisecondArray, values: &Float64Array) -> Option<f64> {
36    compute::sum(values).map(|result| result / values.len() as f64)
37}
38
39/// The minimum value of all points in the specified interval.
40#[range_fn(
41    name = MinOverTime,
42    ret = Float64Array,
43    display_name = prom_min_over_time
44)]
45pub fn min_over_time(_: &TimestampMillisecondArray, values: &Float64Array) -> Option<f64> {
46    let mut valid_values = values.iter().flatten();
47    let mut min = valid_values.next()?;
48    for value in valid_values {
49        if value < min || min.is_nan() {
50            min = value;
51        }
52    }
53    Some(min)
54}
55
56/// The maximum value of all points in the specified interval.
57#[range_fn(
58    name = MaxOverTime,
59    ret = Float64Array,
60    display_name = prom_max_over_time
61)]
62pub fn max_over_time(_: &TimestampMillisecondArray, values: &Float64Array) -> Option<f64> {
63    let mut valid_values = values.iter().flatten();
64    let mut max = valid_values.next()?;
65    for value in valid_values {
66        if value > max || max.is_nan() {
67            max = value;
68        }
69    }
70    Some(max)
71}
72
73/// The sum of all values in the specified interval.
74#[range_fn(
75    name = SumOverTime,
76    ret = Float64Array,
77    display_name = prom_sum_over_time
78)]
79pub fn sum_over_time(_: &TimestampMillisecondArray, values: &Float64Array) -> Option<f64> {
80    compute::sum(values)
81}
82
83/// The count of all values in the specified interval.
84#[range_fn(
85    name = CountOverTime,
86    ret = Float64Array,
87    display_name = prom_count_over_time
88)]
89pub fn count_over_time(_: &TimestampMillisecondArray, values: &Float64Array) -> Option<f64> {
90    if values.is_empty() {
91        None
92    } else {
93        Some(values.len() as f64)
94    }
95}
96
97/// The most recent point value in specified interval.
98#[range_fn(
99    name = LastOverTime,
100    ret = Float64Array,
101    display_name = prom_last_over_time
102)]
103pub fn last_over_time(_: &TimestampMillisecondArray, values: &Float64Array) -> Option<f64> {
104    values.values().last().copied()
105}
106
107/// absent_over_time returns an empty vector if the range vector passed to it has any
108/// elements (floats or native histograms) and a 1-element vector with the value 1 if
109/// the range vector passed to it has no elements.
110#[range_fn(
111    name = AbsentOverTime,
112    ret = Float64Array,
113    display_name = prom_absent_over_time
114)]
115pub fn absent_over_time(_: &TimestampMillisecondArray, values: &Float64Array) -> Option<f64> {
116    if values.is_empty() { Some(1.0) } else { None }
117}
118
119/// the value 1 for any series in the specified interval.
120#[range_fn(
121    name = PresentOverTime,
122    ret = Float64Array,
123    display_name = prom_present_over_time
124)]
125pub fn present_over_time(_: &TimestampMillisecondArray, values: &Float64Array) -> Option<f64> {
126    if values.is_empty() { None } else { Some(1.0) }
127}
128
129/// the population standard variance of the values in the specified interval.
130/// DataFusion's implementation:
131/// <https://github.com/apache/arrow-datafusion/blob/292eb954fc0bad3a1febc597233ba26cb60bda3e/datafusion/physical-expr/src/aggregate/variance.rs#L224-#L241>
132#[range_fn(
133    name = StdvarOverTime,
134    ret = Float64Array,
135    display_name = prom_stdvar_over_time
136)]
137pub fn stdvar_over_time(_: &TimestampMillisecondArray, values: &Float64Array) -> Option<f64> {
138    if values.is_empty() {
139        None
140    } else {
141        let mut count = 0;
142        let mut mean: f64 = 0.0;
143        let mut result: f64 = 0.0;
144        for value in values {
145            let value = value.unwrap();
146            let new_count = count + 1;
147            let delta1 = value - mean;
148            let new_mean = delta1 / new_count as f64 + mean;
149            let delta2 = value - new_mean;
150            let new_result = result + delta1 * delta2;
151
152            count += 1;
153            mean = new_mean;
154            result = new_result;
155        }
156        Some(result / count as f64)
157    }
158}
159
160/// the population standard deviation of the values in the specified interval.
161/// Prometheus's implementation: <https://github.com/prometheus/prometheus/blob/f55ab2217984770aa1eecd0f2d5f54580029b1c0/promql/functions.go#L556-L569>
162#[range_fn(
163    name = StddevOverTime,
164    ret = Float64Array,
165    display_name = prom_stddev_over_time
166)]
167pub fn stddev_over_time(_: &TimestampMillisecondArray, values: &Float64Array) -> Option<f64> {
168    if values.is_empty() {
169        None
170    } else {
171        let mut count = 0.0;
172        let mut mean = 0.0;
173        let mut comp_mean = 0.0;
174        let mut deviations_sum_sq = 0.0;
175        let mut comp_deviations_sum_sq = 0.0;
176        for v in values {
177            count += 1.0;
178            let current_value = v.unwrap();
179            let delta = current_value - (mean + comp_mean);
180            let (new_mean, new_comp_mean) = compensated_sum_inc(delta / count, mean, comp_mean);
181            mean = new_mean;
182            comp_mean = new_comp_mean;
183            let (new_deviations_sum_sq, new_comp_deviations_sum_sq) = compensated_sum_inc(
184                delta * (current_value - (mean + comp_mean)),
185                deviations_sum_sq,
186                comp_deviations_sum_sq,
187            );
188            deviations_sum_sq = new_deviations_sum_sq;
189            comp_deviations_sum_sq = new_comp_deviations_sum_sq;
190        }
191        Some(((deviations_sum_sq + comp_deviations_sum_sq) / count).sqrt())
192    }
193}
194
195#[cfg(test)]
196mod test {
197    use super::*;
198    use crate::functions::test_util::simple_range_udf_runner;
199
200    fn assert_over_time_value(actual: Option<f64>, expected: Option<f64>) {
201        match (actual, expected) {
202            (Some(actual), Some(expected)) if expected.is_nan() => assert!(actual.is_nan()),
203            (Some(actual), Some(expected)) => assert_eq!(actual, expected),
204            (None, None) => {}
205            (actual, expected) => panic!("expected {expected:?}, got {actual:?}"),
206        }
207    }
208
209    fn assert_min_max(
210        values: Vec<Option<f64>>,
211        expected_min: Option<f64>,
212        expected_max: Option<f64>,
213    ) {
214        let timestamps = TimestampMillisecondArray::from(vec![0; values.len()]);
215        let values = Float64Array::from(values);
216
217        assert_over_time_value(min_over_time(&timestamps, &values), expected_min);
218        assert_over_time_value(max_over_time(&timestamps, &values), expected_max);
219    }
220
221    #[test]
222    fn min_max_over_time_ignore_ordinary_nan_when_finite_values_exist() {
223        let ordinary_nan = f64::from_bits(0x7ff8_0000_0000_0000);
224
225        assert_min_max(
226            vec![Some(ordinary_nan), Some(3.0), Some(-2.0)],
227            Some(-2.0),
228            Some(3.0),
229        );
230        assert_min_max(
231            vec![Some(3.0), Some(ordinary_nan), Some(-2.0)],
232            Some(-2.0),
233            Some(3.0),
234        );
235        assert_min_max(
236            vec![Some(-2.0), Some(3.0), Some(ordinary_nan)],
237            Some(-2.0),
238            Some(3.0),
239        );
240        assert_min_max(
241            vec![Some(ordinary_nan), Some(ordinary_nan)],
242            Some(ordinary_nan),
243            Some(ordinary_nan),
244        );
245        assert_min_max(
246            vec![Some(3.0), Some(-2.0), Some(1.0)],
247            Some(-2.0),
248            Some(3.0),
249        );
250        assert_min_max(vec![], None, None);
251        assert_min_max(vec![None, None], None, None);
252    }
253
254    // build timestamp range and value range arrays for test
255    fn build_test_range_arrays() -> (RangeArray, RangeArray) {
256        let ts_array = Arc::new(TimestampMillisecondArray::from_iter(
257            [
258                1000i64, 3000, 5000, 7000, 9000, 11000, 13000, 15000, 17000, 200000, 500000,
259            ]
260            .into_iter()
261            .map(Some),
262        ));
263        let ranges = [
264            (0, 2),
265            (0, 5),
266            (1, 1), // only 1 element
267            (2, 0), // empty range
268            (2, 0), // empty range
269            (3, 3),
270            (4, 3),
271            (5, 3),
272            (8, 1), // only 1 element
273            (9, 0), // empty range
274        ];
275
276        let values_array = Arc::new(Float64Array::from_iter([
277            12.345678, 87.654321, 31.415927, 27.182818, 70.710678, 41.421356, 57.735027, 69.314718,
278            98.019802, 1.98019802, 61.803399,
279        ]));
280
281        let ts_range_array = RangeArray::from_ranges(ts_array, ranges).unwrap();
282        let value_range_array = RangeArray::from_ranges(values_array, ranges).unwrap();
283
284        (ts_range_array, value_range_array)
285    }
286
287    #[test]
288    fn calculate_avg_over_time() {
289        let (ts_array, value_array) = build_test_range_arrays();
290        simple_range_udf_runner(
291            AvgOverTime::scalar_udf(),
292            ts_array,
293            value_array,
294            vec![],
295            vec![
296                Some(49.9999995),
297                Some(45.8618844),
298                Some(87.654321),
299                None,
300                None,
301                Some(46.438284),
302                Some(56.62235366666667),
303                Some(56.15703366666667),
304                Some(98.019802),
305                None,
306            ],
307        );
308    }
309
310    #[test]
311    fn calculate_min_over_time() {
312        let (ts_array, value_array) = build_test_range_arrays();
313        simple_range_udf_runner(
314            MinOverTime::scalar_udf(),
315            ts_array,
316            value_array,
317            vec![],
318            vec![
319                Some(12.345678),
320                Some(12.345678),
321                Some(87.654321),
322                None,
323                None,
324                Some(27.182818),
325                Some(41.421356),
326                Some(41.421356),
327                Some(98.019802),
328                None,
329            ],
330        );
331    }
332
333    #[test]
334    fn calculate_max_over_time() {
335        let (ts_array, value_array) = build_test_range_arrays();
336        simple_range_udf_runner(
337            MaxOverTime::scalar_udf(),
338            ts_array,
339            value_array,
340            vec![],
341            vec![
342                Some(87.654321),
343                Some(87.654321),
344                Some(87.654321),
345                None,
346                None,
347                Some(70.710678),
348                Some(70.710678),
349                Some(69.314718),
350                Some(98.019802),
351                None,
352            ],
353        );
354    }
355
356    #[test]
357    fn calculate_sum_over_time() {
358        let (ts_array, value_array) = build_test_range_arrays();
359        simple_range_udf_runner(
360            SumOverTime::scalar_udf(),
361            ts_array,
362            value_array,
363            vec![],
364            vec![
365                Some(99.999999),
366                Some(229.309422),
367                Some(87.654321),
368                None,
369                None,
370                Some(139.314852),
371                Some(169.867061),
372                Some(168.471101),
373                Some(98.019802),
374                None,
375            ],
376        );
377    }
378
379    #[test]
380    fn calculate_count_over_time() {
381        let (ts_array, value_array) = build_test_range_arrays();
382        simple_range_udf_runner(
383            CountOverTime::scalar_udf(),
384            ts_array,
385            value_array,
386            vec![],
387            vec![
388                Some(2.0),
389                Some(5.0),
390                Some(1.0),
391                None,
392                None,
393                Some(3.0),
394                Some(3.0),
395                Some(3.0),
396                Some(1.0),
397                None,
398            ],
399        );
400    }
401
402    #[test]
403    fn calculate_last_over_time() {
404        let (ts_array, value_array) = build_test_range_arrays();
405        simple_range_udf_runner(
406            LastOverTime::scalar_udf(),
407            ts_array,
408            value_array,
409            vec![],
410            vec![
411                Some(87.654321),
412                Some(70.710678),
413                Some(87.654321),
414                None,
415                None,
416                Some(41.421356),
417                Some(57.735027),
418                Some(69.314718),
419                Some(98.019802),
420                None,
421            ],
422        );
423    }
424
425    #[test]
426    fn calculate_absent_over_time() {
427        let (ts_array, value_array) = build_test_range_arrays();
428        simple_range_udf_runner(
429            AbsentOverTime::scalar_udf(),
430            ts_array,
431            value_array,
432            vec![],
433            vec![
434                None,
435                None,
436                None,
437                Some(1.0),
438                Some(1.0),
439                None,
440                None,
441                None,
442                None,
443                Some(1.0),
444            ],
445        );
446    }
447
448    #[test]
449    fn calculate_present_over_time() {
450        let (ts_array, value_array) = build_test_range_arrays();
451        simple_range_udf_runner(
452            PresentOverTime::scalar_udf(),
453            ts_array,
454            value_array,
455            vec![],
456            vec![
457                Some(1.0),
458                Some(1.0),
459                Some(1.0),
460                None,
461                None,
462                Some(1.0),
463                Some(1.0),
464                Some(1.0),
465                Some(1.0),
466                None,
467            ],
468        );
469    }
470
471    #[test]
472    fn calculate_stdvar_over_time() {
473        let (ts_array, value_array) = build_test_range_arrays();
474        simple_range_udf_runner(
475            StdvarOverTime::scalar_udf(),
476            ts_array,
477            value_array,
478            vec![],
479            vec![
480                Some(1417.8479276253622),
481                Some(808.999919713209),
482                Some(0.0),
483                None,
484                None,
485                Some(328.3638826418587),
486                Some(143.5964181766362),
487                Some(130.91830542386285),
488                Some(0.0),
489                None,
490            ],
491        );
492
493        // add more assertions
494        let ts_array = Arc::new(TimestampMillisecondArray::from_iter(
495            [1000i64, 3000, 5000, 7000, 9000, 11000, 13000, 15000]
496                .into_iter()
497                .map(Some),
498        ));
499        let values_array = Arc::new(Float64Array::from_iter([
500            1.5990505637277868,
501            1.5990505637277868,
502            1.5990505637277868,
503            0.0,
504            8.0,
505            8.0,
506            2.0,
507            3.0,
508        ]));
509        let ranges = [(0, 3), (3, 5)];
510        simple_range_udf_runner(
511            StdvarOverTime::scalar_udf(),
512            RangeArray::from_ranges(ts_array, ranges).unwrap(),
513            RangeArray::from_ranges(values_array, ranges).unwrap(),
514            vec![],
515            vec![Some(0.0), Some(10.559999999999999)],
516        );
517    }
518
519    #[test]
520    fn calculate_std_dev_over_time() {
521        let (ts_array, value_array) = build_test_range_arrays();
522        simple_range_udf_runner(
523            StddevOverTime::scalar_udf(),
524            ts_array,
525            value_array,
526            vec![],
527            vec![
528                Some(37.6543215),
529                Some(28.442923895289123),
530                Some(0.0),
531                None,
532                None,
533                Some(18.12081352042062),
534                Some(11.983172291869804),
535                Some(11.441953741554055),
536                Some(0.0),
537                None,
538            ],
539        );
540
541        // add more assertions
542        let ts_array = Arc::new(TimestampMillisecondArray::from_iter(
543            [1000i64, 3000, 5000, 7000, 9000, 11000, 13000, 15000]
544                .into_iter()
545                .map(Some),
546        ));
547        let values_array = Arc::new(Float64Array::from_iter([
548            1.5990505637277868,
549            1.5990505637277868,
550            1.5990505637277868,
551            0.0,
552            8.0,
553            8.0,
554            2.0,
555            3.0,
556        ]));
557        let ranges = [(0, 3), (3, 5)];
558        simple_range_udf_runner(
559            StddevOverTime::scalar_udf(),
560            RangeArray::from_ranges(ts_array, ranges).unwrap(),
561            RangeArray::from_ranges(values_array, ranges).unwrap(),
562            vec![],
563            vec![Some(0.0), Some(3.249615361854384)],
564        );
565    }
566}