Datasource/CloudWatch: Results of CloudWatch Logs stats queries are now grouped (#24396)
* Datasource/CloudWatch: Results of CloudWatch Logs stats queries are now grouped
This commit is contained in:
@@ -4,6 +4,7 @@ import (
|
||||
"context"
|
||||
"fmt"
|
||||
"sort"
|
||||
"strconv"
|
||||
|
||||
"github.com/aws/aws-sdk-go/aws"
|
||||
"github.com/aws/aws-sdk-go/aws/awserr"
|
||||
@@ -29,6 +30,34 @@ func (e *CloudWatchExecutor) executeLogActions(ctx context.Context, queryContext
|
||||
return err
|
||||
}
|
||||
|
||||
// When a query of the form "stats ... by ..." is made, we want to return
|
||||
// one series per group defined in the query, but due to the format
|
||||
// the query response is in, there does not seem to be a way to tell
|
||||
// by the response alone if/how the results should be grouped.
|
||||
// Because of this, if the frontend sees that a "stats ... by ..." query is being made
|
||||
// the "groupResults" parameter is sent along with the query to the backend so that we
|
||||
// can correctly group the CloudWatch logs response.
|
||||
if query.Model.Get("groupResults").MustBool() && len(dataframe.Fields) > 0 {
|
||||
groupingFields := findGroupingFields(dataframe.Fields)
|
||||
|
||||
groupedFrames, err := groupResults(dataframe, groupingFields)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
encodedFrames := make([][]byte, 0)
|
||||
for _, frame := range groupedFrames {
|
||||
dataframeEnc, err := frame.MarshalArrow()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
encodedFrames = append(encodedFrames, dataframeEnc)
|
||||
}
|
||||
|
||||
resultChan <- &tsdb.QueryResult{RefId: query.RefId, Dataframes: encodedFrames}
|
||||
return nil
|
||||
}
|
||||
|
||||
dataframeEnc, err := dataframe.MarshalArrow()
|
||||
if err != nil {
|
||||
return err
|
||||
@@ -56,6 +85,23 @@ func (e *CloudWatchExecutor) executeLogActions(ctx context.Context, queryContext
|
||||
return response, nil
|
||||
}
|
||||
|
||||
func findGroupingFields(fields []*data.Field) []string {
|
||||
groupingFields := make([]string, 0)
|
||||
for _, field := range fields {
|
||||
if field.Type().Numeric() || field.Type() == data.FieldTypeNullableTime || field.Type() == data.FieldTypeTime {
|
||||
continue
|
||||
}
|
||||
|
||||
if _, err := strconv.ParseFloat(*field.At(0).(*string), 64); err == nil {
|
||||
continue
|
||||
}
|
||||
|
||||
groupingFields = append(groupingFields, field.Name)
|
||||
}
|
||||
|
||||
return groupingFields
|
||||
}
|
||||
|
||||
func (e *CloudWatchExecutor) executeLogAction(ctx context.Context, queryContext *tsdb.TsdbQuery, query *tsdb.Query) (*data.Frame, error) {
|
||||
parameters := query.Model
|
||||
subType := query.Model.Get("subtype").MustString()
|
||||
|
||||
@@ -73,3 +73,48 @@ func logsResultsToDataframes(response *cloudwatchlogs.GetQueryResultsOutput) (*d
|
||||
|
||||
return frame, nil
|
||||
}
|
||||
|
||||
func groupResults(results *data.Frame, groupingFieldNames []string) ([]*data.Frame, error) {
|
||||
groupingFields := make([]*data.Field, 0)
|
||||
|
||||
for _, field := range results.Fields {
|
||||
for _, groupingField := range groupingFieldNames {
|
||||
if field.Name == groupingField {
|
||||
groupingFields = append(groupingFields, field)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
rowLength, err := results.RowLen()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
groupedDataFrames := make(map[string]*data.Frame)
|
||||
for i := 0; i < rowLength; i++ {
|
||||
groupKey := generateGroupKey(groupingFields, i)
|
||||
if _, exists := groupedDataFrames[groupKey]; !exists {
|
||||
newFrame := results.EmptyCopy()
|
||||
newFrame.Name = groupKey
|
||||
groupedDataFrames[groupKey] = newFrame
|
||||
}
|
||||
|
||||
groupedDataFrames[groupKey].AppendRow(results.RowCopy(i)...)
|
||||
}
|
||||
|
||||
newDataFrames := make([]*data.Frame, 0, len(groupedDataFrames))
|
||||
for _, dataFrame := range groupedDataFrames {
|
||||
newDataFrames = append(newDataFrames, dataFrame)
|
||||
}
|
||||
|
||||
return newDataFrames, nil
|
||||
}
|
||||
|
||||
func generateGroupKey(fields []*data.Field, row int) string {
|
||||
groupKey := ""
|
||||
for _, field := range fields {
|
||||
groupKey += *field.At(row).(*string)
|
||||
}
|
||||
|
||||
return groupKey
|
||||
}
|
||||
|
||||
@@ -159,3 +159,141 @@ func TestLogsResultsToDataframes(t *testing.T) {
|
||||
assert.Equal(t, expectedDataframe.Meta, dataframes.Meta)
|
||||
assert.ElementsMatch(t, expectedDataframe.Fields, dataframes.Fields)
|
||||
}
|
||||
|
||||
func TestGroupKeyGeneration(t *testing.T) {
|
||||
logField := data.NewField("@log", data.Labels{}, []*string{
|
||||
aws.String("fakelog-a"),
|
||||
aws.String("fakelog-b"),
|
||||
aws.String("fakelog-c"),
|
||||
})
|
||||
|
||||
streamField := data.NewField("stream", data.Labels{}, []*string{
|
||||
aws.String("stream-a"),
|
||||
aws.String("stream-b"),
|
||||
aws.String("stream-c"),
|
||||
})
|
||||
|
||||
fakeFields := []*data.Field{logField, streamField}
|
||||
expectedKeys := []string{"fakelog-astream-a", "fakelog-bstream-b", "fakelog-cstream-c"}
|
||||
|
||||
assert.Equal(t, expectedKeys[0], generateGroupKey(fakeFields, 0))
|
||||
assert.Equal(t, expectedKeys[1], generateGroupKey(fakeFields, 1))
|
||||
assert.Equal(t, expectedKeys[2], generateGroupKey(fakeFields, 2))
|
||||
}
|
||||
|
||||
func TestGroupingResults(t *testing.T) {
|
||||
timeA, _ := time.Parse("2006-01-02 15:04:05.000", "2020-03-02 15:04:05.000")
|
||||
timeB, _ := time.Parse("2006-01-02 15:04:05.000", "2020-03-02 16:04:05.000")
|
||||
timeC, _ := time.Parse("2006-01-02 15:04:05.000", "2020-03-02 17:04:05.000")
|
||||
timeVals := []*time.Time{
|
||||
&timeA, &timeA, &timeA, &timeB, &timeB, &timeB, &timeC, &timeC, &timeC,
|
||||
}
|
||||
timeField := data.NewField("@timestamp", data.Labels{}, timeVals)
|
||||
|
||||
logField := data.NewField("@log", data.Labels{}, []*string{
|
||||
aws.String("fakelog-a"),
|
||||
aws.String("fakelog-b"),
|
||||
aws.String("fakelog-c"),
|
||||
aws.String("fakelog-a"),
|
||||
aws.String("fakelog-b"),
|
||||
aws.String("fakelog-c"),
|
||||
aws.String("fakelog-a"),
|
||||
aws.String("fakelog-b"),
|
||||
aws.String("fakelog-c"),
|
||||
})
|
||||
|
||||
countField := data.NewField("count", data.Labels{}, []*string{
|
||||
aws.String("100"),
|
||||
aws.String("150"),
|
||||
aws.String("20"),
|
||||
aws.String("34"),
|
||||
aws.String("57"),
|
||||
aws.String("62"),
|
||||
aws.String("105"),
|
||||
aws.String("200"),
|
||||
aws.String("99"),
|
||||
})
|
||||
|
||||
fakeDataFrame := &data.Frame{
|
||||
Name: "CloudWatchLogsResponse",
|
||||
Fields: []*data.Field{
|
||||
timeField,
|
||||
logField,
|
||||
countField,
|
||||
},
|
||||
RefID: "",
|
||||
}
|
||||
|
||||
groupedTimeVals := []*time.Time{
|
||||
&timeA, &timeB, &timeC,
|
||||
}
|
||||
groupedTimeField := data.NewField("@timestamp", data.Labels{}, groupedTimeVals)
|
||||
groupedLogFieldA := data.NewField("@log", data.Labels{}, []*string{
|
||||
aws.String("fakelog-a"),
|
||||
aws.String("fakelog-a"),
|
||||
aws.String("fakelog-a"),
|
||||
})
|
||||
|
||||
groupedCountFieldA := data.NewField("count", data.Labels{}, []*string{
|
||||
aws.String("100"),
|
||||
aws.String("34"),
|
||||
aws.String("105"),
|
||||
})
|
||||
|
||||
groupedLogFieldB := data.NewField("@log", data.Labels{}, []*string{
|
||||
aws.String("fakelog-b"),
|
||||
aws.String("fakelog-b"),
|
||||
aws.String("fakelog-b"),
|
||||
})
|
||||
|
||||
groupedCountFieldB := data.NewField("count", data.Labels{}, []*string{
|
||||
aws.String("150"),
|
||||
aws.String("57"),
|
||||
aws.String("200"),
|
||||
})
|
||||
|
||||
groupedLogFieldC := data.NewField("@log", data.Labels{}, []*string{
|
||||
aws.String("fakelog-c"),
|
||||
aws.String("fakelog-c"),
|
||||
aws.String("fakelog-c"),
|
||||
})
|
||||
|
||||
groupedCountFieldC := data.NewField("count", data.Labels{}, []*string{
|
||||
aws.String("20"),
|
||||
aws.String("62"),
|
||||
aws.String("99"),
|
||||
})
|
||||
|
||||
expectedGroupedFrames := []*data.Frame{
|
||||
{
|
||||
Name: "fakelog-a",
|
||||
Fields: []*data.Field{
|
||||
groupedTimeField,
|
||||
groupedLogFieldA,
|
||||
groupedCountFieldA,
|
||||
},
|
||||
RefID: "",
|
||||
},
|
||||
{
|
||||
Name: "fakelog-b",
|
||||
Fields: []*data.Field{
|
||||
groupedTimeField,
|
||||
groupedLogFieldB,
|
||||
groupedCountFieldB,
|
||||
},
|
||||
RefID: "",
|
||||
},
|
||||
{
|
||||
Name: "fakelog-c",
|
||||
Fields: []*data.Field{
|
||||
groupedTimeField,
|
||||
groupedLogFieldC,
|
||||
groupedCountFieldC,
|
||||
},
|
||||
RefID: "",
|
||||
},
|
||||
}
|
||||
|
||||
groupedResults, _ := groupResults(fakeDataFrame, []string{"@log"})
|
||||
assert.ElementsMatch(t, expectedGroupedFrames, groupedResults)
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user