Tempo: Fix instant query streaming (#108924)
* Don't use streaming for instant queries * wip * Only return latest instant query result * Always disable streaming for alerting queries * lint
This commit is contained in:
@@ -20,6 +20,10 @@ import (
|
||||
|
||||
const MetricsPathPrefix = "metrics/"
|
||||
|
||||
type PartialTempoQuery struct {
|
||||
MetricsQueryType *dataquery.MetricsQueryType
|
||||
}
|
||||
|
||||
func (s *Service) runMetricsStream(ctx context.Context, req *backend.RunStreamRequest, sender *backend.StreamSender, datasource *Datasource) error {
|
||||
ctx, span := tracing.DefaultTracer().Start(ctx, "datasource.tempo.runMetricsStream")
|
||||
defer span.End()
|
||||
@@ -35,6 +39,15 @@ func (s *Service) runMetricsStream(ctx context.Context, req *backend.RunStreamRe
|
||||
return err
|
||||
}
|
||||
|
||||
tempoQuery := &PartialTempoQuery{}
|
||||
err = json.Unmarshal(req.Data, tempoQuery)
|
||||
if err != nil {
|
||||
response.Error = fmt.Errorf("error unmarshaling Tempo query model: %v", err)
|
||||
span.RecordError(response.Error)
|
||||
span.SetStatus(codes.Error, response.Error.Error())
|
||||
return err
|
||||
}
|
||||
|
||||
var qrr *tempopb.QueryRangeRequest
|
||||
err = json.Unmarshal(req.Data, &qrr)
|
||||
if err != nil {
|
||||
@@ -56,6 +69,24 @@ func (s *Service) runMetricsStream(ctx context.Context, req *backend.RunStreamRe
|
||||
// Ideally this would be pushed higher, so it's set once for all rpc calls, but we have only one now.
|
||||
ctx = metadata.AppendToOutgoingContext(ctx, "User-Agent", backend.UserAgentFromContext(ctx).String())
|
||||
|
||||
if isInstantQuery(tempoQuery.MetricsQueryType) {
|
||||
instantQuery := &tempopb.QueryInstantRequest{
|
||||
Query: qrr.Query,
|
||||
Start: qrr.Start,
|
||||
End: qrr.End,
|
||||
}
|
||||
|
||||
stream, err := datasource.StreamingClient.MetricsQueryInstant(ctx, instantQuery)
|
||||
if err != nil {
|
||||
span.RecordError(err)
|
||||
span.SetStatus(codes.Error, err.Error())
|
||||
s.logger.Error("Error Search()", "err", err)
|
||||
return err
|
||||
}
|
||||
|
||||
return s.processInstantMetricsStream(ctx, stream, sender)
|
||||
}
|
||||
|
||||
stream, err := datasource.StreamingClient.MetricsQueryRange(ctx, qrr)
|
||||
if err != nil {
|
||||
span.RecordError(err)
|
||||
@@ -101,3 +132,38 @@ func (s *Service) processMetricsStream(ctx context.Context, query string, stream
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *Service) processInstantMetricsStream(ctx context.Context, stream tempopb.StreamingQuerier_MetricsQueryInstantClient, sender StreamSender) error {
|
||||
ctx, span := tracing.DefaultTracer().Start(ctx, "datasource.tempo.processStream")
|
||||
defer span.End()
|
||||
messageCount := 0
|
||||
for {
|
||||
msg, err := stream.Recv()
|
||||
messageCount++
|
||||
span.SetAttributes(attribute.Int("message_count", messageCount))
|
||||
if errors.Is(err, io.EOF) {
|
||||
if err := s.sendResponse(ctx, nil, nil, dataquery.SearchStreamingStateDone, sender); err != nil {
|
||||
span.RecordError(err)
|
||||
span.SetStatus(codes.Error, err.Error())
|
||||
return err
|
||||
}
|
||||
break
|
||||
}
|
||||
if err != nil {
|
||||
s.logger.Error("Error receiving message", "err", err)
|
||||
span.RecordError(err)
|
||||
span.SetStatus(codes.Error, err.Error())
|
||||
return err
|
||||
}
|
||||
|
||||
transformed := traceql.TransformInstantMetricsResponse(*msg)
|
||||
|
||||
if err := s.sendResponse(ctx, transformed, msg.Metrics, dataquery.SearchStreamingStateStreaming, sender); err != nil {
|
||||
span.RecordError(err)
|
||||
span.SetStatus(codes.Error, err.Error())
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -9,7 +9,6 @@ import (
|
||||
"time"
|
||||
|
||||
"github.com/grafana/grafana-plugin-sdk-go/data"
|
||||
"github.com/grafana/grafana/pkg/tsdb/tempo/kinds/dataquery"
|
||||
"github.com/grafana/tempo/pkg/tempopb"
|
||||
v1 "github.com/grafana/tempo/pkg/tempopb/common/v1"
|
||||
)
|
||||
@@ -61,7 +60,7 @@ func TransformMetricsResponse(query string, resp tempopb.QueryRangeResponse) []*
|
||||
return append(frames, exemplarFrames...)
|
||||
}
|
||||
|
||||
func TransformInstantMetricsResponse(query *dataquery.TempoQuery, resp tempopb.QueryInstantResponse) []*data.Frame {
|
||||
func TransformInstantMetricsResponse(resp tempopb.QueryInstantResponse) []*data.Frame {
|
||||
frames := make([]*data.Frame, len(resp.Series))
|
||||
|
||||
for i, series := range resp.Series {
|
||||
|
||||
@@ -5,7 +5,6 @@ import (
|
||||
"time"
|
||||
|
||||
"github.com/grafana/grafana-plugin-sdk-go/data"
|
||||
"github.com/grafana/grafana/pkg/tsdb/tempo/kinds/dataquery"
|
||||
"github.com/grafana/tempo/pkg/tempopb"
|
||||
v1 "github.com/grafana/tempo/pkg/tempopb/common/v1"
|
||||
"github.com/stretchr/testify/assert"
|
||||
@@ -113,7 +112,6 @@ func TestTransformMetricsResponse_MultipleSeries(t *testing.T) {
|
||||
}
|
||||
|
||||
func TestTransformInstantMetricsResponse(t *testing.T) {
|
||||
query := &dataquery.TempoQuery{}
|
||||
resp := tempopb.QueryInstantResponse{
|
||||
Series: []*tempopb.InstantSeries{
|
||||
{
|
||||
@@ -123,7 +121,7 @@ func TestTransformInstantMetricsResponse(t *testing.T) {
|
||||
},
|
||||
}
|
||||
|
||||
frames := TransformInstantMetricsResponse(query, resp)
|
||||
frames := TransformInstantMetricsResponse(resp)
|
||||
|
||||
assert.Len(t, frames, 1)
|
||||
frame := frames[0]
|
||||
|
||||
@@ -97,7 +97,7 @@ func (s *Service) runTraceQlQueryMetrics(ctx context.Context, pCtx backend.Plugi
|
||||
return res, err
|
||||
}
|
||||
|
||||
frames := traceql.TransformInstantMetricsResponse(tempoQuery, queryResponse)
|
||||
frames := traceql.TransformInstantMetricsResponse(queryResponse)
|
||||
result.Frames = frames
|
||||
} else {
|
||||
var queryResponse tempopb.QueryRangeResponse
|
||||
|
||||
Reference in New Issue
Block a user