1use 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#[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#[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#[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#[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#[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#[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#[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#[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#[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#[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(×tamps, &values), expected_min);
218 assert_over_time_value(max_over_time(×tamps, &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 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), (2, 0), (2, 0), (3, 3),
270 (4, 3),
271 (5, 3),
272 (8, 1), (9, 0), ];
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 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 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}