Elasticsearch: Fix using of individual query time ranges when querying (#84201)
* WIP - proof of concept * Update pkg/tsdb/elasticsearch/client/client.go Co-authored-by: Sven Grossmann <sven.grossmann@grafana.com> * update and add test * lint * Fix lint * Bring back logging when creating client --------- Co-authored-by: Sven Grossmann <sven.grossmann@grafana.com>
This commit is contained in:
co-authored by
Sven Grossmann
parent
8ed4d5127f
commit
2a1a5145d0
@@ -13,7 +13,6 @@ import (
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/grafana/grafana-plugin-sdk-go/backend"
|
||||
"go.opentelemetry.io/otel/attribute"
|
||||
"go.opentelemetry.io/otel/codes"
|
||||
"go.opentelemetry.io/otel/trace"
|
||||
@@ -56,7 +55,7 @@ type Client interface {
|
||||
}
|
||||
|
||||
// NewClient creates a new elasticsearch client
|
||||
var NewClient = func(ctx context.Context, ds *DatasourceInfo, timeRange backend.TimeRange, logger log.Logger, tracer tracing.Tracer) (Client, error) {
|
||||
var NewClient = func(ctx context.Context, ds *DatasourceInfo, logger log.Logger, tracer tracing.Tracer) (Client, error) {
|
||||
logger = logger.New("entity", "client")
|
||||
|
||||
ip, err := newIndexPattern(ds.Interval, ds.Database)
|
||||
@@ -65,19 +64,14 @@ var NewClient = func(ctx context.Context, ds *DatasourceInfo, timeRange backend.
|
||||
return nil, err
|
||||
}
|
||||
|
||||
indices, err := ip.GetIndices(timeRange)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
logger.Debug("Creating new client", "configuredFields", fmt.Sprintf("%#v", ds.ConfiguredFields), "indices", strings.Join(indices, ", "), "interval", ds.Interval, "index", ds.Database)
|
||||
logger.Debug("Creating new client", "configuredFields", fmt.Sprintf("%#v", ds.ConfiguredFields), "interval", ds.Interval, "index", ds.Database)
|
||||
|
||||
return &baseClientImpl{
|
||||
logger: logger,
|
||||
ctx: ctx,
|
||||
ds: ds,
|
||||
configuredFields: ds.ConfiguredFields,
|
||||
indices: indices,
|
||||
timeRange: timeRange,
|
||||
indexPattern: ip,
|
||||
tracer: tracer,
|
||||
}, nil
|
||||
}
|
||||
@@ -86,8 +80,7 @@ type baseClientImpl struct {
|
||||
ctx context.Context
|
||||
ds *DatasourceInfo
|
||||
configuredFields ConfiguredFields
|
||||
indices []string
|
||||
timeRange backend.TimeRange
|
||||
indexPattern IndexPattern
|
||||
logger log.Logger
|
||||
tracer tracing.Tracer
|
||||
}
|
||||
@@ -239,11 +232,16 @@ func (c *baseClientImpl) createMultiSearchRequests(searchRequests []*SearchReque
|
||||
multiRequests := []*multiRequest{}
|
||||
|
||||
for _, searchReq := range searchRequests {
|
||||
indices, err := c.indexPattern.GetIndices(searchReq.TimeRange)
|
||||
if err != nil {
|
||||
c.logger.Error("Failed to get indices from index pattern", "error", err)
|
||||
continue
|
||||
}
|
||||
mr := multiRequest{
|
||||
header: map[string]any{
|
||||
"search_type": "query_then_fetch",
|
||||
"ignore_unavailable": true,
|
||||
"index": strings.Join(c.indices, ","),
|
||||
"index": strings.Join(indices, ","),
|
||||
},
|
||||
body: searchReq,
|
||||
interval: searchReq.Interval,
|
||||
|
||||
@@ -68,7 +68,7 @@ func TestClient_ExecuteMultisearch(t *testing.T) {
|
||||
To: to,
|
||||
}
|
||||
|
||||
c, err := NewClient(context.Background(), &ds, timeRange, log.New("test", "test"), tracing.InitializeTracerForTest())
|
||||
c, err := NewClient(context.Background(), &ds, log.New("test", "test"), tracing.InitializeTracerForTest())
|
||||
require.NoError(t, err)
|
||||
require.NotNil(t, c)
|
||||
|
||||
@@ -76,7 +76,7 @@ func TestClient_ExecuteMultisearch(t *testing.T) {
|
||||
ts.Close()
|
||||
})
|
||||
|
||||
ms, err := createMultisearchForTest(t, c)
|
||||
ms, err := createMultisearchForTest(t, c, timeRange)
|
||||
require.NoError(t, err)
|
||||
res, err := c.ExecuteMultisearch(ms)
|
||||
require.NoError(t, err)
|
||||
@@ -111,6 +111,79 @@ func TestClient_ExecuteMultisearch(t *testing.T) {
|
||||
assert.Equal(t, 200, res.Status)
|
||||
require.Len(t, res.Responses, 1)
|
||||
})
|
||||
|
||||
t.Run("Given a fake http client, 2 queries and a client with response", func(t *testing.T) {
|
||||
var requestBody *bytes.Buffer
|
||||
ts := httptest.NewServer(http.HandlerFunc(func(rw http.ResponseWriter, r *http.Request) {
|
||||
buf, err := io.ReadAll(r.Body)
|
||||
require.NoError(t, err)
|
||||
|
||||
requestBody = bytes.NewBuffer(buf)
|
||||
|
||||
rw.Header().Set("Content-Type", "application/x-ndjson")
|
||||
_, err = rw.Write([]byte(
|
||||
`{
|
||||
"responses": [
|
||||
{
|
||||
"hits": { "hits": [], "max_score": 0, "total": { "value": 4656, "relation": "eq"} },
|
||||
"status": 200
|
||||
}
|
||||
]
|
||||
}`))
|
||||
require.NoError(t, err)
|
||||
rw.WriteHeader(200)
|
||||
}))
|
||||
|
||||
configuredFields := ConfiguredFields{
|
||||
TimeField: "testtime",
|
||||
LogMessageField: "line",
|
||||
LogLevelField: "lvl",
|
||||
}
|
||||
|
||||
ds := DatasourceInfo{
|
||||
URL: ts.URL,
|
||||
HTTPClient: ts.Client(),
|
||||
Database: "[metrics-]YYYY.MM.DD",
|
||||
ConfiguredFields: configuredFields,
|
||||
Interval: "Daily",
|
||||
MaxConcurrentShardRequests: 6,
|
||||
IncludeFrozen: true,
|
||||
XPack: true,
|
||||
}
|
||||
|
||||
from := time.Date(2018, 5, 15, 17, 50, 0, 0, time.UTC)
|
||||
to := time.Date(2018, 5, 15, 17, 55, 0, 0, time.UTC)
|
||||
timeRange := backend.TimeRange{
|
||||
From: from,
|
||||
To: to,
|
||||
}
|
||||
|
||||
from2 := time.Date(2018, 5, 17, 17, 50, 0, 0, time.UTC)
|
||||
to2 := time.Date(2018, 5, 17, 17, 55, 0, 0, time.UTC)
|
||||
timeRange2 := backend.TimeRange{
|
||||
From: from2,
|
||||
To: to2,
|
||||
}
|
||||
|
||||
c, err := NewClient(context.Background(), &ds, log.New("test", "test"), tracing.InitializeTracerForTest())
|
||||
require.NoError(t, err)
|
||||
require.NotNil(t, c)
|
||||
|
||||
t.Cleanup(func() {
|
||||
ts.Close()
|
||||
})
|
||||
|
||||
ms, err := createMultisearchWithMultipleQueriesForTest(t, c, timeRange, timeRange2)
|
||||
require.NoError(t, err)
|
||||
_, err = c.ExecuteMultisearch(ms)
|
||||
require.NoError(t, err)
|
||||
|
||||
require.NotNil(t, requestBody)
|
||||
|
||||
bodyString := requestBody.String()
|
||||
require.Contains(t, bodyString, "metrics-2018.05.15")
|
||||
require.Contains(t, bodyString, "metrics-2018.05.17")
|
||||
})
|
||||
}
|
||||
|
||||
func TestClient_Index(t *testing.T) {
|
||||
@@ -190,7 +263,7 @@ func TestClient_Index(t *testing.T) {
|
||||
To: to,
|
||||
}
|
||||
|
||||
c, err := NewClient(context.Background(), &ds, timeRange, log.New("test", "test"), tracing.InitializeTracerForTest())
|
||||
c, err := NewClient(context.Background(), &ds, log.New("test", "test"), tracing.InitializeTracerForTest())
|
||||
require.NoError(t, err)
|
||||
require.NotNil(t, c)
|
||||
|
||||
@@ -198,7 +271,7 @@ func TestClient_Index(t *testing.T) {
|
||||
ts.Close()
|
||||
})
|
||||
|
||||
ms, err := createMultisearchForTest(t, c)
|
||||
ms, err := createMultisearchForTest(t, c, timeRange)
|
||||
require.NoError(t, err)
|
||||
_, err = c.ExecuteMultisearch(ms)
|
||||
require.NoError(t, err)
|
||||
@@ -217,11 +290,11 @@ func TestClient_Index(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func createMultisearchForTest(t *testing.T, c Client) (*MultiSearchRequest, error) {
|
||||
func createMultisearchForTest(t *testing.T, c Client, timeRange backend.TimeRange) (*MultiSearchRequest, error) {
|
||||
t.Helper()
|
||||
|
||||
msb := c.MultiSearch()
|
||||
s := msb.Search(15 * time.Second)
|
||||
s := msb.Search(15*time.Second, timeRange)
|
||||
s.Agg().DateHistogram("2", "@timestamp", func(a *DateHistogramAgg, ab AggBuilder) {
|
||||
a.FixedInterval = "$__interval"
|
||||
|
||||
@@ -231,3 +304,28 @@ func createMultisearchForTest(t *testing.T, c Client) (*MultiSearchRequest, erro
|
||||
})
|
||||
return msb.Build()
|
||||
}
|
||||
|
||||
func createMultisearchWithMultipleQueriesForTest(t *testing.T, c Client, firstTimeRange backend.TimeRange, secondTimeRange backend.TimeRange) (*MultiSearchRequest, error) {
|
||||
t.Helper()
|
||||
|
||||
msb := c.MultiSearch()
|
||||
s1 := msb.Search(15*time.Second, firstTimeRange)
|
||||
s1.Agg().DateHistogram("2", "@timestamp", func(a *DateHistogramAgg, ab AggBuilder) {
|
||||
a.FixedInterval = "$__interval"
|
||||
|
||||
ab.Metric("1", "avg", "@hostname", func(a *MetricAggregation) {
|
||||
a.Settings["script"] = "$__interval_ms*@hostname"
|
||||
})
|
||||
})
|
||||
|
||||
s2 := msb.Search(15*time.Second, secondTimeRange)
|
||||
s2.Agg().DateHistogram("2", "@timestamp", func(a *DateHistogramAgg, ab AggBuilder) {
|
||||
a.FixedInterval = "$__interval"
|
||||
|
||||
ab.Metric("1", "avg", "@hostname", func(a *MetricAggregation) {
|
||||
a.Settings["script"] = "$__interval_ms*@hostname"
|
||||
})
|
||||
})
|
||||
|
||||
return msb.Build()
|
||||
}
|
||||
|
||||
@@ -18,11 +18,11 @@ const (
|
||||
intervalYearly = "yearly"
|
||||
)
|
||||
|
||||
type indexPattern interface {
|
||||
type IndexPattern interface {
|
||||
GetIndices(timeRange backend.TimeRange) ([]string, error)
|
||||
}
|
||||
|
||||
var newIndexPattern = func(interval string, pattern string) (indexPattern, error) {
|
||||
var newIndexPattern = func(interval string, pattern string) (IndexPattern, error) {
|
||||
if interval == noInterval {
|
||||
return &staticIndexPattern{indexName: pattern}, nil
|
||||
}
|
||||
|
||||
@@ -3,6 +3,8 @@ package es
|
||||
import (
|
||||
"encoding/json"
|
||||
"time"
|
||||
|
||||
"github.com/grafana/grafana-plugin-sdk-go/backend"
|
||||
)
|
||||
|
||||
// SearchRequest represents a search request
|
||||
@@ -14,6 +16,7 @@ type SearchRequest struct {
|
||||
Query *Query
|
||||
Aggs AggArray
|
||||
CustomProps map[string]interface{}
|
||||
TimeRange backend.TimeRange
|
||||
}
|
||||
|
||||
// MarshalJSON returns the JSON encoding of the request.
|
||||
|
||||
@@ -3,6 +3,8 @@ package es
|
||||
import (
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/grafana/grafana-plugin-sdk-go/backend"
|
||||
)
|
||||
|
||||
const (
|
||||
@@ -22,15 +24,17 @@ type SearchRequestBuilder struct {
|
||||
queryBuilder *QueryBuilder
|
||||
aggBuilders []AggBuilder
|
||||
customProps map[string]any
|
||||
timeRange backend.TimeRange
|
||||
}
|
||||
|
||||
// NewSearchRequestBuilder create a new search request builder
|
||||
func NewSearchRequestBuilder(interval time.Duration) *SearchRequestBuilder {
|
||||
func NewSearchRequestBuilder(interval time.Duration, timeRange backend.TimeRange) *SearchRequestBuilder {
|
||||
builder := &SearchRequestBuilder{
|
||||
interval: interval,
|
||||
sort: make(map[string]any),
|
||||
customProps: make(map[string]any),
|
||||
aggBuilders: make([]AggBuilder, 0),
|
||||
timeRange: timeRange,
|
||||
}
|
||||
return builder
|
||||
}
|
||||
@@ -39,6 +43,7 @@ func NewSearchRequestBuilder(interval time.Duration) *SearchRequestBuilder {
|
||||
func (b *SearchRequestBuilder) Build() (*SearchRequest, error) {
|
||||
sr := SearchRequest{
|
||||
Index: b.index,
|
||||
TimeRange: b.timeRange,
|
||||
Interval: b.interval,
|
||||
Size: b.size,
|
||||
Sort: b.sort,
|
||||
@@ -164,8 +169,8 @@ func NewMultiSearchRequestBuilder() *MultiSearchRequestBuilder {
|
||||
}
|
||||
|
||||
// Search initiates and returns a new search request builder
|
||||
func (m *MultiSearchRequestBuilder) Search(interval time.Duration) *SearchRequestBuilder {
|
||||
b := NewSearchRequestBuilder(interval)
|
||||
func (m *MultiSearchRequestBuilder) Search(interval time.Duration, timeRange backend.TimeRange) *SearchRequestBuilder {
|
||||
b := NewSearchRequestBuilder(interval, timeRange)
|
||||
m.requestBuilders = append(m.requestBuilders, b)
|
||||
return b
|
||||
}
|
||||
|
||||
@@ -7,14 +7,21 @@ import (
|
||||
|
||||
"github.com/stretchr/testify/require"
|
||||
|
||||
"github.com/grafana/grafana-plugin-sdk-go/backend"
|
||||
"github.com/grafana/grafana/pkg/components/simplejson"
|
||||
)
|
||||
|
||||
func TestSearchRequest(t *testing.T) {
|
||||
timeField := "@timestamp"
|
||||
from := time.Date(2018, 5, 10, 17, 50, 0, 0, time.UTC)
|
||||
to := time.Date(2018, 5, 12, 17, 55, 0, 0, time.UTC)
|
||||
timeRange := backend.TimeRange{
|
||||
From: from,
|
||||
To: to,
|
||||
}
|
||||
|
||||
setup := func() *SearchRequestBuilder {
|
||||
return NewSearchRequestBuilder(15 * time.Second)
|
||||
return NewSearchRequestBuilder(15*time.Second, timeRange)
|
||||
}
|
||||
|
||||
t.Run("When building search request", func(t *testing.T) {
|
||||
@@ -398,9 +405,15 @@ func TestSearchRequest(t *testing.T) {
|
||||
}
|
||||
|
||||
func TestMultiSearchRequest(t *testing.T) {
|
||||
from := time.Date(2018, 5, 10, 17, 50, 0, 0, time.UTC)
|
||||
to := time.Date(2018, 5, 12, 17, 55, 0, 0, time.UTC)
|
||||
timeRange := backend.TimeRange{
|
||||
From: from,
|
||||
To: to,
|
||||
}
|
||||
t.Run("When adding one search request", func(t *testing.T) {
|
||||
b := NewMultiSearchRequestBuilder()
|
||||
b.Search(15 * time.Second)
|
||||
b.Search(15*time.Second, timeRange)
|
||||
|
||||
t.Run("When building search request should contain one search request", func(t *testing.T) {
|
||||
mr, err := b.Build()
|
||||
@@ -411,8 +424,8 @@ func TestMultiSearchRequest(t *testing.T) {
|
||||
|
||||
t.Run("When adding two search requests", func(t *testing.T) {
|
||||
b := NewMultiSearchRequestBuilder()
|
||||
b.Search(15 * time.Second)
|
||||
b.Search(15 * time.Second)
|
||||
b.Search(15*time.Second, timeRange)
|
||||
b.Search(15*time.Second, timeRange)
|
||||
|
||||
t.Run("When building search request should contain two search requests", func(t *testing.T) {
|
||||
mr, err := b.Build()
|
||||
|
||||
@@ -53,9 +53,9 @@ func (e *elasticsearchDataQuery) execute() (*backend.QueryDataResponse, error) {
|
||||
|
||||
ms := e.client.MultiSearch()
|
||||
|
||||
from := e.dataQueries[0].TimeRange.From.UnixNano() / int64(time.Millisecond)
|
||||
to := e.dataQueries[0].TimeRange.To.UnixNano() / int64(time.Millisecond)
|
||||
for _, q := range queries {
|
||||
from := q.TimeRange.From.UnixNano() / int64(time.Millisecond)
|
||||
to := q.TimeRange.To.UnixNano() / int64(time.Millisecond)
|
||||
if err := e.processQuery(q, ms, from, to); err != nil {
|
||||
mq, _ := json.Marshal(q)
|
||||
e.logger.Error("Failed to process query to multisearch request builder", "error", err, "query", string(mq), "queriesLength", len(queries), "duration", time.Since(start), "stage", es.StagePrepareRequest)
|
||||
@@ -88,7 +88,7 @@ func (e *elasticsearchDataQuery) processQuery(q *Query, ms *es.MultiSearchReques
|
||||
}
|
||||
|
||||
defaultTimeField := e.client.GetConfiguredFields().TimeField
|
||||
b := ms.Search(q.Interval)
|
||||
b := ms.Search(q.Interval, q.TimeRange)
|
||||
b.Size(0)
|
||||
filters := b.Query().Bool().Filter()
|
||||
filters.AddDateRangeFilter(defaultTimeField, to, from, es.DateFormatEpochMS)
|
||||
|
||||
@@ -65,7 +65,7 @@ func queryData(ctx context.Context, queries []backend.DataQuery, dsInfo *es.Data
|
||||
return &backend.QueryDataResponse{}, fmt.Errorf("query contains no queries")
|
||||
}
|
||||
|
||||
client, err := es.NewClient(ctx, dsInfo, queries[0].TimeRange, logger, tracer)
|
||||
client, err := es.NewClient(ctx, dsInfo, logger, tracer)
|
||||
if err != nil {
|
||||
return &backend.QueryDataResponse{}, err
|
||||
}
|
||||
|
||||
@@ -3,6 +3,7 @@ package elasticsearch
|
||||
import (
|
||||
"time"
|
||||
|
||||
"github.com/grafana/grafana-plugin-sdk-go/backend"
|
||||
"github.com/grafana/grafana/pkg/components/simplejson"
|
||||
)
|
||||
|
||||
@@ -16,6 +17,7 @@ type Query struct {
|
||||
IntervalMs int64
|
||||
RefID string
|
||||
MaxDataPoints int64
|
||||
TimeRange backend.TimeRange
|
||||
}
|
||||
|
||||
// BucketAgg represents a bucket aggregation of the time series query model of the datasource
|
||||
|
||||
@@ -42,6 +42,7 @@ func parseQuery(tsdbQuery []backend.DataQuery, logger log.Logger) ([]*Query, err
|
||||
IntervalMs: intervalMs,
|
||||
RefID: q.RefID,
|
||||
MaxDataPoints: q.MaxDataPoints,
|
||||
TimeRange: q.TimeRange,
|
||||
})
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user