Removed tables - refactored processAggregationDocs func

This commit is contained in:
dsotirakis
2021-06-04 17:26:59 +03:00
parent 600300cf89
commit 23f80f42a5
2 changed files with 152 additions and 113 deletions
+79 -36
View File
@@ -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) {
+73 -77
View File
@@ -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{