Elasticsearch: Handle keyed filters buckets and emit frames (#113478)

This commit is contained in:
Adam Yeats
2025-11-06 17:20:08 -06:00
committed by GitHub
parent f468597ad8
commit b9b1028b91
2 changed files with 335 additions and 80 deletions
+199 -80
View File
@@ -748,78 +748,191 @@ func processMetrics(esAgg *simplejson.Json, target *Query, query *backend.DataRe
return nil
}
func processAggregationDocs(esAgg *simplejson.Json, aggDef *BucketAgg, target *Query,
queryResult *backend.DataResponse, props map[string]string) error {
propKeys := createPropKeys(props)
frames := data.Frames{}
fields := createFields(queryResult.Frames, propKeys)
for _, v := range esAgg.Get("buckets").MustArray() {
bucket := simplejson.NewFromAny(v)
var values []interface{}
found := false
for _, field := range fields {
for _, propKey := range propKeys {
if field.Name == propKey {
value := props[propKey]
field.Append(&value)
}
}
if field.Name == aggDef.Field {
found = true
if key, err := bucket.Get("key").String(); err == nil {
field.Append(&key)
} else {
f, err := bucket.Get("key").Float64()
if err != nil {
return fmt.Errorf("error appending bucket key to existing field with name %s: %w", field.Name, err)
}
field.Append(&f)
}
}
}
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 fmt.Errorf("error appending bucket key to new field with name %s: %w", aggDef.Field, err)
}
aggDefField = extractDataField(aggDef.Field, &f)
aggDefField.Append(&f)
}
fields = append(fields, aggDefField)
}
for _, metric := range target.Metrics {
switch metric.Type {
case countType:
addMetricValueToFields(&fields, values, getMetricName(metric.Type), castToFloat(bucket.Get("doc_count")))
case extendedStatsType:
addExtendedStatsToFields(&fields, bucket, metric, values)
case percentilesType:
addPercentilesToFields(&fields, bucket, metric, values)
case topMetricsType:
addTopMetricsToFields(&fields, bucket, metric, values)
default:
addOtherMetricsToFields(&fields, bucket, metric, values, target)
}
}
var dataFields []*data.Field
dataFields = append(dataFields, fields...)
frames = data.Frames{
&data.Frame{
Fields: dataFields,
}}
// ensurePropFields guarantees all property columns exist even if prior frames lacked them
func ensurePropFields(fields *[]*data.Field, keys []string) {
have := map[string]bool{}
for _, f := range *fields {
have[f.Name] = true
}
queryResult.Frames = frames
for _, k := range keys {
if !have[k] {
d := ""
f := extractDataField(k, &d)
*fields = append(*fields, f)
}
}
}
// appendPropsRow appends one row of property values; skipKey avoids double-append
func appendPropsRow(fields *[]*data.Field, props map[string]string, propKeys []string, skipKey string) {
for _, f := range *fields {
for _, pk := range propKeys {
if pk == skipKey {
continue
}
if f.Name == pk {
val := props[pk]
f.Append(&val)
}
}
}
}
// appendMetrics appends all metric values for a single bucket/row
func appendMetrics(fields *[]*data.Field, bucket *simplejson.Json, target *Query) {
var values []interface{}
for _, metric := range target.Metrics {
switch metric.Type {
case countType:
addMetricValueToFields(fields, values, getMetricName(metric.Type), castToFloat(bucket.Get("doc_count")))
case extendedStatsType:
addExtendedStatsToFields(fields, bucket, metric, values)
case percentilesType:
addPercentilesToFields(fields, bucket, metric, values)
case topMetricsType:
addTopMetricsToFields(fields, bucket, metric, values)
default:
addOtherMetricsToFields(fields, bucket, metric, values, target)
}
}
}
// appendKeyColumnString appends a string key to an existing field or creates it
func appendKeyColumnString(fields *[]*data.Field, fieldName, key string) {
for _, f := range *fields {
if f.Name == fieldName {
k := key
f.Append(&k)
return
}
}
k := key
f := extractDataField(fieldName, &k)
f.Append(&k)
*fields = append(*fields, f)
}
// appendBucketKeyValue appends the bucket's "key" (string or number) to fieldName
func appendBucketKeyValue(fields *[]*data.Field, fieldName string, bucket *simplejson.Json) error {
for _, f := range *fields {
if f.Name == fieldName {
if s, err := bucket.Get("key").String(); err == nil {
f.Append(&s)
return nil
}
num, err := bucket.Get("key").Float64()
if err != nil {
return fmt.Errorf("error appending bucket key to existing field %q: %w", fieldName, err)
}
f.Append(&num)
return nil
}
}
// field not present yet
if s, err := bucket.Get("key").String(); err == nil {
f := extractDataField(fieldName, &s)
f.Append(&s)
*fields = append(*fields, f)
return nil
}
num, err := bucket.Get("key").Float64()
if err != nil {
return fmt.Errorf("error appending bucket key to new field %q: %w", fieldName, err)
}
f := extractDataField(fieldName, &num)
f.Append(&num)
*fields = append(*fields, f)
return nil
}
func processAggregationDocs(
esAgg *simplejson.Json,
aggDef *BucketAgg,
target *Query,
queryResult *backend.DataResponse,
props map[string]string,
) error {
propKeys := createPropKeys(props)
buckets := esAgg.Get("buckets")
if arr := buckets.MustArray(); len(arr) > 0 {
fields := createFields(queryResult.Frames, propKeys)
ensurePropFields(&fields, propKeys)
for _, v := range arr {
bucket := simplejson.NewFromAny(v)
appendPropsRow(&fields, props, propKeys, "")
if aggDef.Field != "" {
if err := appendBucketKeyValue(&fields, aggDef.Field, bucket); err != nil {
return err
}
}
appendMetrics(&fields, bucket, target)
}
queryResult.Frames = data.Frames{&data.Frame{Fields: fields}}
return nil
}
if m := buckets.MustMap(); len(m) > 0 {
// default key column to "filter" for leaf filters
keyFieldName := aggDef.Field
if keyFieldName == "" {
keyFieldName = "filter"
}
// ensure "filter" exists among props
hasFilter := false
for _, pk := range propKeys {
if pk == "filter" {
hasFilter = true
break
}
}
if !hasFilter {
propKeys = append(propKeys, "filter")
}
fields := createFields(queryResult.Frames, propKeys)
ensurePropFields(&fields, propKeys)
keys := make([]string, 0, len(m))
for k := range m {
keys = append(keys, k)
}
sort.Strings(keys)
for _, k := range keys {
bucket := simplejson.NewFromAny(m[k])
locProps := make(map[string]string, len(props)+1)
for kk, vv := range props {
locProps[kk] = vv
}
locProps["filter"] = k
// avoid double-append when the key column is "filter"
skip := ""
if keyFieldName == "filter" {
skip = "filter"
}
appendPropsRow(&fields, locProps, propKeys, skip)
appendKeyColumnString(&fields, keyFieldName, k)
appendMetrics(&fields, bucket, target)
}
queryResult.Frames = data.Frames{&data.Frame{Fields: fields}}
return nil
}
// no buckets present
queryResult.Frames = data.Frames{}
return nil
}
@@ -1223,17 +1336,23 @@ func setLogsCustomMeta(frame *data.Frame, searchWords map[string]bool, limit int
func createFields(frames data.Frames, propKeys []string) []*data.Field {
var fields []*data.Field
// Otherwise use the fields from frames
if frames != nil {
for _, frame := range frames {
fields = append(fields, frame.Fields...)
}
// If we have no frames, we create fields from propKeys
} else {
for _, propKey := range propKeys {
fields = append(fields, data.NewField(propKey, nil, []*string{}))
have := map[string]bool{}
// collect existing fields
for _, frame := range frames {
for _, f := range frame.Fields {
fields = append(fields, f)
have[f.Name] = true
}
}
// add missing prop fields
for _, pk := range propKeys {
if !have[pk] {
fields = append(fields, data.NewField(pk, nil, []*string{}))
}
}
return fields
}
@@ -3648,6 +3648,142 @@ func TestTrimEdges(t *testing.T) {
requireFrameLength(t, frames[0], 1)
}
func TestFiltersAggregation_KeyedBuckets(t *testing.T) {
t.Run("Leaf filters (keyed buckets) returns a table with filter | Count", func(t *testing.T) {
targets := map[string]string{
"A": `{
"metrics": [{ "type": "count", "id": "1" }],
"bucketAggs": [{
"type": "filters",
"id": "2",
"settings": {
"filters": [
{ "label": "a 0-1 min", "query": "duration_seconds:[0 TO 60}" },
{ "label": "b 1-5 min", "query": "duration_seconds:[60 TO 300}" }
]
}
}]
}`,
}
// ES returns a keyed map for filters buckets (labels -> bucket)
response := `{
"responses": [{
"aggregations": {
"2": {
"buckets": {
"a 0-1 min": { "doc_count": 12 },
"b 1-5 min": { "doc_count": 39 }
}
}
}
}]
}`
result, err := parseTestResponse(targets, response, false)
require.NoError(t, err)
require.Len(t, result.Responses, 1)
res := result.Responses["A"]
require.NotNil(t, res)
require.NoError(t, res.Error)
frames := res.Frames
require.Len(t, frames, 1)
frame := frames[0]
// expect exactly 2 rows (one per filter bucket) and 2 columns: filter | Count
requireFrameLength(t, frame, 2)
require.Len(t, frame.Fields, 2)
// build field map for stable assertions
fieldMap := map[string]*data.Field{}
for _, f := range frame.Fields {
fieldMap[f.Name] = f
}
require.Contains(t, fieldMap, "filter")
require.Contains(t, fieldMap, "Count")
// keys are sorted lexicographically in the parser, so "a 0-1 min" then "b 1-5 min"
requireStringAt(t, "a 0-1 min", fieldMap["filter"], 0)
requireStringAt(t, "b 1-5 min", fieldMap["filter"], 1)
requireFloatAt(t, 12, fieldMap["Count"], 0)
requireFloatAt(t, 39, fieldMap["Count"], 1)
})
t.Run("Filters -> Terms keeps the filter column and yields filter | type | Count", func(t *testing.T) {
targets := map[string]string{
"A": `{
"metrics": [{ "type": "count", "id": "1" }],
"bucketAggs": [
{
"type": "filters",
"id": "2",
"settings": {
"filters": [
{ "label": "A", "query": "duration_seconds:[0 TO 60}" },
{ "label": "B", "query": "duration_seconds:[60 TO 300}" }
]
}
},
{ "type": "terms", "field": "type", "id": "3" }
]
}`,
}
response := `{
"responses": [{
"aggregations": {
"2": {
"buckets": {
"A": { "3": { "buckets": [ { "key": "pull_request", "doc_count": 10 } ] } },
"B": { "3": { "buckets": [ { "key": "pull_request", "doc_count": 5 } ] } }
}
}
}
}]
}`
result, err := parseTestResponse(targets, response, false)
require.NoError(t, err)
require.Len(t, result.Responses, 1)
res := result.Responses["A"]
require.NotNil(t, res)
require.NoError(t, res.Error)
frames := res.Frames
require.Len(t, frames, 1)
frame := frames[0]
// expect 2 rows (A, B) and 3 columns: filter | type | Count
requireFrameLength(t, frame, 2)
require.Len(t, frame.Fields, 3)
fieldMap := map[string]*data.Field{}
for _, f := range frame.Fields {
fieldMap[f.Name] = f
}
require.Contains(t, fieldMap, "filter")
require.Contains(t, fieldMap, "type")
require.Contains(t, fieldMap, "Count")
// filters are sorted lexicographically: A, B
requireStringAt(t, "A", fieldMap["filter"], 0)
requireStringAt(t, "B", fieldMap["filter"], 1)
// terms field "type" should repeat "pull_request" on each row
requireStringAt(t, "pull_request", fieldMap["type"], 0)
requireStringAt(t, "pull_request", fieldMap["type"], 1)
requireFloatAt(t, 10, fieldMap["Count"], 0)
requireFloatAt(t, 5, fieldMap["Count"], 1)
})
}
func parseTestResponse(tsdbQueries map[string]string, responseBody string, keepLabelsInResponse bool) (*backend.QueryDataResponse, error) {
from := time.Date(2018, 5, 15, 17, 50, 0, 0, time.UTC)
to := time.Date(2018, 5, 15, 17, 55, 0, 0, time.UTC)