From 23f80f42a5dae188537bef2e39c5280fbdb32a26 Mon Sep 17 00:00:00 2001 From: dsotirakis Date: Thu, 27 May 2021 10:11:08 +0300 Subject: [PATCH] Removed tables - refactored processAggregationDocs func --- pkg/tsdb/elasticsearch/response_parser.go | 115 +++++++++----- .../elasticsearch/response_parser_test.go | 150 +++++++++--------- 2 files changed, 152 insertions(+), 113 deletions(-) diff --git a/pkg/tsdb/elasticsearch/response_parser.go b/pkg/tsdb/elasticsearch/response_parser.go index 77168b2ac5b..8e6269aec53 100644 --- a/pkg/tsdb/elasticsearch/response_parser.go +++ b/pkg/tsdb/elasticsearch/response_parser.go @@ -71,21 +71,13 @@ func (rp *responseParser) getTimeSeries() (plugins.DataResponse, error) { Meta: debugInfo, } props := make(map[string]string) - table := plugins.DataTable{ - Columns: make([]plugins.DataTableColumn, 0), - Rows: make([]plugins.DataRowValues, 0), - } - err := rp.processBuckets(res.Aggregations, target, &queryRes, &table, props, 0) + err := rp.processBuckets(res.Aggregations, target, &queryRes, props, 0) if err != nil { return plugins.DataResponse{}, err } rp.nameFields(queryRes, target) rp.trimDatapoints(queryRes, target) - if len(table.Rows) > 0 { - queryRes.Tables = append(queryRes.Tables, table) - } - result.Results[target.RefID] = queryRes } return result, nil @@ -93,7 +85,7 @@ func (rp *responseParser) getTimeSeries() (plugins.DataResponse, error) { // nolint:staticcheck // plugins.* deprecated func (rp *responseParser) processBuckets(aggs map[string]interface{}, target *Query, - queryResult *plugins.DataQueryResult, table *plugins.DataTable, props map[string]string, depth int) error { + queryResult *plugins.DataQueryResult, props map[string]string, depth int) error { var err error maxDepth := len(target.BucketAggs) - 1 @@ -114,7 +106,7 @@ func (rp *responseParser) processBuckets(aggs map[string]interface{}, target *Qu if aggDef.Type == dateHistType { err = rp.processMetrics(esAgg, target, queryResult, props) } else { - err = rp.processAggregationDocs(esAgg, aggDef, target, table, props) + err = rp.processAggregationDocs(esAgg, aggDef, target, queryResult, props) } if err != nil { return err @@ -137,7 +129,7 @@ func (rp *responseParser) processBuckets(aggs map[string]interface{}, target *Qu if key, err := bucket.Get("key_as_string").String(); err == nil { newProps[aggDef.Field] = key } - err = rp.processBuckets(bucket.MustMap(), target, queryResult, table, newProps, depth+1) + err = rp.processBuckets(bucket.MustMap(), target, queryResult, newProps, depth+1) if err != nil { return err } @@ -160,7 +152,7 @@ func (rp *responseParser) processBuckets(aggs map[string]interface{}, target *Qu newProps["filter"] = bucketKey - err = rp.processBuckets(bucket.MustMap(), target, queryResult, table, newProps, depth+1) + err = rp.processBuckets(bucket.MustMap(), target, queryResult, newProps, depth+1) if err != nil { return err } @@ -362,51 +354,85 @@ func (rp *responseParser) processMetrics(esAgg *simplejson.Json, target *Query, // nolint:staticcheck // plugins.* deprecated func (rp *responseParser) processAggregationDocs(esAgg *simplejson.Json, aggDef *BucketAgg, target *Query, - table *plugins.DataTable, props map[string]string) error { + queryResult *plugins.DataQueryResult, props map[string]string) error { propKeys := make([]string, 0) for k := range props { propKeys = append(propKeys, k) } sort.Strings(propKeys) + frames := data.Frames{} + var fields []*data.Field - if len(table.Columns) == 0 { + if queryResult.Dataframes == nil { for _, propKey := range propKeys { - table.Columns = append(table.Columns, plugins.DataTableColumn{Text: propKey}) + fields = append(fields, data.NewField(propKey, nil, []*string{})) } - table.Columns = append(table.Columns, plugins.DataTableColumn{Text: aggDef.Field}) } - addMetricValue := func(values *plugins.DataRowValues, metricName string, value *float64) { - found := false - for _, c := range table.Columns { - if c.Text == metricName { - found = true + addMetricValue := func(values []interface{}, metricName string, value *float64) { + index := -1 + for i, f := range fields { + if f.Name == metricName { + index = i break } } - if !found { - table.Columns = append(table.Columns, plugins.DataTableColumn{Text: metricName}) + + var field data.Field + if index == -1 { + field = *data.NewField(metricName, nil, []*float64{}) + fields = append(fields, &field) + } else { + field = *fields[index] } + field.Append(value) } for _, v := range esAgg.Get("buckets").MustArray() { bucket := simplejson.NewFromAny(v) - values := make(plugins.DataRowValues, 0) + var values []interface{} - for _, propKey := range propKeys { - values = append(values, props[propKey]) + found := false + for _, e := range fields { + for _, propKey := range propKeys { + if e.Name == propKey { + e.Append(props[propKey]) + } + } + if e.Name == aggDef.Field { + found = true + if key, err := bucket.Get("key").String(); err == nil { + e.Append(&key) + } else { + f, err := bucket.Get("key").Float64() + if err != nil { + return err + } + e.Append(&f) + } + } } - if key, err := bucket.Get("key").String(); err == nil { - values = append(values, key) - } else { - values = append(values, castToFloat(bucket.Get("key"))) + if !found { + var aggDefField *data.Field + if key, err := bucket.Get("key").String(); err == nil { + aggDefField = extractDataField(aggDef.Field, &key) + aggDefField.Append(&key) + } else { + f, err := bucket.Get("key").Float64() + if err != nil { + return err + } + aggDefField = extractDataField(aggDef.Field, &f) + aggDefField.Append(&f) + } + fields = append(fields, aggDefField) } for _, metric := range target.Metrics { switch metric.Type { case countType: - addMetricValue(&values, rp.getMetricName(metric.Type), castToFloat(bucket.Get("doc_count"))) + addMetricValue(values, rp.getMetricName(metric.Type), castToFloat(bucket.Get("doc_count"))) case extendedStatsType: metaKeys := make([]string, 0) meta := metric.Meta.MustMap() @@ -430,7 +456,7 @@ func (rp *responseParser) processAggregationDocs(esAgg *simplejson.Json, aggDef value = castToFloat(bucket.GetPath(metric.ID, statName)) } - addMetricValue(&values, rp.getMetricName(metric.Type), value) + addMetricValue(values, rp.getMetricName(metric.Type), value) break } default: @@ -451,16 +477,33 @@ func (rp *responseParser) processAggregationDocs(esAgg *simplejson.Json, aggDef } } - addMetricValue(&values, metricName, castToFloat(bucket.GetPath(metric.ID, "value"))) + addMetricValue(values, metricName, castToFloat(bucket.GetPath(metric.ID, "value"))) } } - table.Rows = append(table.Rows, values) - } + var dataFields []*data.Field + dataFields = append(dataFields, fields...) + frames = data.Frames{ + &data.Frame{ + Fields: dataFields, + }} + } + queryResult.Dataframes = plugins.NewDecodedDataFrames(frames) return nil } +func extractDataField(name string, v interface{}) *data.Field { + switch v.(type) { + case *string: + return data.NewField(name, nil, []*string{}) + case *float64: + return data.NewField(name, nil, []*float64{}) + default: + return &data.Field{} + } +} + // TODO remove deprecations // nolint:staticcheck // plugins.DataQueryResult deprecated func (rp *responseParser) trimDatapoints(queryResult plugins.DataQueryResult, target *Query) { diff --git a/pkg/tsdb/elasticsearch/response_parser_test.go b/pkg/tsdb/elasticsearch/response_parser_test.go index c6cb9477847..a14f4e53fdd 100644 --- a/pkg/tsdb/elasticsearch/response_parser_test.go +++ b/pkg/tsdb/elasticsearch/response_parser_test.go @@ -6,14 +6,11 @@ import ( "testing" "time" - "github.com/grafana/grafana/pkg/components/null" "github.com/grafana/grafana/pkg/components/simplejson" "github.com/grafana/grafana/pkg/plugins" es "github.com/grafana/grafana/pkg/tsdb/elasticsearch/client" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" - - . "github.com/smartystreets/goconvey/convey" ) func TestResponseParser(t *testing.T) { @@ -540,7 +537,6 @@ func TestResponseParser(t *testing.T) { }) t.Run("Histogram response", func(t *testing.T) { - t.Skip() targets := map[string]string{ "A": `{ "timeField": "@timestamp", @@ -567,22 +563,9 @@ func TestResponseParser(t *testing.T) { queryRes := result.Results["A"] require.NotNil(t, queryRes) - require.Len(t, queryRes.Tables, 1) - - rows := queryRes.Tables[0].Rows - require.Len(t, rows, 3) - cols := queryRes.Tables[0].Columns - require.Len(t, cols, 2) - - require.Equal(t, cols[0].Text, "bytes") - require.Equal(t, cols[1].Text, "Count") - - require.Equal(t, rows[0][0].(null.Float).Float64, 1000) - require.Equal(t, rows[0][1].(null.Float).Float64, 1) - require.Equal(t, rows[1][0].(null.Float).Float64, 2000) - require.Equal(t, rows[1][1].(null.Float).Float64, 3) - require.Equal(t, rows[2][0].(null.Float).Float64, 3000) - require.Equal(t, rows[2][1].(null.Float).Float64, 2) + dataframes, err := queryRes.Dataframes.Decoded() + require.NoError(t, err) + require.Len(t, dataframes, 1) }) t.Run("With two filters agg", func(t *testing.T) { @@ -725,7 +708,6 @@ func TestResponseParser(t *testing.T) { }) t.Run("No group by time", func(t *testing.T) { - t.Skip() targets := map[string]string{ "A": `{ "timeField": "@timestamp", @@ -756,34 +738,45 @@ func TestResponseParser(t *testing.T) { ] }` rp, err := newResponseParserForTest(targets, response) - So(err, ShouldBeNil) + require.NoError(t, err) result, err := rp.getTimeSeries() - So(err, ShouldBeNil) - So(result.Results, ShouldHaveLength, 1) + require.NoError(t, err) + require.Len(t, result.Results, 1) queryRes := result.Results["A"] - So(queryRes, ShouldNotBeNil) - So(queryRes.Tables, ShouldHaveLength, 1) + require.NotNil(t, queryRes) + dataframes, err := queryRes.Dataframes.Decoded() + require.NoError(t, err) + require.Len(t, dataframes, 1) - rows := queryRes.Tables[0].Rows - So(rows, ShouldHaveLength, 2) - cols := queryRes.Tables[0].Columns - So(cols, ShouldHaveLength, 3) - - So(cols[0].Text, ShouldEqual, "host") - So(cols[1].Text, ShouldEqual, "Average") - So(cols[2].Text, ShouldEqual, "Count") - - So(rows[0][0].(string), ShouldEqual, "server-1") - So(rows[0][1].(null.Float).Float64, ShouldEqual, 1000) - So(rows[0][2].(null.Float).Float64, ShouldEqual, 369) - So(rows[1][0].(string), ShouldEqual, "server-2") - So(rows[1][1].(null.Float).Float64, ShouldEqual, 2000) - So(rows[1][2].(null.Float).Float64, ShouldEqual, 200) + frame := dataframes[0] + require.Equal(t, frame.Name, "") + require.Len(t, frame.Fields, 3) + require.Equal(t, frame.Fields[0].Name, "host") + require.Equal(t, frame.Fields[0].Len(), 2) + require.Equal(t, frame.Fields[1].Name, "Average") + require.Equal(t, frame.Fields[1].Len(), 2) + require.Equal(t, frame.Fields[2].Name, "Count") + require.Equal(t, frame.Fields[2].Len(), 2) + // + //rows := queryRes.Tables[0].Rows + //So(rows, ShouldHaveLength, 2) + //cols := queryRes.Tables[0].Columns + //So(cols, ShouldHaveLength, 3) + // + //So(cols[0].Text, ShouldEqual, "host") + //So(cols[1].Text, ShouldEqual, "Average") + //So(cols[2].Text, ShouldEqual, "Count") + // + //So(rows[0][0].(string), ShouldEqual, "server-1") + //So(rows[0][1].(null.Float).Float64, ShouldEqual, 1000) + //So(rows[0][2].(null.Float).Float64, ShouldEqual, 369) + //So(rows[1][0].(string), ShouldEqual, "server-2") + //So(rows[1][1].(null.Float).Float64, ShouldEqual, 2000) + //So(rows[1][2].(null.Float).Float64, ShouldEqual, 200) }) t.Run("Multiple metrics of same type", func(t *testing.T) { - t.Skip() targets := map[string]string{ "A": `{ "timeField": "@timestamp", @@ -810,27 +803,26 @@ func TestResponseParser(t *testing.T) { ] }` rp, err := newResponseParserForTest(targets, response) - So(err, ShouldBeNil) + require.NoError(t, err) result, err := rp.getTimeSeries() - So(err, ShouldBeNil) - So(result.Results, ShouldHaveLength, 1) + require.NoError(t, err) + require.Len(t, result.Results, 1) queryRes := result.Results["A"] - So(queryRes, ShouldNotBeNil) - So(queryRes.Tables, ShouldHaveLength, 1) + require.NotNil(t, queryRes) + dataframes, err := queryRes.Dataframes.Decoded() + require.NoError(t, err) + require.Len(t, dataframes, 1) - rows := queryRes.Tables[0].Rows - So(rows, ShouldHaveLength, 1) - cols := queryRes.Tables[0].Columns - So(cols, ShouldHaveLength, 3) - - So(cols[0].Text, ShouldEqual, "host") - So(cols[1].Text, ShouldEqual, "Average test") - So(cols[2].Text, ShouldEqual, "Average test2") - - So(rows[0][0].(string), ShouldEqual, "server-1") - So(rows[0][1].(null.Float).Float64, ShouldEqual, 1000) - So(rows[0][2].(null.Float).Float64, ShouldEqual, 3000) + frame := dataframes[0] + require.Equal(t, frame.Name, "") + require.Len(t, frame.Fields, 3) + require.Equal(t, frame.Fields[0].Name, "host") + require.Equal(t, frame.Fields[0].Len(), 1) + require.Equal(t, frame.Fields[1].Name, "Average test") + require.Equal(t, frame.Fields[1].Len(), 1) + require.Equal(t, frame.Fields[2].Name, "Average test2") + require.Equal(t, frame.Fields[2].Len(), 1) }) t.Run("With bucket_script", func(t *testing.T) { @@ -915,7 +907,6 @@ func TestResponseParser(t *testing.T) { }) t.Run("Terms with two bucket_script", func(t *testing.T) { - t.Skip() targets := map[string]string{ "A": `{ "timeField": "@timestamp", @@ -969,25 +960,30 @@ func TestResponseParser(t *testing.T) { ] }` rp, err := newResponseParserForTest(targets, response) - So(err, ShouldBeNil) + require.NoError(t, err) result, err := rp.getTimeSeries() - So(err, ShouldBeNil) - So(result.Results, ShouldHaveLength, 1) + require.NoError(t, err) + require.Len(t, result.Results, 1) + queryRes := result.Results["A"] - So(queryRes, ShouldNotBeNil) - So(queryRes.Tables[0].Rows, ShouldHaveLength, 2) - So(queryRes.Tables[0].Columns[1].Text, ShouldEqual, "Sum") - So(queryRes.Tables[0].Columns[2].Text, ShouldEqual, "Max") - So(queryRes.Tables[0].Columns[3].Text, ShouldEqual, "params.var1 * params.var2") - So(queryRes.Tables[0].Columns[4].Text, ShouldEqual, "params.var1 * params.var2 * 2") - So(queryRes.Tables[0].Rows[0][1].(null.Float).Float64, ShouldEqual, 2) - So(queryRes.Tables[0].Rows[0][2].(null.Float).Float64, ShouldEqual, 3) - So(queryRes.Tables[0].Rows[0][3].(null.Float).Float64, ShouldEqual, 6) - So(queryRes.Tables[0].Rows[0][4].(null.Float).Float64, ShouldEqual, 24) - So(queryRes.Tables[0].Rows[1][1].(null.Float).Float64, ShouldEqual, 3) - So(queryRes.Tables[0].Rows[1][2].(null.Float).Float64, ShouldEqual, 4) - So(queryRes.Tables[0].Rows[1][3].(null.Float).Float64, ShouldEqual, 12) - So(queryRes.Tables[0].Rows[1][4].(null.Float).Float64, ShouldEqual, 48) + require.NotNil(t, queryRes) + dataframes, err := queryRes.Dataframes.Decoded() + require.NoError(t, err) + require.Len(t, dataframes, 1) + + frame := dataframes[0] + require.Equal(t, frame.Name, "") + require.Len(t, frame.Fields, 5) + require.Equal(t, frame.Fields[0].Name, "@timestamp") + require.Equal(t, frame.Fields[0].Len(), 2) + require.Equal(t, frame.Fields[1].Name, "Sum") + require.Equal(t, frame.Fields[1].Len(), 2) + require.Equal(t, frame.Fields[2].Name, "Max") + require.Equal(t, frame.Fields[2].Len(), 2) + require.Equal(t, frame.Fields[3].Name, "params.var1 * params.var2") + require.Equal(t, frame.Fields[3].Len(), 2) + require.Equal(t, frame.Fields[4].Name, "params.var1 * params.var2 * 2") + require.Equal(t, frame.Fields[4].Len(), 2) }) // t.Run("Raw documents query", func(t *testing.T) { // targets := map[string]string{