From c5f3ce13717853c07dbb9fb2ebd27dd976bcc97b Mon Sep 17 00:00:00 2001 From: ismail simsek Date: Wed, 22 Nov 2023 19:27:42 +0100 Subject: [PATCH] InfluxDB: Parse data for table view to have parity with frontend parser (#78365) * Use TimeSeriesWide format for table response * fix group by query result parsing * handle labels * provide a test where result has no tags * parsing results without time column * clean the code * remove the comment line * more cleaning * lint --- pkg/tsdb/influxdb/influxql/response_parser.go | 128 +++- .../influxdb/influxql/response_parser_test.go | 671 +++++++++++++++++- 2 files changed, 793 insertions(+), 6 deletions(-) diff --git a/pkg/tsdb/influxdb/influxql/response_parser.go b/pkg/tsdb/influxdb/influxql/response_parser.go index d4116b147d8..1c09d67e028 100644 --- a/pkg/tsdb/influxdb/influxql/response_parser.go +++ b/pkg/tsdb/influxdb/influxql/response_parser.go @@ -49,9 +49,13 @@ func parse(buf io.Reader, statusCode int, query *models.Query) *backend.DataResp result := response.Results[0] if result.Error != "" { return &backend.DataResponse{Error: fmt.Errorf(result.Error)} - } else { - return &backend.DataResponse{Frames: transformRows(result.Series, *query)} } + + if query.ResultFormat == "table" { + return &backend.DataResponse{Frames: transformRowsForTable(result.Series, *query)} + } + + return &backend.DataResponse{Frames: transformRowsForTimeSeries(result.Series, *query)} } func parseJSON(buf io.Reader) (models.Response, error) { @@ -65,7 +69,125 @@ func parseJSON(buf io.Reader) (models.Response, error) { return response, err } -func transformRows(rows []models.Row, query models.Query) data.Frames { +func transformRowsForTable(rows []models.Row, query models.Query) data.Frames { + if len(rows) == 0 { + return make([]*data.Frame, 0) + } + + frames := make([]*data.Frame, 0, 1) + + newFrame := data.NewFrame(rows[0].Name) + newFrame.Meta = &data.FrameMeta{ + ExecutedQueryString: query.RawQuery, + PreferredVisualization: getVisType(query.ResultFormat), + } + + conLen := len(rows[0].Columns) + if rows[0].Columns[0] == "time" { + newFrame.Fields = append(newFrame.Fields, newTimeField(rows)) + } else { + newFrame.Fields = append(newFrame.Fields, newValueFields(rows, nil, 0, 1)...) + } + + newFrame.Fields = append(newFrame.Fields, newTagField(rows, nil)...) + newFrame.Fields = append(newFrame.Fields, newValueFields(rows, nil, 1, conLen)...) + + frames = append(frames, newFrame) + return frames +} + +func newTimeField(rows []models.Row) *data.Field { + var timeArray []time.Time + for _, row := range rows { + for _, valuePair := range row.Values { + timestamp, timestampErr := parseTimestamp(valuePair[0]) + // we only add this row if the timestamp is valid + if timestampErr != nil { + continue + } + + timeArray = append(timeArray, timestamp) + } + } + + timeField := data.NewField("Time", nil, timeArray) + return timeField +} + +func newTagField(rows []models.Row, labels data.Labels) []*data.Field { + fields := make([]*data.Field, 0, len(rows[0].Tags)) + + for key := range rows[0].Tags { + tagField := data.NewField(key, labels, []*string{}) + for _, row := range rows { + for range row.Values { + value := row.Tags[key] + tagField.Append(&value) + } + } + tagField.SetConfig(&data.FieldConfig{DisplayNameFromDS: key}) + fields = append(fields, tagField) + } + + return fields +} + +func newValueFields(rows []models.Row, labels data.Labels, colIdxStart, colIdxEnd int) []*data.Field { + fields := make([]*data.Field, 0) + + for colIdx := colIdxStart; colIdx < colIdxEnd; colIdx++ { + var valueField *data.Field + var floatArray []*float64 + var stringArray []*string + var boolArray []*bool + + for _, row := range rows { + valType := typeof(row.Values, colIdx) + + for _, valuePair := range row.Values { + switch valType { + case "string": + value, ok := valuePair[colIdx].(string) + if ok { + stringArray = append(stringArray, &value) + } else { + stringArray = append(stringArray, nil) + } + case "json.Number": + value := parseNumber(valuePair[colIdx]) + floatArray = append(floatArray, value) + case "bool": + value, ok := valuePair[colIdx].(bool) + if ok { + boolArray = append(boolArray, &value) + } else { + boolArray = append(boolArray, nil) + } + case "null": + floatArray = append(floatArray, nil) + } + } + + switch valType { + case "string": + valueField = data.NewField(row.Columns[colIdx], labels, stringArray) + case "json.Number": + valueField = data.NewField(row.Columns[colIdx], labels, floatArray) + case "bool": + valueField = data.NewField(row.Columns[colIdx], labels, boolArray) + case "null": + valueField = data.NewField(row.Columns[colIdx], labels, floatArray) + } + + valueField.SetConfig(&data.FieldConfig{DisplayNameFromDS: row.Columns[colIdx]}) + } + fields = append(fields, valueField) + } + + return fields +} + +func transformRowsForTimeSeries(rows []models.Row, query models.Query) data.Frames { // pre-allocate frames - this can save many allocations cols := 0 for _, row := range rows { diff --git a/pkg/tsdb/influxdb/influxql/response_parser_test.go b/pkg/tsdb/influxdb/influxql/response_parser_test.go index a71071bcf26..851a00b99a2 100644 --- a/pkg/tsdb/influxdb/influxql/response_parser_test.go +++ b/pkg/tsdb/influxdb/influxql/response_parser_test.go @@ -810,6 +810,97 @@ func TestResponseParser_Parse_RetentionPolicy(t *testing.T) { }) } +func TestResponseParser_table_format(t *testing.T) { + t.Run("test table result format parsing", func(t *testing.T) { + resp := ResponseParse(prepare(tableResultFormatInfluxResponse1), 200, &models.Query{RefID: "A", RawQuery: `a nice query`, ResultFormat: "table"}) + assert.Equal(t, 1, len(resp.Frames)) + assert.Equal(t, "a nice query", resp.Frames[0].Meta.ExecutedQueryString) + assert.Equal(t, 3, len(resp.Frames[0].Fields)) + for i := range resp.Frames[0].Fields { + assert.Equal(t, resp.Frames[0].Fields[0].Len(), resp.Frames[0].Fields[i].Len()) + } + assert.Equal(t, "Time", resp.Frames[0].Fields[0].Name) + assert.Equal(t, "usage_idle", resp.Frames[0].Fields[1].Name) + assert.Equal(t, toPtr(99.09456740445926), resp.Frames[0].Fields[1].At(2)) + assert.Equal(t, "usage_iowait", resp.Frames[0].Fields[2].Name) + }) + + t.Run("test table result format parsing with grouping", func(t *testing.T) { + resp := ResponseParse(prepare(tableResultFormatInfluxResponse2), 200, &models.Query{RefID: "A", RawQuery: `a nice query`, ResultFormat: "table"}) + assert.Equal(t, 1, len(resp.Frames)) + assert.Equal(t, "a nice query", resp.Frames[0].Meta.ExecutedQueryString) + assert.Equal(t, 7, len(resp.Frames[0].Fields)) + for i := range resp.Frames[0].Fields { + assert.Equal(t, resp.Frames[0].Fields[0].Len(), resp.Frames[0].Fields[i].Len()) + } + assert.Equal(t, "Time", resp.Frames[0].Fields[0].Name) + assert.Equal(t, "cpu", resp.Frames[0].Fields[1].Name) + assert.Equal(t, toPtr("cpu-total"), resp.Frames[0].Fields[1].At(0)) + assert.Equal(t, toPtr("cpu0"), resp.Frames[0].Fields[1].At(1)) + assert.Equal(t, toPtr("cpu9"), resp.Frames[0].Fields[1].At(10)) + assert.Equal(t, "mean", resp.Frames[0].Fields[2].Name) + assert.Equal(t, "min", resp.Frames[0].Fields[3].Name) + assert.Equal(t, "p90", resp.Frames[0].Fields[4].Name) + assert.Equal(t, "p95", resp.Frames[0].Fields[5].Name) + assert.Equal(t, "max", resp.Frames[0].Fields[6].Name) + }) + + t.Run("parse result as table group by tag", func(t *testing.T) { + resp := ResponseParse(prepare(tableResultFormatInfluxResponse3), 200, &models.Query{RefID: "A", RawQuery: `a nice query`, ResultFormat: "table"}) + assert.Equal(t, 1, len(resp.Frames)) + assert.Equal(t, "a nice query", resp.Frames[0].Meta.ExecutedQueryString) + for i := range resp.Frames[0].Fields { + assert.Equal(t, resp.Frames[0].Fields[0].Len(), resp.Frames[0].Fields[i].Len()) + } + assert.Equal(t, "Time", resp.Frames[0].Fields[0].Name) + assert.Equal(t, "cpu", resp.Frames[0].Fields[1].Name) + assert.Equal(t, resp.Frames[0].Fields[1].Name, resp.Frames[0].Fields[1].Config.DisplayNameFromDS) + assert.Equal(t, toPtr("cpu-total"), resp.Frames[0].Fields[1].At(0)) + assert.Equal(t, toPtr("cpu0"), resp.Frames[0].Fields[1].At(5)) + assert.Equal(t, toPtr("cpu2"), resp.Frames[0].Fields[1].At(12)) + assert.Equal(t, "mean", resp.Frames[0].Fields[2].Name) + assert.Equal(t, resp.Frames[0].Fields[2].Name, resp.Frames[0].Fields[2].Config.DisplayNameFromDS) + }) + + t.Run("parse result without tags as table", func(t *testing.T) { + resp := ResponseParse(prepare(tableResultFormatInfluxResponse4), 200, &models.Query{RefID: "A", RawQuery: `a nice query`, ResultFormat: "table"}) + assert.Equal(t, 1, len(resp.Frames)) + assert.Equal(t, "a nice query", resp.Frames[0].Meta.ExecutedQueryString) + for i := range resp.Frames[0].Fields { + assert.Equal(t, resp.Frames[0].Fields[0].Len(), resp.Frames[0].Fields[i].Len()) + } + assert.Equal(t, "Time", resp.Frames[0].Fields[0].Name) + assert.Equal(t, "mean", resp.Frames[0].Fields[1].Name) + assert.Equal(t, resp.Frames[0].Fields[1].Name, resp.Frames[0].Fields[1].Config.DisplayNameFromDS) + }) + + t.Run("parse show measurements response as table", func(t *testing.T) { + resp := ResponseParse(prepare(showMeasurementsResponse), 200, &models.Query{RefID: "A", RawQuery: `a nice query`, ResultFormat: "table"}) + assert.Equal(t, 1, len(resp.Frames)) + assert.Equal(t, "a nice query", resp.Frames[0].Meta.ExecutedQueryString) + for i := range resp.Frames[0].Fields { + assert.Equal(t, resp.Frames[0].Fields[0].Len(), resp.Frames[0].Fields[i].Len()) + } + assert.Equal(t, 1, len(resp.Frames[0].Fields)) + assert.Equal(t, "name", resp.Frames[0].Fields[0].Name) + }) + + t.Run("parse retention policy response as table", func(t *testing.T) { + resp := ResponseParse(prepare(showRetentionPolicyResponse), 200, &models.Query{RefID: "A", RawQuery: `a nice query`, ResultFormat: "table"}) + assert.Equal(t, 1, len(resp.Frames)) + assert.Equal(t, "a nice query", resp.Frames[0].Meta.ExecutedQueryString) + for i := range resp.Frames[0].Fields { + assert.Equal(t, resp.Frames[0].Fields[0].Len(), resp.Frames[0].Fields[i].Len()) + } + assert.Equal(t, 5, len(resp.Frames[0].Fields)) + assert.Equal(t, "name", resp.Frames[0].Fields[0].Name) + assert.Equal(t, "duration", resp.Frames[0].Fields[1].Name) + assert.Equal(t, "shardGroupDuration", resp.Frames[0].Fields[2].Name) + assert.Equal(t, "replicaN", resp.Frames[0].Fields[3].Name) + assert.Equal(t, "default", resp.Frames[0].Fields[4].Name) + }) +} + func TestResponseParser_Parse(t *testing.T) { tests := []struct { name string @@ -912,10 +1003,9 @@ func TestResponseParser_Parse(t *testing.T) { ] }`, f: func(t *testing.T, got backend.DataResponse) { - assert.Equal(t, "domain", got.Frames[0].Name) + assert.Equal(t, "Annotation", got.Frames[0].Name) assert.Equal(t, "domain", got.Frames[0].Fields[1].Config.DisplayNameFromDS) - assert.Equal(t, "ASD", got.Frames[2].Name) - assert.Equal(t, "ASD", got.Frames[2].Fields[1].Config.DisplayNameFromDS) + assert.Equal(t, "type", got.Frames[0].Fields[2].Config.DisplayNameFromDS) assert.Equal(t, tableVisType, got.Frames[0].Meta.PreferredVisualization) }, }, @@ -930,3 +1020,578 @@ func TestResponseParser_Parse(t *testing.T) { }) } } + +func toPtr[T any](v T) *T { + return &v +} + +const tableResultFormatInfluxResponse1 = `{ + "results": [ + { + "statement_id": 0, + "series": [ + { + "name": "cpu", + "columns": [ + "time", + "usage_idle", + "usage_iowait" + ], + "values": [ + [ + 1700090120000, + 99.0255173802101, + 0.020092425155804713 + ], + [ + 1700090120000, + 99.29718875523953, + 0 + ], + [ + 1700090120000, + 99.09456740445926, + 0 + ], + [ + 1700090120000, + 99.39455095864957, + 0 + ], + [ + 1700090120000, + 99.09729187566201, + 0 + ] + ] + } + ] + } + ] +}` + +const tableResultFormatInfluxResponse2 = `{ + "results": [ + { + "statement_id": 0, + "series": [ + { + "name": "cpu", + "tags": { + "cpu": "cpu-total" + }, + "columns": [ + "time", + "mean", + "min", + "p90", + "p95", + "max" + ], + "values": [ + [ + 1700046000000, + 99.06348570053983, + 97.3214285712978, + 99.2066680055868, + 99.24812030075188, + 99.31809065366402 + ] + ] + }, + { + "name": "cpu", + "tags": { + "cpu": "cpu0" + }, + "columns": [ + "time", + "mean", + "min", + "p90", + "p95", + "max" + ], + "values": [ + [ + 1700046000000, + 98.99671817733766, + 96.65991902847126, + 99.29364278499536, + 99.29718875523953, + 99.59839357421622 + ] + ] + }, + { + "name": "cpu", + "tags": { + "cpu": "cpu1" + }, + "columns": [ + "time", + "mean", + "min", + "p90", + "p95", + "max" + ], + "values": [ + [ + 1700046000000, + 99.03148357927465, + 96.67673715996412, + 99.39698492464545, + 99.39759036146867, + 99.59798994966731 + ] + ] + }, + { + "name": "cpu", + "tags": { + "cpu": "cpu2" + }, + "columns": [ + "time", + "mean", + "min", + "p90", + "p95", + "max" + ], + "values": [ + [ + 1700046000000, + 99.03087433486812, + 96.03658536600605, + 99.29859719431953, + 99.39759036146867, + 99.59879638908582 + ] + ] + }, + { + "name": "cpu", + "tags": { + "cpu": "cpu3" + }, + "columns": [ + "time", + "mean", + "min", + "p90", + "p95", + "max" + ], + "values": [ + [ + 1700046000000, + 99.0796957137731, + 97.37903225797402, + 99.39698492464723, + 99.39879759521435, + 99.4984954865762 + ] + ] + }, + { + "name": "cpu", + "tags": { + "cpu": "cpu4" + }, + "columns": [ + "time", + "mean", + "min", + "p90", + "p95", + "max" + ], + "values": [ + [ + 1700046000000, + 99.09573460946685, + 97.57330637016123, + 99.39759036146867, + 99.49698189117252, + 99.59839357450608 + ] + ] + }, + { + "name": "cpu", + "tags": { + "cpu": "cpu5" + }, + "columns": [ + "time", + "mean", + "min", + "p90", + "p95", + "max" + ], + "values": [ + [ + 1700046000000, + 99.0690883079725, + 96.65991902847126, + 99.39698492464545, + 99.39819458377468, + 99.59798994995865 + ] + ] + }, + { + "name": "cpu", + "tags": { + "cpu": "cpu6" + }, + "columns": [ + "time", + "mean", + "min", + "p90", + "p95", + "max" + ], + "values": [ + [ + 1700046000000, + 99.06475215715605, + 97.37108190081956, + 99.39698492464545, + 99.39879759521259, + 99.69879518073434 + ] + ] + }, + { + "name": "cpu", + "tags": { + "cpu": "cpu7" + }, + "columns": [ + "time", + "mean", + "min", + "p90", + "p95", + "max" + ], + "values": [ + [ + 1700046000000, + 99.06204005079694, + 97.7596741344093, + 99.39637826964127, + 99.39759036147042, + 99.59879638908698 + ] + ] + }, + { + "name": "cpu", + "tags": { + "cpu": "cpu8" + }, + "columns": [ + "time", + "mean", + "min", + "p90", + "p95", + "max" + ], + "values": [ + [ + 1700046000000, + 99.0999818796052, + 96.56565656568982, + 99.39698492464723, + 99.39819458377468, + 99.59758551299777 + ] + ] + }, + { + "name": "cpu", + "tags": { + "cpu": "cpu9" + }, + "columns": [ + "time", + "mean", + "min", + "p90", + "p95", + "max" + ], + "values": [ + [ + 1700046000000, + 99.10477313534511, + 96.8463886063268, + 99.39759036146867, + 99.39819458377468, + 99.59839357421622 + ] + ] + } + ] + } + ] +} +` + +const tableResultFormatInfluxResponse3 = `{ + "results": [ + { + "statement_id": 0, + "series": [ + { + "name": "cpu", + "tags": { + "cpu": "cpu-total" + }, + "columns": [ + "time", + "mean" + ], + "values": [ + [ + 1700046000000, + 99.06919189833442 + ], + [ + 1700047200000, + 99.13105510262923 + ], + [ + 1700048400000, + 98.99236330721192 + ], + [ + 1700049600000, + 98.80510091380069 + ] + ] + }, + { + "name": "cpu", + "tags": { + "cpu": "cpu0" + }, + "columns": [ + "time", + "mean" + ], + "values": [ + [ + 1700046000000, + 99.01372119142576 + ], + [ + 1700047200000, + 99.00430308480553 + ], + [ + 1700048400000, + 98.9737996641964 + ], + [ + 1700049600000, + 98.79638916754935 + ] + ] + }, + { + "name": "cpu", + "tags": { + "cpu": "cpu1" + }, + "columns": [ + "time", + "mean" + ], + "values": [ + [ + 1700046000000, + 99.04949983158023 + ], + [ + 1700047200000, + 99.06989461231551 + ], + [ + 1700048400000, + 98.97954813782476 + ], + [ + 1700049600000, + 98.49246231161365 + ] + ] + }, + { + "name": "cpu", + "tags": { + "cpu": "cpu2" + }, + "columns": [ + "time", + "mean" + ], + "values": [ + [ + 1700046000000, + 99.11296419686643 + ], + [ + 1700047200000, + 99.01817278917116 + ], + [ + 1700048400000, + 98.96847021232013 + ], + [ + 1700049600000, + 98.192771084406 + ] + ] + }, + { + "name": "cpu", + "tags": { + "cpu": "cpu3" + }, + "columns": [ + "time", + "mean" + ], + "values": [ + [ + 1700046000000, + 99.0742704326151 + ], + [ + 1700047200000, + 99.17835628293322 + ], + [ + 1700048400000, + 98.98968994907334 + ], + [ + 1700049600000, + 98.69215291745849 + ] + ] + } + ] + } + ] +} +` + +const tableResultFormatInfluxResponse4 = `{ + "results": [ + { + "statement_id": 0, + "series": [ + { + "name": "cpu", + "columns": [ + "time", + "mean" + ], + "values": [ + [ + 1700046000000, + 99.0693929754458 + ], + [ + 1700047200000, + 99.13073313839024 + ], + [ + 1700048400000, + 98.99278645182834 + ], + [ + 1700049600000, + 98.77818123433566 + ] + ] + } + ] + } + ] +}` + +const showMeasurementsResponse = `{ + "results": [ + { + "statement_id": 0, + "series": [ + { + "name": "measurements", + "columns": [ + "name" + ], + "values": [ + [ + "cpu" + ], + [ + "disk" + ], + [ + "diskio" + ], + [ + "kernel" + ] + ] + } + ] + } + ] +}` + +const showRetentionPolicyResponse = `{ + "results": [ + { + "statement_id": 0, + "series": [ + { + "columns": [ + "name", + "duration", + "shardGroupDuration", + "replicaN", + "default" + ], + "values": [ + [ + "default", + "0s", + "168h0m0s", + 1, + true + ], + [ + "autogen", + "0s", + "168h0m0s", + 1, + false + ] + ] + } + ] + } + ] +}`