From 512faa7c9c2cbbb35afaadc66f0071f025855b77 Mon Sep 17 00:00:00 2001 From: Kyle Brandt Date: Thu, 1 Apr 2021 14:33:35 -0400 Subject: [PATCH] SSE: Support time series frames with non-nullable float64 values (#32608) * SSE: fix reduce to handle non-null * add type for data.Field that is float64 or *float64 * resample fix for non-null input case, add couple non-null tests --- pkg/expr/classic/classic.go | 1 + pkg/expr/classic/reduce.go | 139 ++++++++++++++---------------- pkg/expr/mathexp/reduce.go | 70 ++++++++------- pkg/expr/mathexp/reduce_test.go | 13 +++ pkg/expr/mathexp/resample.go | 11 +-- pkg/expr/mathexp/resample_test.go | 28 ++++++ pkg/expr/mathexp/types.go | 20 +++++ 7 files changed, 167 insertions(+), 115 deletions(-) diff --git a/pkg/expr/classic/classic.go b/pkg/expr/classic/classic.go index f4e7b503055..45f9436acc0 100644 --- a/pkg/expr/classic/classic.go +++ b/pkg/expr/classic/classic.go @@ -86,6 +86,7 @@ func (ccc *ConditionsCmd) Execute(ctx context.Context, vars mathexp.Vars) (mathe } reducedNum := c.Reducer.Reduce(series) + // TODO handle error / no data signals thisCondNoDataFound := reducedNum.GetFloat64Value() == nil diff --git a/pkg/expr/classic/reduce.go b/pkg/expr/classic/reduce.go index be12bd98b7b..39991a2e262 100644 --- a/pkg/expr/classic/reduce.go +++ b/pkg/expr/classic/reduce.go @@ -4,7 +4,6 @@ import ( "math" "sort" - "github.com/grafana/grafana-plugin-sdk-go/data" "github.com/grafana/grafana/pkg/expr/mathexp" ) @@ -35,44 +34,42 @@ func (cr classicReducer) Reduce(series mathexp.Series) mathexp.Number { allNull := true vF := series.Frame.Fields[series.ValueIdx] + ff := mathexp.Float64Field(*vF) switch cr { case "avg": validPointsCount := 0 - for i := 0; i < vF.Len(); i++ { - if f, ok := vF.At(i).(*float64); ok { - if nilOrNaN(f) { - continue - } - value += *f - validPointsCount++ - allNull = false + for i := 0; i < ff.Len(); i++ { + f := ff.GetValue(i) + if nilOrNaN(f) { + continue } + value += *f + validPointsCount++ + allNull = false } if validPointsCount > 0 { value /= float64(validPointsCount) } case "sum": - for i := 0; i < vF.Len(); i++ { - if f, ok := vF.At(i).(*float64); ok { - if nilOrNaN(f) { - continue - } - value += *f - allNull = false + for i := 0; i < ff.Len(); i++ { + f := ff.GetValue(i) + if nilOrNaN(f) { + continue } + value += *f + allNull = false } case "min": value = math.MaxFloat64 - for i := 0; i < vF.Len(); i++ { - if f, ok := vF.At(i).(*float64); ok { - if nilOrNaN(f) { - continue - } - allNull = false - if value > *f { - value = *f - } + for i := 0; i < ff.Len(); i++ { + f := ff.GetValue(i) + if nilOrNaN(f) { + continue + } + allNull = false + if value > *f { + value = *f } } if allNull { @@ -80,43 +77,40 @@ func (cr classicReducer) Reduce(series mathexp.Series) mathexp.Number { } case "max": value = -math.MaxFloat64 - for i := 0; i < vF.Len(); i++ { - if f, ok := vF.At(i).(*float64); ok { - if nilOrNaN(f) { - continue - } - allNull = false - if value < *f { - value = *f - } + for i := 0; i < ff.Len(); i++ { + f := ff.GetValue(i) + if nilOrNaN(f) { + continue + } + allNull = false + if value < *f { + value = *f } } if allNull { value = 0 } case "count": - value = float64(vF.Len()) + value = float64(ff.Len()) allNull = false case "last": - for i := vF.Len() - 1; i >= 0; i-- { - if f, ok := vF.At(i).(*float64); ok { - if !nilOrNaN(f) { - value = *f - allNull = false - break - } + for i := ff.Len() - 1; i >= 0; i-- { + f := ff.GetValue(i) + if !nilOrNaN(f) { + value = *f + allNull = false + break } } case "median": var values []float64 - for i := 0; i < vF.Len(); i++ { - if f, ok := vF.At(i).(*float64); ok { - if nilOrNaN(f) { - continue - } - allNull = false - values = append(values, *f) + for i := 0; i < ff.Len(); i++ { + f := ff.GetValue(i) + if nilOrNaN(f) { + continue } + allNull = false + values = append(values, *f) } if len(values) >= 1 { sort.Float64s(values) @@ -128,21 +122,20 @@ func (cr classicReducer) Reduce(series mathexp.Series) mathexp.Number { } } case "diff": - allNull, value = calculateDiff(vF, allNull, value, diff) + allNull, value = calculateDiff(ff, allNull, value, diff) case "diff_abs": - allNull, value = calculateDiff(vF, allNull, value, diffAbs) + allNull, value = calculateDiff(ff, allNull, value, diffAbs) case "percent_diff": - allNull, value = calculateDiff(vF, allNull, value, percentDiff) + allNull, value = calculateDiff(ff, allNull, value, percentDiff) case "percent_diff_abs": - allNull, value = calculateDiff(vF, allNull, value, percentDiffAbs) + allNull, value = calculateDiff(ff, allNull, value, percentDiffAbs) case "count_non_null": - for i := 0; i < vF.Len(); i++ { - if f, ok := vF.At(i).(*float64); ok { - if nilOrNaN(f) { - continue - } - value++ + for i := 0; i < ff.Len(); i++ { + f := ff.GetValue(i) + if nilOrNaN(f) { + continue } + value++ } if value > 0 { @@ -158,30 +151,28 @@ func (cr classicReducer) Reduce(series mathexp.Series) mathexp.Number { return num } -func calculateDiff(vF *data.Field, allNull bool, value float64, fn func(float64, float64) float64) (bool, float64) { +func calculateDiff(ff mathexp.Float64Field, allNull bool, value float64, fn func(float64, float64) float64) (bool, float64) { var ( first float64 i int ) // get the newest point - for i = vF.Len() - 1; i >= 0; i-- { - if f, ok := vF.At(i).(*float64); ok { - if !nilOrNaN(f) { - first = *f - allNull = false - break - } + for i = ff.Len() - 1; i >= 0; i-- { + f := ff.GetValue(i) + if !nilOrNaN(f) { + first = *f + allNull = false + break } } if i >= 1 { // get the oldest point - for i := 0; i < vF.Len(); i++ { - if f, ok := vF.At(i).(*float64); ok { - if !nilOrNaN(f) { - value = fn(first, *f) - allNull = false - break - } + for i := 0; i < ff.Len(); i++ { + f := ff.GetValue(i) + if !nilOrNaN(f) { + value = fn(first, *f) + allNull = false + break } } } diff --git a/pkg/expr/mathexp/reduce.go b/pkg/expr/mathexp/reduce.go index 34a5a862058..6591ad23cb6 100644 --- a/pkg/expr/mathexp/reduce.go +++ b/pkg/expr/mathexp/reduce.go @@ -7,67 +7,64 @@ import ( "github.com/grafana/grafana-plugin-sdk-go/data" ) -func Sum(v *data.Field) *float64 { +func Sum(fv *Float64Field) *float64 { var sum float64 - for i := 0; i < v.Len(); i++ { - if f, ok := v.At(i).(*float64); ok { - if f == nil || math.IsNaN(*f) { - nan := math.NaN() - return &nan - } - sum += *f + for i := 0; i < fv.Len(); i++ { + f := fv.GetValue(i) + if f == nil || math.IsNaN(*f) { + nan := math.NaN() + return &nan } + sum += *f } return &sum } -func Avg(v *data.Field) *float64 { - sum := Sum(v) - f := *sum / float64(v.Len()) +func Avg(fv *Float64Field) *float64 { + sum := Sum(fv) + f := *sum / float64(fv.Len()) return &f } -func Min(fv *data.Field) *float64 { +func Min(fv *Float64Field) *float64 { var f float64 if fv.Len() == 0 { nan := math.NaN() return &nan } for i := 0; i < fv.Len(); i++ { - if v, ok := fv.At(i).(*float64); ok { - if v == nil || math.IsNaN(*v) { - nan := math.NaN() - return &nan - } - if i == 0 || *v < f { - f = *v - } + v := fv.GetValue(i) + if v == nil || math.IsNaN(*v) { + nan := math.NaN() + return &nan + } + if i == 0 || *v < f { + f = *v } } return &f } -func Max(fv *data.Field) *float64 { +func Max(fv *Float64Field) *float64 { var f float64 if fv.Len() == 0 { nan := math.NaN() return &nan } for i := 0; i < fv.Len(); i++ { - if v, ok := fv.At(i).(*float64); ok { - if v == nil || math.IsNaN(*v) { - nan := math.NaN() - return &nan - } - if i == 0 || *v > f { - f = *v - } + v := fv.GetValue(i) + if v == nil || math.IsNaN(*v) { + nan := math.NaN() + return &nan + } + if i == 0 || *v > f { + f = *v } } return &f } -func Count(fv *data.Field) *float64 { +func Count(fv *Float64Field) *float64 { f := float64(fv.Len()) return &f } @@ -80,18 +77,19 @@ func (s Series) Reduce(refID, rFunc string) (Number, error) { } number := NewNumber(refID, l) var f *float64 - fVec := s.Frame.Fields[1] + fVec := s.Frame.Fields[s.ValueIdx] + floatField := Float64Field(*fVec) switch rFunc { case "sum": - f = Sum(fVec) + f = Sum(&floatField) case "mean": - f = Avg(fVec) + f = Avg(&floatField) case "min": - f = Min(fVec) + f = Min(&floatField) case "max": - f = Max(fVec) + f = Max(&floatField) case "count": - f = Count(fVec) + f = Count(&floatField) default: return number, fmt.Errorf("reduction %v not implemented", rFunc) } diff --git a/pkg/expr/mathexp/reduce_test.go b/pkg/expr/mathexp/reduce_test.go index 18d2f8e96fd..4a0efdce7ed 100644 --- a/pkg/expr/mathexp/reduce_test.go +++ b/pkg/expr/mathexp/reduce_test.go @@ -178,6 +178,19 @@ func TestSeriesReduce(t *testing.T) { }, }, }, + { + name: "mean series (non-null-value)", + red: "mean", + varToReduce: "A", + vars: aSeriesNoNull, + errIs: require.NoError, + resultsIs: require.Equal, + results: Results{ + []Value{ + makeNumber("", nil, float64Pointer(1.5)), + }, + }, + }, { name: "count empty series", red: "count", diff --git a/pkg/expr/mathexp/resample.go b/pkg/expr/mathexp/resample.go index b7773bfa066..8f5cade603a 100644 --- a/pkg/expr/mathexp/resample.go +++ b/pkg/expr/mathexp/resample.go @@ -14,7 +14,7 @@ func (s Series) Resample(refID string, interval time.Duration, downsampler strin if newSeriesLength <= 0 { return s, fmt.Errorf("the series cannot be sampled further; the time range is shorter than the interval") } - resampled := NewSeries(refID, s.GetLabels(), s.TimeIdx, s.TimeIsNullable, s.ValueIdx, s.ValueIsNullable, newSeriesLength+1) + resampled := NewSeries(refID, s.GetLabels(), s.TimeIdx, true, s.ValueIdx, true, newSeriesLength+1) bookmark := 0 var lastSeen *float64 idx := 0 @@ -57,16 +57,17 @@ func (s Series) Resample(refID string, interval time.Duration, downsampler strin } } else { // downsampling fVec := data.NewField("", s.GetLabels(), vals) + ff := Float64Field(*fVec) var tmp *float64 switch downsampler { case "sum": - tmp = Sum(fVec) + tmp = Sum(&ff) case "mean": - tmp = Avg(fVec) + tmp = Avg(&ff) case "min": - tmp = Min(fVec) + tmp = Min(&ff) case "max": - tmp = Max(fVec) + tmp = Max(&ff) default: return s, fmt.Errorf("downsampling %v not implemented", downsampler) } diff --git a/pkg/expr/mathexp/resample_test.go b/pkg/expr/mathexp/resample_test.go index 772f3e75c97..204125ca82f 100644 --- a/pkg/expr/mathexp/resample_test.go +++ b/pkg/expr/mathexp/resample_test.go @@ -77,6 +77,34 @@ func TestResampleSeries(t *testing.T) { unixTimePointer(15, 0), float64Pointer(2), }), }, + { + name: "resample series: downsampling (mean / pad) (no-nullable)", + interval: time.Second * 5, + downsampler: "mean", + upsampler: "pad", + timeRange: backend.TimeRange{ + From: time.Unix(0, 0), + To: time.Unix(16, 0), + }, + seriesToResample: makeNoNullSeries("", nil, noNullTP{ + time.Unix(2, 0), 2, + }, noNullTP{ + time.Unix(4, 0), 3, + }, noNullTP{ + time.Unix(7, 0), 1, + }, noNullTP{ + time.Unix(9, 0), 2, + }), + series: makeSeriesNullableTime("", nil, nullTimeTP{ + unixTimePointer(0, 0), nil, + }, nullTimeTP{ + unixTimePointer(5, 0), float64Pointer(2.5), + }, nullTimeTP{ + unixTimePointer(10, 0), float64Pointer(1.5), + }, nullTimeTP{ + unixTimePointer(15, 0), float64Pointer(2), + }), + }, { name: "resample series: downsampling (max / fillna)", interval: time.Second * 5, diff --git a/pkg/expr/mathexp/types.go b/pkg/expr/mathexp/types.go index d60e5df06e7..c0fbd56695c 100644 --- a/pkg/expr/mathexp/types.go +++ b/pkg/expr/mathexp/types.go @@ -123,3 +123,23 @@ func (n Number) GetMeta() interface{} { func (n Number) SetMeta(v interface{}) { n.Frame.SetMeta(&data.FrameMeta{Custom: v}) } + +// FloatField is a *float64 or a float64 data.Field with methods to always +// get a *float64. +type Float64Field data.Field + +// GetValue returns the value at idx as *float64. +func (ff *Float64Field) GetValue(idx int) *float64 { + field := data.Field(*ff) + if field.Type() == data.FieldTypeNullableFloat64 { + return field.At(idx).(*float64) + } + f := field.At(idx).(float64) + return &f +} + +// Len returns the the length of the field. +func (ff *Float64Field) Len() int { + df := data.Field(*ff) + return df.Len() +}