response_parser.go 6.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258
  1. package elasticsearch
  2. import (
  3. "errors"
  4. "fmt"
  5. "github.com/grafana/grafana/pkg/components/null"
  6. "github.com/grafana/grafana/pkg/components/simplejson"
  7. "github.com/grafana/grafana/pkg/tsdb"
  8. "regexp"
  9. "strconv"
  10. "strings"
  11. )
  12. type ElasticsearchResponseParser struct {
  13. Responses []Response
  14. Targets []*Query
  15. }
  16. func (rp *ElasticsearchResponseParser) getTimeSeries() *tsdb.QueryResult {
  17. queryRes := tsdb.NewQueryResult()
  18. for i, res := range rp.Responses {
  19. target := rp.Targets[i]
  20. props := make(map[string]string)
  21. series := make([]*tsdb.TimeSeries, 0)
  22. rp.processBuckets(res.Aggregations, target, &series, props, 0)
  23. rp.nameSeries(&series, target)
  24. queryRes.Series = append(queryRes.Series, series...)
  25. }
  26. return queryRes
  27. }
  28. func (rp *ElasticsearchResponseParser) processBuckets(aggs map[string]interface{}, target *Query, series *[]*tsdb.TimeSeries, props map[string]string, depth int) error {
  29. var err error
  30. maxDepth := len(target.BucketAggs) - 1
  31. for aggId, v := range aggs {
  32. aggDef, _ := findAgg(target, aggId)
  33. esAgg := simplejson.NewFromAny(v)
  34. if aggDef == nil {
  35. continue
  36. }
  37. if depth == maxDepth {
  38. if aggDef.Type == "date_histogram" {
  39. err = rp.processMetrics(esAgg, target, series, props)
  40. if err != nil {
  41. return err
  42. }
  43. } else {
  44. return fmt.Errorf("not support type:%s", aggDef.Type)
  45. }
  46. } else {
  47. for i, b := range esAgg.Get("buckets").MustArray() {
  48. bucket := simplejson.NewFromAny(b)
  49. newProps := props
  50. if key, err := bucket.Get("key").String(); err == nil {
  51. newProps[aggDef.Field] = key
  52. } else {
  53. props["filter"] = strconv.Itoa(i)
  54. }
  55. if key, err := bucket.Get("key_as_string").String(); err == nil {
  56. props[aggDef.Field] = key
  57. }
  58. rp.processBuckets(bucket.MustMap(), target, series, newProps, depth+1)
  59. }
  60. }
  61. }
  62. return nil
  63. }
  64. func (rp *ElasticsearchResponseParser) processMetrics(esAgg *simplejson.Json, target *Query, series *[]*tsdb.TimeSeries, props map[string]string) error {
  65. for _, metric := range target.Metrics {
  66. if metric.Hide {
  67. continue
  68. }
  69. switch metric.Type {
  70. case "count":
  71. newSeries := tsdb.TimeSeries{}
  72. for _, v := range esAgg.Get("buckets").MustArray() {
  73. bucket := simplejson.NewFromAny(v)
  74. value := castToNullFloat(bucket.Get("doc_count"))
  75. key := castToNullFloat(bucket.Get("key"))
  76. newSeries.Points = append(newSeries.Points, tsdb.TimePoint{value, key})
  77. }
  78. newSeries.Tags = props
  79. newSeries.Tags["metric"] = "count"
  80. *series = append(*series, &newSeries)
  81. case "percentiles":
  82. buckets := esAgg.Get("buckets").MustArray()
  83. if len(buckets) == 0 {
  84. break
  85. }
  86. firstBucket := simplejson.NewFromAny(buckets[0])
  87. percentiles := firstBucket.GetPath(metric.ID, "values").MustMap()
  88. for percentileName := range percentiles {
  89. newSeries := tsdb.TimeSeries{}
  90. newSeries.Tags = props
  91. newSeries.Tags["metric"] = "p" + percentileName
  92. newSeries.Tags["field"] = metric.Field
  93. for _, v := range buckets {
  94. bucket := simplejson.NewFromAny(v)
  95. value := castToNullFloat(bucket.GetPath(metric.ID, "values", percentileName))
  96. key := castToNullFloat(bucket.Get("key"))
  97. newSeries.Points = append(newSeries.Points, tsdb.TimePoint{value, key})
  98. }
  99. *series = append(*series, &newSeries)
  100. }
  101. default:
  102. newSeries := tsdb.TimeSeries{}
  103. newSeries.Tags = props
  104. newSeries.Tags["metric"] = metric.Type
  105. newSeries.Tags["field"] = metric.Field
  106. for _, v := range esAgg.Get("buckets").MustArray() {
  107. bucket := simplejson.NewFromAny(v)
  108. key := castToNullFloat(bucket.Get("key"))
  109. valueObj, err := bucket.Get(metric.ID).Map()
  110. if err != nil {
  111. break
  112. }
  113. var value null.Float
  114. if _, ok := valueObj["normalized_value"]; ok {
  115. value = castToNullFloat(bucket.GetPath(metric.ID, "normalized_value"))
  116. } else {
  117. value = castToNullFloat(bucket.GetPath(metric.ID, "value"))
  118. }
  119. newSeries.Points = append(newSeries.Points, tsdb.TimePoint{value, key})
  120. }
  121. *series = append(*series, &newSeries)
  122. }
  123. }
  124. return nil
  125. }
  126. func (rp *ElasticsearchResponseParser) nameSeries(seriesList *[]*tsdb.TimeSeries, target *Query) {
  127. set := make(map[string]string)
  128. for _, v := range *seriesList {
  129. if metricType, exists := v.Tags["metric"]; exists {
  130. if _, ok := set[metricType]; !ok {
  131. set[metricType] = ""
  132. }
  133. }
  134. }
  135. metricTypeCount := len(set)
  136. for _, series := range *seriesList {
  137. series.Name = rp.getSeriesName(series, target, metricTypeCount)
  138. }
  139. }
  140. func (rp *ElasticsearchResponseParser) getSeriesName(series *tsdb.TimeSeries, target *Query, metricTypeCount int) string {
  141. metricType := series.Tags["metric"]
  142. metricName := rp.getMetricName(metricType)
  143. delete(series.Tags, "metric")
  144. field := ""
  145. if v, ok := series.Tags["field"]; ok {
  146. field = v
  147. delete(series.Tags, "field")
  148. }
  149. if target.Alias != "" {
  150. var re = regexp.MustCompile(`{{([\s\S]+?)}}`)
  151. for _, match := range re.FindAllString(target.Alias, -1) {
  152. group := match[2 : len(match)-2]
  153. if strings.HasPrefix(group, "term ") {
  154. if term, ok := series.Tags["term "]; ok {
  155. strings.Replace(target.Alias, match, term, 1)
  156. }
  157. }
  158. if v, ok := series.Tags[group]; ok {
  159. strings.Replace(target.Alias, match, v, 1)
  160. }
  161. switch group {
  162. case "metric":
  163. strings.Replace(target.Alias, match, metricName, 1)
  164. case "field":
  165. strings.Replace(target.Alias, match, field, 1)
  166. }
  167. }
  168. }
  169. // todo, if field and pipelineAgg
  170. if field != "" && isPipelineAgg(metricType) {
  171. found := false
  172. for _, metric := range target.Metrics {
  173. if metric.ID == field {
  174. metricName += " " + describeMetric(metric.Type, field)
  175. found = true
  176. }
  177. }
  178. if !found {
  179. metricName = "Unset"
  180. }
  181. } else if field != "" {
  182. metricName += " " + field
  183. }
  184. if len(series.Tags) == 0 {
  185. return metricName
  186. }
  187. name := ""
  188. for _, v := range series.Tags {
  189. name += v + " "
  190. }
  191. if metricTypeCount == 1 {
  192. return strings.TrimSpace(name)
  193. }
  194. return strings.TrimSpace(name) + " " + metricName
  195. }
  196. func (rp *ElasticsearchResponseParser) getMetricName(metric string) string {
  197. if text, ok := metricAggType[metric]; ok {
  198. return text
  199. }
  200. if text, ok := extendedStats[metric]; ok {
  201. return text
  202. }
  203. return metric
  204. }
  205. func castToNullFloat(j *simplejson.Json) null.Float {
  206. f, err := j.Float64()
  207. if err == nil {
  208. return null.FloatFrom(f)
  209. }
  210. s, err := j.String()
  211. if err == nil {
  212. v, _ := strconv.ParseFloat(s, 64)
  213. return null.FloatFromPtr(&v)
  214. }
  215. return null.NewFloat(0, false)
  216. }
  217. func findAgg(target *Query, aggId string) (*BucketAgg, error) {
  218. for _, v := range target.BucketAggs {
  219. if aggId == v.ID {
  220. return v, nil
  221. }
  222. }
  223. return nil, errors.New("can't found aggDef, aggID:" + aggId)
  224. }