Skip to main content

promql/functions/
changes.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
15//! Implementation of [`changes`](https://prometheus.io/docs/prometheus/latest/querying/functions/#changes) in PromQL. Refer to the [original
16//! implementation](https://github.com/prometheus/prometheus/blob/main/promql/functions.go#L1023-L1040).
17
18use datafusion::logical_expr::ScalarUDF;
19
20use crate::functions::edge_count::{self, EdgeKind};
21
22/// used to count the number of value changes that occur within a specific time range
23#[derive(Debug)]
24pub struct Changes {}
25
26impl Changes {
27    pub const fn name() -> &'static str {
28        "prom_changes"
29    }
30
31    pub fn scalar_udf() -> ScalarUDF {
32        edge_count::scalar_udf(Self::name(), EdgeKind::Changes)
33    }
34}
35
36#[cfg(test)]
37mod test {
38    use std::sync::Arc;
39
40    use datafusion::arrow::array::{Float64Array, TimestampMillisecondArray};
41    use datafusion::arrow::buffer::NullBuffer;
42    use datatypes::arrow::array::Array;
43
44    use super::*;
45    use crate::functions::test_util::{
46        self, STALE_NAN, TinyPrng, assert_execution_error, build_test_range_arrays,
47        invoke_range_udf, simple_range_udf_runner,
48    };
49    use crate::range_array::RangeArray;
50
51    fn changes_oracle(values: &[f64]) -> Option<f64> {
52        let (first, rest) = values.split_first()?;
53        let mut changes = 0;
54        let mut previous = first;
55        for current in rest {
56            if current != previous && !(current.is_nan() && previous.is_nan()) {
57                changes += 1;
58            }
59            previous = current;
60        }
61        Some(changes as f64)
62    }
63
64    #[test]
65    fn calculate_changes() {
66        let timestamps = vec![
67            1000i64, 3000, 5000, 7000, 9000, 11000, 13000, 15000, 17000, 200000, 500000,
68        ];
69        let ranges = vec![
70            (0, 1),
71            (0, 4),
72            (0, 6),
73            (0, 10),
74            (0, 0), // empty range
75        ];
76
77        // assertion 1
78        let values_1 = vec![1.0, 2.0, 3.0, 0.0, 1.0, 0.0, 0.0, 1.0, 2.0, 0.0];
79        let (ts_array_1, value_array_1) =
80            build_test_range_arrays(timestamps.clone(), values_1, ranges.clone());
81        simple_range_udf_runner(
82            Changes::scalar_udf(),
83            ts_array_1,
84            value_array_1,
85            vec![],
86            vec![Some(0.0), Some(3.0), Some(5.0), Some(8.0), None],
87        );
88
89        // assertion 2
90        let values_2 = vec![1.0, 2.0, 3.0, 4.0, 5.0, 1.0, 2.0, 3.0, 4.0, 5.0];
91        let (ts_array_2, value_array_2) =
92            build_test_range_arrays(timestamps.clone(), values_2, ranges.clone());
93        simple_range_udf_runner(
94            Changes::scalar_udf(),
95            ts_array_2,
96            value_array_2,
97            vec![],
98            vec![Some(0.0), Some(3.0), Some(5.0), Some(9.0), None],
99        );
100
101        // assertion 3
102        let values_3 = vec![0.0, 0.0, 0.0, 0.0, 0.0, 1.0, 1.0, 1.0, 1.0, 1.0];
103        let (ts_array_3, value_array_3) = build_test_range_arrays(timestamps, values_3, ranges);
104        simple_range_udf_runner(
105            Changes::scalar_udf(),
106            ts_array_3,
107            value_array_3,
108            vec![],
109            vec![Some(0.0), Some(0.0), Some(1.0), Some(1.0), None],
110        );
111    }
112
113    #[test]
114    fn changes_range_array_oracle_edge_cases() {
115        let values = vec![
116            Some(-0.0),
117            Some(0.0),
118            Some(f64::INFINITY),
119            Some(f64::INFINITY),
120            Some(f64::NEG_INFINITY),
121            Some(f64::NAN),
122            Some(STALE_NAN),
123            Some(f64::NAN),
124            Some(3.0),
125            Some(1.0),
126            Some(2.0),
127            Some(3.0),
128            Some(1.0),
129            Some(3.0),
130            Some(0.0),
131        ];
132        let raw_values = values.iter().map(|value| value.unwrap()).collect();
133        let expected = test_util::run_oracle_ranges(
134            values,
135            raw_values,
136            vec![(5, 0), (2, 1), (7, 5), (0, 5), (11, 5), (16, 4)],
137            vec![(0, 0), (0, 1), (4, 5), (0, 5), (8, 5), (1, 4)],
138            changes_oracle,
139            Changes::scalar_udf(),
140        );
141        assert_eq!(
142            expected,
143            vec![None, Some(0.0), Some(2.0), Some(2.0), Some(4.0), Some(2.0)]
144        );
145
146        let expected = test_util::run_oracle_ranges(
147            vec![Some(2.0), None, Some(2.0)],
148            vec![2.0, 0.0, 2.0],
149            vec![(9, 3), (1, 1)],
150            vec![(0, 3), (1, 1)],
151            changes_oracle,
152            Changes::scalar_udf(),
153        );
154        assert_eq!(expected, vec![Some(2.0), Some(0.0)]);
155    }
156
157    #[test]
158    fn changes_range_array_seeded_differential() {
159        let mut prng = TinyPrng(0x2f6e_2b1d_834a_90c5);
160        let raw_values = (0..48)
161            .map(|_| match prng.next_index(12) {
162                0 => -0.0,
163                1 => 0.0,
164                2 => -2.0,
165                3 => -1.0,
166                4 => 1.0,
167                5 => 2.0,
168                6 => f64::INFINITY,
169                7 => f64::NEG_INFINITY,
170                8 | 9 => f64::NAN,
171                _ => STALE_NAN,
172            })
173            .collect::<Vec<_>>();
174        let values = raw_values.iter().copied().map(Some).collect();
175        let mut timestamp_ranges = Vec::new();
176        let mut value_ranges = Vec::new();
177        for _ in 0..32 {
178            let length = prng.next_index(13) as u32;
179            timestamp_ranges.push((prng.next_index(65 - length as usize) as u32, length));
180            value_ranges.push((prng.next_index(49 - length as usize) as u32, length));
181        }
182
183        test_util::run_oracle_ranges(
184            values,
185            raw_values,
186            timestamp_ranges,
187            value_ranges,
188            changes_oracle,
189            Changes::scalar_udf(),
190        );
191    }
192
193    #[test]
194    fn changes_range_array_mismatch_errors() {
195        let timestamps = Arc::new(TimestampMillisecondArray::from_iter([
196            Some(0),
197            Some(1),
198            Some(2),
199        ]));
200        let values = Arc::new(Float64Array::from_iter([1.0, 2.0, 3.0]));
201        let error = invoke_range_udf(
202            Changes::scalar_udf(),
203            RangeArray::from_ranges(timestamps.clone(), [(0, 1), (1, 1)]).unwrap(),
204            RangeArray::from_ranges(values.clone(), [(0, 1)]).unwrap(),
205        )
206        .unwrap_err();
207        assert_execution_error(
208            error,
209            "RangeArray have different lengths in PromQL function prom_changes: array1=2, array2=1",
210        );
211
212        let error = invoke_range_udf(
213            Changes::scalar_udf(),
214            RangeArray::from_ranges(timestamps, [(0, 1), (1, 2)]).unwrap(),
215            RangeArray::from_ranges(values, [(0, 1), (1, 1)]).unwrap(),
216        )
217        .unwrap_err();
218        assert_execution_error(
219            error,
220            "RangeArray's element 1 have different lengths in PromQL function prom_changes: array1=2, array2=1",
221        );
222    }
223
224    #[test]
225    fn changes_range_array_boundaries_and_raw_nulls() {
226        // The 0.0 -> 1.0 edge enters the range; both identical windows contain no changes.
227        let (timestamps, values) = build_test_range_arrays(
228            vec![0, 1, 2, 3],
229            vec![0.0, 1.0, 1.0, 1.0],
230            vec![(1, 3), (1, 3)],
231        );
232        simple_range_udf_runner(
233            Changes::scalar_udf(),
234            timestamps,
235            values,
236            vec![],
237            vec![Some(0.0), Some(0.0)],
238        );
239
240        let (timestamps, values) = build_test_range_arrays(
241            vec![0, 1, 2, 3],
242            vec![1.0, 2.0, 3.0, 4.0],
243            vec![(0, 0), (2, 0), (4, 0)],
244        );
245        simple_range_udf_runner(
246            Changes::scalar_udf(),
247            timestamps,
248            values,
249            vec![],
250            vec![None, None, None],
251        );
252
253        let (timestamps, values) = build_test_range_arrays(
254            vec![0, 1, 2, 3],
255            vec![1.0, 2.0, 3.0, 4.0],
256            vec![(0, 1), (2, 1), (3, 1)],
257        );
258        simple_range_udf_runner(
259            Changes::scalar_udf(),
260            timestamps,
261            values,
262            vec![],
263            vec![Some(0.0), Some(0.0), Some(0.0)],
264        );
265
266        let values = Arc::new(Float64Array::new(
267            vec![10.0, 7.0, 10.0].into(),
268            Some(NullBuffer::from(vec![true, false, true])),
269        ));
270        assert!(!values.is_valid(1));
271        assert_eq!(values.value(1), 7.0);
272        let timestamps = Arc::new(TimestampMillisecondArray::from_iter([
273            Some(0),
274            Some(1),
275            Some(2),
276        ]));
277        simple_range_udf_runner(
278            Changes::scalar_udf(),
279            RangeArray::from_ranges(timestamps, [(0, 3)]).unwrap(),
280            RangeArray::from_ranges(values, [(0, 3)]).unwrap(),
281            vec![],
282            vec![Some(2.0)],
283        );
284    }
285}