1mod aggr_over_time;
16mod changes;
17mod deriv;
18mod double_exponential_smoothing;
19mod edge_count;
20mod extrapolate_rate;
21mod idelta;
22mod native_histogram;
23mod predict_linear;
24mod quantile;
25mod quantile_aggr;
26mod resets;
27mod round;
28#[cfg(test)]
29mod test_util;
30mod vector_matching;
31
32pub use aggr_over_time::{
33 AbsentOverTime, AvgOverTime, CountOverTime, LastOverTime, MaxOverTime, MinOverTime,
34 PresentOverTime, StddevOverTime, StdvarOverTime, SumOverTime,
35};
36pub use changes::Changes;
37use datafusion::arrow::array::{
38 ArrayRef, DictionaryArray, Float64Array, TimestampMillisecondArray,
39};
40use datafusion::error::DataFusionError;
41use datafusion::physical_plan::ColumnarValue;
42use datatypes::arrow::array::Array;
43use datatypes::arrow::datatypes::{DataType, Int64Type};
44pub use deriv::Deriv;
45pub use double_exponential_smoothing::DoubleExponentialSmoothing;
46pub use extrapolate_rate::{Delta, Increase, Rate};
47pub use idelta::IDelta;
48pub use native_histogram::{
49 MixedRange, NativeHistogramAbsentOverTime, NativeHistogramAdd, NativeHistogramAggAvg,
50 NativeHistogramAggSum, NativeHistogramAvg, NativeHistogramAvgOverTime, NativeHistogramChanges,
51 NativeHistogramCount, NativeHistogramCountOverTime, NativeHistogramDelta,
52 NativeHistogramDivScalar, NativeHistogramDrop, NativeHistogramEq, NativeHistogramFraction,
53 NativeHistogramIDelta, NativeHistogramIRate, NativeHistogramIncrease,
54 NativeHistogramLastOverTime, NativeHistogramMulScalar, NativeHistogramNeg,
55 NativeHistogramNotEq, NativeHistogramPresentOverTime, NativeHistogramQuantile,
56 NativeHistogramRate, NativeHistogramResets, NativeHistogramScalarMul, NativeHistogramStddev,
57 NativeHistogramStdvar, NativeHistogramSub, NativeHistogramSum, NativeHistogramSumOverTime,
58 NativeHistogramToString, PromqlFloatToString,
59};
60pub use predict_linear::PredictLinear;
61pub use quantile::QuantileOverTime;
62pub use quantile_aggr::{QUANTILE_NAME, quantile_udaf};
63pub use resets::Resets;
64pub use round::Round;
65pub use vector_matching::{MatchGroupViolation, UniqueMatchGroup};
66
67use crate::range_array::RangeArray;
68
69pub(crate) fn extract_array(columnar_value: &ColumnarValue) -> Result<ArrayRef, DataFusionError> {
73 match columnar_value {
74 ColumnarValue::Array(array) => Ok(array.clone()),
75 ColumnarValue::Scalar(scalar) => Ok(scalar.to_array_of_size(1)?),
76 }
77}
78
79pub(crate) fn extract_range_dict(
81 columnar_value: &ColumnarValue,
82 func_name: &str,
83 arg_name: &str,
84 expected_value_type: &DataType,
85) -> Result<DictionaryArray<Int64Type>, DataFusionError> {
86 let array = extract_array(columnar_value)?;
87 let dict = array
88 .as_any()
89 .downcast_ref::<DictionaryArray<Int64Type>>()
90 .ok_or_else(|| {
91 DataFusionError::Execution(format!(
92 "{func_name}: expect {arg_name} as DictionaryArray<Int64>, found {}",
93 array.data_type()
94 ))
95 })?
96 .clone();
97
98 if &dict.value_type() != expected_value_type {
99 return Err(DataFusionError::Execution(format!(
100 "{func_name}: expect {arg_name} values of type {expected_value_type}, found {}",
101 dict.value_type()
102 )));
103 }
104
105 RangeArray::try_new(dict.clone()).map_err(DataFusionError::from)?;
106 Ok(dict)
107}
108
109pub(crate) fn extract_range_array(
111 columnar_value: &ColumnarValue,
112) -> Result<RangeArray, DataFusionError> {
113 let array = extract_array(columnar_value)?;
114 let dict = array
115 .as_any()
116 .downcast_ref::<DictionaryArray<Int64Type>>()
117 .ok_or_else(|| {
118 DataFusionError::Execution(format!(
119 "expected DictionaryArray<Int64>, found {}",
120 array.data_type()
121 ))
122 })?
123 .clone();
124 RangeArray::try_new(dict).map_err(DataFusionError::from)
125}
126
127pub(crate) fn compensated_sum_inc(inc: f64, sum: f64, mut compensation: f64) -> (f64, f64) {
134 let new_sum = sum + inc;
135 if sum.abs() >= inc.abs() {
136 compensation += (sum - new_sum) + inc;
137 } else {
138 compensation += (inc - new_sum) + sum;
139 }
140 (new_sum, compensation)
141}
142
143pub(crate) fn linear_regression(
147 times: &TimestampMillisecondArray,
148 values: &Float64Array,
149 intercept_time: i64,
150) -> (Option<f64>, Option<f64>) {
151 linear_regression_slice(times.values(), values, 0, values.len(), intercept_time)
152}
153
154pub(crate) fn linear_regression_slice(
155 times: &[i64],
156 values: &Float64Array,
157 offset: usize,
158 len: usize,
159 intercept_time: i64,
160) -> (Option<f64>, Option<f64>) {
161 linear_regression_slices(times, offset, values, offset, len, intercept_time)
162}
163
164pub(crate) fn linear_regression_slices(
165 times: &[i64],
166 time_offset: usize,
167 values: &Float64Array,
168 value_offset: usize,
169 len: usize,
170 intercept_time: i64,
171) -> (Option<f64>, Option<f64>) {
172 let raw_values = values.values();
173 let has_nulls = values.null_count() > 0;
174 let mut count: f64 = 0.0;
175 let mut sum_x: f64 = 0.0;
176 let mut sum_y: f64 = 0.0;
177 let mut sum_xy: f64 = 0.0;
178 let mut sum_x2: f64 = 0.0;
179 let mut comp_x: f64 = 0.0;
180 let mut comp_y: f64 = 0.0;
181 let mut comp_xy: f64 = 0.0;
182 let mut comp_x2: f64 = 0.0;
183
184 let mut const_y = true;
185 let mut init_y = None;
186
187 for i in 0..len {
188 let time_idx = time_offset + i;
189 let value_idx = value_offset + i;
190 if has_nulls && values.is_null(value_idx) {
191 continue;
192 }
193 let value = raw_values[value_idx];
194 let time = times[time_idx] as f64;
195 let initial = init_y.get_or_insert(value);
196 if const_y && count > 0.0 && value != *initial {
197 const_y = false;
198 }
199 count += 1.0;
200 let x = (time - intercept_time as f64) / 1e3f64;
201 (sum_x, comp_x) = compensated_sum_inc(x, sum_x, comp_x);
202 (sum_y, comp_y) = compensated_sum_inc(value, sum_y, comp_y);
203 (sum_xy, comp_xy) = compensated_sum_inc(x * value, sum_xy, comp_xy);
204 (sum_x2, comp_x2) = compensated_sum_inc(x * x, sum_x2, comp_x2);
205 }
206
207 if count < 2.0 {
208 return (None, None);
209 }
210
211 if const_y {
212 let init_y = init_y.unwrap();
213 if !init_y.is_finite() {
214 return (None, None);
215 }
216 return (Some(0.0), Some(init_y));
217 }
218
219 sum_x += comp_x;
220 sum_y += comp_y;
221 sum_xy += comp_xy;
222 sum_x2 += comp_x2;
223
224 let cov_xy = sum_xy - sum_x * sum_y / count;
225 let var_x = sum_x2 - sum_x * sum_x / count;
226
227 let slope = cov_xy / var_x;
228 let intercept = sum_y / count - slope * sum_x / count;
229
230 (Some(slope), Some(intercept))
231}
232
233#[cfg(test)]
234mod test {
235 use std::sync::Arc;
236
237 use datafusion::physical_plan::ColumnarValue;
238 use datatypes::arrow::array::Int64Array;
239 use datatypes::arrow::datatypes::Int64Type;
240
241 use super::*;
242 use crate::range_array::RangeArray;
243
244 #[test]
245 fn calculate_linear_regression_none() {
246 let ts_array = TimestampMillisecondArray::from_iter(
247 [
248 0i64, 300, 600, 900, 1200, 1500, 1800, 2100, 2400, 2700, 3000,
249 ]
250 .into_iter()
251 .map(Some),
252 );
253 let values_array = Float64Array::from_iter([
254 1.0 / 0.0,
255 1.0 / 0.0,
256 1.0 / 0.0,
257 1.0 / 0.0,
258 1.0 / 0.0,
259 1.0 / 0.0,
260 1.0 / 0.0,
261 1.0 / 0.0,
262 1.0 / 0.0,
263 1.0 / 0.0,
264 ]);
265 let (slope, intercept) = linear_regression(&ts_array, &values_array, ts_array.value(0));
266 assert_eq!(slope, None);
267 assert_eq!(intercept, None);
268 }
269
270 #[test]
271 fn calculate_linear_regression_value_is_const() {
272 let ts_array = TimestampMillisecondArray::from_iter(
273 [
274 0i64, 300, 600, 900, 1200, 1500, 1800, 2100, 2400, 2700, 3000,
275 ]
276 .into_iter()
277 .map(Some),
278 );
279 let values_array =
280 Float64Array::from_iter([10.0, 10.0, 10.0, 10.0, 10.0, 10.0, 10.0, 10.0, 10.0, 10.0]);
281 let (slope, intercept) = linear_regression(&ts_array, &values_array, ts_array.value(0));
282 assert_eq!(slope, Some(0.0));
283 assert_eq!(intercept, Some(10.0));
284 }
285
286 #[test]
287 fn calculate_linear_regression() {
288 let ts_array = TimestampMillisecondArray::from_iter(
289 [
290 0i64, 300, 600, 900, 1200, 1500, 1800, 2100, 2400, 2700, 3000,
291 ]
292 .into_iter()
293 .map(Some),
294 );
295 let values_array = Float64Array::from_iter([
296 0.0, 10.0, 20.0, 30.0, 40.0, 0.0, 10.0, 20.0, 30.0, 40.0, 50.0,
297 ]);
298 let (slope, intercept) = linear_regression(&ts_array, &values_array, ts_array.value(0));
299 assert_eq!(slope, Some(10.606060606060607));
300 assert_eq!(intercept, Some(6.818181818181815));
301
302 let (slope, intercept) = linear_regression(&ts_array, &values_array, 3000);
303 assert_eq!(slope, Some(10.606060606060607));
304 assert_eq!(intercept, Some(38.63636363636364));
305 }
306
307 #[test]
308 fn calculate_linear_regression_value_have_none() {
309 let ts_array = TimestampMillisecondArray::from_iter(
310 [
311 0i64, 300, 600, 900, 1200, 1350, 1500, 1800, 2100, 2400, 2550, 2700, 3000,
312 ]
313 .into_iter()
314 .map(Some),
315 );
316 let values_array: Float64Array = [
317 Some(0.0),
318 Some(10.0),
319 Some(20.0),
320 Some(30.0),
321 Some(40.0),
322 None,
323 Some(0.0),
324 Some(10.0),
325 Some(20.0),
326 Some(30.0),
327 None,
328 Some(40.0),
329 Some(50.0),
330 ]
331 .into_iter()
332 .collect();
333 let (slope, intercept) = linear_regression(&ts_array, &values_array, ts_array.value(0));
334 assert_eq!(slope, Some(10.606060606060607));
335 assert_eq!(intercept, Some(6.818181818181815));
336 }
337
338 #[test]
339 fn calculate_linear_regression_value_all_none() {
340 let ts_array = TimestampMillisecondArray::from_iter([0i64, 300, 600].into_iter().map(Some));
341 let values_array: Float64Array = [None, None, None].into_iter().collect();
342 let (slope, intercept) = linear_regression(&ts_array, &values_array, ts_array.value(0));
343 assert_eq!(slope, None);
344 assert_eq!(intercept, None);
345 }
346
347 #[test]
349 fn test_kahan_sum() {
350 let inputs = vec![1.0, 10.0f64.powf(100.0), 1.0, -10.0f64.powf(100.0)];
351
352 let mut sum = 0.0;
353 let mut c = 0f64;
354
355 for v in inputs {
356 (sum, c) = compensated_sum_inc(v, sum, c);
357 }
358 assert_eq!(sum + c, 2.0)
359 }
360
361 #[test]
362 fn extract_range_array_rejects_external_dictionary_with_null_keys() {
363 let keys = Int64Array::from_iter([Some(0), None]);
364 let values = Arc::new(Float64Array::from_iter([1.0, 2.0]));
365 let dict = DictionaryArray::<Int64Type>::try_new(keys, values).unwrap();
366
367 let err = extract_range_array(&ColumnarValue::Array(Arc::new(dict))).unwrap_err();
368 assert!(err.to_string().contains("Empty range is not expected"));
369 }
370
371 #[test]
372 fn extract_range_array_accepts_internal_packed_ranges() {
373 let values = Arc::new(Float64Array::from_iter([1.0, 2.0, 3.0]));
374 let range_array = RangeArray::from_ranges(values, [(0, 2), (1, 2)]).unwrap();
375
376 let extracted =
377 extract_range_array(&ColumnarValue::Array(Arc::new(range_array.into_dict()))).unwrap();
378
379 assert_eq!(extracted.get_offset_length(0), Some((0, 2)));
380 assert_eq!(extracted.get_offset_length(1), Some((1, 2)));
381 }
382}