Prometheus: Streaming JSON parser performance improvements (#48792)
This commit is contained in:
@@ -34,39 +34,46 @@ func TestMatrixResponses(t *testing.T) {
|
||||
}
|
||||
|
||||
for _, test := range tt {
|
||||
t.Run(test.name, func(t *testing.T) {
|
||||
queryFileName := filepath.Join("../testdata", test.filepath+".query.json")
|
||||
responseFileName := filepath.Join("../testdata", test.filepath+".result.json")
|
||||
goldenFileName := filepath.Join("../testdata", test.filepath+".result.streaming.golden")
|
||||
enableWideSeries := false
|
||||
queryFileName := filepath.Join("../testdata", test.filepath+".query.json")
|
||||
responseFileName := filepath.Join("../testdata", test.filepath+".result.json")
|
||||
goldenFileName := filepath.Join("../testdata", test.filepath+".result.streaming.golden")
|
||||
t.Run(test.name, goldenScenario(test.name, queryFileName, responseFileName, goldenFileName, enableWideSeries))
|
||||
enableWideSeries = true
|
||||
goldenFileName = filepath.Join("../testdata", test.filepath+".result.streaming-wide.golden")
|
||||
t.Run(test.name, goldenScenario(test.name, queryFileName, responseFileName, goldenFileName, enableWideSeries))
|
||||
}
|
||||
}
|
||||
|
||||
query, err := loadStoredQuery(queryFileName)
|
||||
func goldenScenario(name, queryFileName, responseFileName, goldenFileName string, wide bool) func(t *testing.T) {
|
||||
return func(t *testing.T) {
|
||||
query, err := loadStoredQuery(queryFileName)
|
||||
require.NoError(t, err)
|
||||
|
||||
responseBytes, err := os.ReadFile(responseFileName)
|
||||
require.NoError(t, err)
|
||||
|
||||
result, err := runQuery(responseBytes, query, wide)
|
||||
require.NoError(t, err)
|
||||
require.Len(t, result.Responses, 1)
|
||||
|
||||
dr, found := result.Responses["A"]
|
||||
require.True(t, found)
|
||||
|
||||
actual, err := json.MarshalIndent(&dr, "", " ")
|
||||
require.NoError(t, err)
|
||||
|
||||
// nolint:gosec
|
||||
// We can ignore the gosec G304 because this is a test with static defined paths
|
||||
expected, err := ioutil.ReadFile(goldenFileName + ".json")
|
||||
if err != nil || update {
|
||||
err = os.WriteFile(goldenFileName+".json", actual, 0600)
|
||||
require.NoError(t, err)
|
||||
}
|
||||
|
||||
responseBytes, err := os.ReadFile(responseFileName)
|
||||
require.NoError(t, err)
|
||||
require.JSONEq(t, string(expected), string(actual))
|
||||
|
||||
result, err := runQuery(responseBytes, query)
|
||||
require.NoError(t, err)
|
||||
require.Len(t, result.Responses, 1)
|
||||
|
||||
dr, found := result.Responses["A"]
|
||||
require.True(t, found)
|
||||
|
||||
actual, err := json.MarshalIndent(&dr, "", " ")
|
||||
require.NoError(t, err)
|
||||
|
||||
// nolint:gosec
|
||||
// We can ignore the gosec G304 because this is a test with static defined paths
|
||||
expected, err := ioutil.ReadFile(goldenFileName + ".json")
|
||||
if err != nil || update {
|
||||
err = os.WriteFile(goldenFileName+".json", actual, 0600)
|
||||
require.NoError(t, err)
|
||||
}
|
||||
|
||||
require.JSONEq(t, string(expected), string(actual))
|
||||
|
||||
require.NoError(t, experimental.CheckGoldenDataResponse(goldenFileName+".txt", &dr, update))
|
||||
})
|
||||
require.NoError(t, experimental.CheckGoldenDataResponse(goldenFileName+".txt", &dr, update))
|
||||
}
|
||||
}
|
||||
|
||||
@@ -123,8 +130,8 @@ func loadStoredQuery(fileName string) (*backend.QueryDataRequest, error) {
|
||||
}, nil
|
||||
}
|
||||
|
||||
func runQuery(response []byte, q *backend.QueryDataRequest) (*backend.QueryDataResponse, error) {
|
||||
tCtx := setup()
|
||||
func runQuery(response []byte, q *backend.QueryDataRequest, wide bool) (*backend.QueryDataResponse, error) {
|
||||
tCtx := setup(wide)
|
||||
res := &http.Response{
|
||||
StatusCode: 200,
|
||||
Body: ioutil.NopCloser(bytes.NewReader(response)),
|
||||
|
||||
@@ -22,7 +22,7 @@ import (
|
||||
// - go tool pprof -http=localhost:6061 memprofile.out
|
||||
func BenchmarkJson(b *testing.B) {
|
||||
body, q := createJsonTestData(1642000000, 1, 300, 400)
|
||||
tCtx := setup()
|
||||
tCtx := setup(true)
|
||||
b.ResetTimer()
|
||||
for n := 0; n < b.N; n++ {
|
||||
res := http.Response{
|
||||
|
||||
@@ -41,6 +41,7 @@ type QueryData struct {
|
||||
ID int64
|
||||
URL string
|
||||
TimeInterval string
|
||||
enableWideSeries bool
|
||||
}
|
||||
|
||||
func New(
|
||||
@@ -75,6 +76,7 @@ func New(
|
||||
TimeInterval: timeInterval,
|
||||
ID: settings.ID,
|
||||
URL: settings.URL,
|
||||
enableWideSeries: features.IsEnabled(featuremgmt.FlagPrometheusWideSeries),
|
||||
}, nil
|
||||
}
|
||||
|
||||
|
||||
@@ -58,7 +58,7 @@ func TestPrometheus_parseTimeSeriesResponse(t *testing.T) {
|
||||
},
|
||||
}
|
||||
|
||||
tctx := setup()
|
||||
tctx := setup(true)
|
||||
|
||||
qm := models.QueryModel{
|
||||
LegendFormat: "legend {{app}}",
|
||||
@@ -119,19 +119,17 @@ func TestPrometheus_parseTimeSeriesResponse(t *testing.T) {
|
||||
},
|
||||
JSON: b,
|
||||
}
|
||||
tctx := setup()
|
||||
tctx := setup(true)
|
||||
res, err := execute(tctx, query, result)
|
||||
require.NoError(t, err)
|
||||
|
||||
require.Len(t, res, 1)
|
||||
//require.Equal(t, "legend Application", res[0].Name)
|
||||
require.Len(t, res[0].Fields, 2)
|
||||
require.Len(t, res[0].Fields[0].Labels, 0)
|
||||
require.Equal(t, "Time", res[0].Fields[0].Name)
|
||||
require.Len(t, res[0].Fields[1].Labels, 2)
|
||||
require.Equal(t, "app=Application, tag2=tag2", res[0].Fields[1].Labels.String())
|
||||
require.Equal(t, "Value", res[0].Fields[1].Name)
|
||||
require.Equal(t, "legend Application", res[0].Fields[1].Config.DisplayNameFromDS)
|
||||
require.Equal(t, "legend Application", res[0].Fields[1].Name)
|
||||
|
||||
// Ensure the timestamps are UTC zoned
|
||||
testValue := res[0].Fields[0].At(0)
|
||||
@@ -167,7 +165,7 @@ func TestPrometheus_parseTimeSeriesResponse(t *testing.T) {
|
||||
},
|
||||
JSON: b,
|
||||
}
|
||||
tctx := setup()
|
||||
tctx := setup(true)
|
||||
res, err := execute(tctx, query, result)
|
||||
|
||||
require.NoError(t, err)
|
||||
@@ -176,8 +174,8 @@ func TestPrometheus_parseTimeSeriesResponse(t *testing.T) {
|
||||
require.Equal(t, time.Unix(1, 0).UTC(), res[0].Fields[0].At(0))
|
||||
require.Equal(t, time.Unix(4, 0).UTC(), res[0].Fields[0].At(1))
|
||||
require.Equal(t, res[0].Fields[1].Len(), 2)
|
||||
require.Equal(t, float64(1), res[0].Fields[1].At(0).(float64))
|
||||
require.Equal(t, float64(4), res[0].Fields[1].At(1).(float64))
|
||||
require.Equal(t, float64(1), *res[0].Fields[1].At(0).(*float64))
|
||||
require.Equal(t, float64(4), *res[0].Fields[1].At(1).(*float64))
|
||||
})
|
||||
|
||||
t.Run("matrix response with from alerting missed data points should be parsed correctly", func(t *testing.T) {
|
||||
@@ -209,19 +207,17 @@ func TestPrometheus_parseTimeSeriesResponse(t *testing.T) {
|
||||
},
|
||||
JSON: b,
|
||||
}
|
||||
tctx := setup()
|
||||
tctx := setup(true)
|
||||
res, err := execute(tctx, query, result)
|
||||
|
||||
require.NoError(t, err)
|
||||
require.Len(t, res, 1)
|
||||
require.Equal(t, res[0].Name, "{app=\"Application\", tag2=\"tag2\"}")
|
||||
require.Len(t, res[0].Fields, 2)
|
||||
require.Len(t, res[0].Fields[0].Labels, 0)
|
||||
require.Equal(t, res[0].Fields[0].Name, "Time")
|
||||
require.Len(t, res[0].Fields[1].Labels, 2)
|
||||
require.Equal(t, res[0].Fields[1].Labels.String(), "app=Application, tag2=tag2")
|
||||
require.Equal(t, res[0].Fields[1].Name, "Value")
|
||||
require.Equal(t, res[0].Fields[1].Config.DisplayNameFromDS, "{app=\"Application\", tag2=\"tag2\"}")
|
||||
require.Equal(t, "{app=\"Application\", tag2=\"tag2\"}", res[0].Fields[1].Name)
|
||||
})
|
||||
|
||||
t.Run("matrix response with NaN value should be changed to null", func(t *testing.T) {
|
||||
@@ -252,12 +248,12 @@ func TestPrometheus_parseTimeSeriesResponse(t *testing.T) {
|
||||
JSON: b,
|
||||
}
|
||||
|
||||
tctx := setup()
|
||||
tctx := setup(true)
|
||||
res, err := execute(tctx, query, result)
|
||||
require.NoError(t, err)
|
||||
|
||||
require.Equal(t, res[0].Fields[1].Name, "Value")
|
||||
require.True(t, math.IsNaN(res[0].Fields[1].At(0).(float64)))
|
||||
require.Equal(t, "{app=\"Application\"}", res[0].Fields[1].Name)
|
||||
require.True(t, math.IsNaN(*res[0].Fields[1].At(0).(*float64)))
|
||||
})
|
||||
|
||||
t.Run("vector response should be parsed normally", func(t *testing.T) {
|
||||
@@ -281,20 +277,18 @@ func TestPrometheus_parseTimeSeriesResponse(t *testing.T) {
|
||||
query := backend.DataQuery{
|
||||
JSON: b,
|
||||
}
|
||||
tctx := setup()
|
||||
tctx := setup(true)
|
||||
res, err := execute(tctx, query, qr)
|
||||
require.NoError(t, err)
|
||||
|
||||
require.Len(t, res, 1)
|
||||
require.Equal(t, res[0].Name, "legend Application")
|
||||
require.Len(t, res[0].Fields, 2)
|
||||
require.Len(t, res[0].Fields[0].Labels, 0)
|
||||
require.Equal(t, res[0].Fields[0].Name, "Time")
|
||||
require.Equal(t, res[0].Fields[0].Name, "Time")
|
||||
require.Len(t, res[0].Fields[1].Labels, 2)
|
||||
require.Equal(t, res[0].Fields[1].Labels.String(), "app=Application, tag2=tag2")
|
||||
require.Equal(t, res[0].Fields[1].Name, "Value")
|
||||
require.Equal(t, res[0].Fields[1].Config.DisplayNameFromDS, "legend Application")
|
||||
require.Equal(t, "legend Application", res[0].Fields[1].Name)
|
||||
|
||||
// Ensure the timestamps are UTC zoned
|
||||
testValue := res[0].Fields[0].At(0)
|
||||
@@ -321,17 +315,15 @@ func TestPrometheus_parseTimeSeriesResponse(t *testing.T) {
|
||||
query := backend.DataQuery{
|
||||
JSON: b,
|
||||
}
|
||||
tctx := setup()
|
||||
tctx := setup(true)
|
||||
res, err := execute(tctx, query, qr)
|
||||
require.NoError(t, err)
|
||||
|
||||
require.Len(t, res, 1)
|
||||
require.Equal(t, res[0].Name, "1")
|
||||
require.Len(t, res[0].Fields, 2)
|
||||
require.Len(t, res[0].Fields[0].Labels, 0)
|
||||
require.Equal(t, res[0].Fields[0].Name, "Time")
|
||||
require.Equal(t, res[0].Fields[1].Name, "Value")
|
||||
require.Equal(t, res[0].Fields[1].Config.DisplayNameFromDS, "1")
|
||||
require.Equal(t, "1", res[0].Fields[1].Name)
|
||||
|
||||
// Ensure the timestamps are UTC zoned
|
||||
testValue := res[0].Fields[0].At(0)
|
||||
@@ -397,7 +389,7 @@ type testContext struct {
|
||||
queryData *querydata.QueryData
|
||||
}
|
||||
|
||||
func setup() *testContext {
|
||||
func setup(wideFrames bool) *testContext {
|
||||
tracer, err := tracing.InitializeTracerForTest()
|
||||
if err != nil {
|
||||
panic(err)
|
||||
@@ -414,7 +406,7 @@ func setup() *testContext {
|
||||
queryData, _ := querydata.New(
|
||||
httpProvider,
|
||||
setting.NewCfg(),
|
||||
&fakeFeatureToggles{enabled: true},
|
||||
&fakeFeatureToggles{flags: map[string]bool{"prometheusStreamingJSONParser": true, "prometheusWideSeries": wideFrames}},
|
||||
tracer,
|
||||
backend.DataSourceInstanceSettings{URL: "http://localhost:9090", JSONData: json.RawMessage(`{"timeInterval": "15s"}`)},
|
||||
&fakeLogger{},
|
||||
@@ -427,11 +419,11 @@ func setup() *testContext {
|
||||
}
|
||||
|
||||
type fakeFeatureToggles struct {
|
||||
enabled bool
|
||||
flags map[string]bool
|
||||
}
|
||||
|
||||
func (f *fakeFeatureToggles) IsEnabled(feature string) bool {
|
||||
return f.enabled
|
||||
return f.flags[feature]
|
||||
}
|
||||
|
||||
type fakeHttpClientProvider struct {
|
||||
|
||||
@@ -22,20 +22,27 @@ func (s *QueryData) parseResponse(ctx context.Context, q *models.Query, res *htt
|
||||
}()
|
||||
|
||||
iter := jsoniter.Parse(jsoniter.ConfigDefault, res.Body, 1024)
|
||||
r := converter.ReadPrometheusStyleResult(iter)
|
||||
r := converter.ReadPrometheusStyleResult(iter, converter.Options{
|
||||
MatrixWideSeries: s.enableWideSeries,
|
||||
VectorWideSeries: s.enableWideSeries,
|
||||
})
|
||||
if r == nil {
|
||||
return nil, fmt.Errorf("received empty response from prometheus")
|
||||
}
|
||||
|
||||
// The ExecutedQueryString can be viewed in QueryInspector in UI
|
||||
for _, frame := range r.Frames {
|
||||
addMetadataToFrame(q, frame)
|
||||
if s.enableWideSeries {
|
||||
addMetadataToWideFrame(q, frame)
|
||||
} else {
|
||||
addMetadataToMultiFrame(q, frame)
|
||||
}
|
||||
}
|
||||
|
||||
return r, nil
|
||||
}
|
||||
|
||||
func addMetadataToFrame(q *models.Query, frame *data.Frame) {
|
||||
func addMetadataToMultiFrame(q *models.Query, frame *data.Frame) {
|
||||
if frame.Meta == nil {
|
||||
frame.Meta = &data.FrameMeta{}
|
||||
}
|
||||
@@ -43,16 +50,32 @@ func addMetadataToFrame(q *models.Query, frame *data.Frame) {
|
||||
if len(frame.Fields) < 2 {
|
||||
return
|
||||
}
|
||||
frame.Name = getName(q, frame)
|
||||
frame.Name = getName(q, frame.Fields[1])
|
||||
frame.Fields[0].Config = &data.FieldConfig{Interval: float64(q.Step.Milliseconds())}
|
||||
if frame.Name != "" {
|
||||
frame.Fields[1].Config = &data.FieldConfig{DisplayNameFromDS: frame.Name}
|
||||
}
|
||||
}
|
||||
|
||||
func addMetadataToWideFrame(q *models.Query, frame *data.Frame) {
|
||||
if frame.Meta == nil {
|
||||
frame.Meta = &data.FrameMeta{}
|
||||
}
|
||||
frame.Meta.ExecutedQueryString = executedQueryString(q)
|
||||
if len(frame.Fields) < 2 {
|
||||
return
|
||||
}
|
||||
frame.Fields[0].Config = &data.FieldConfig{Interval: float64(q.Step.Milliseconds())}
|
||||
for _, f := range frame.Fields {
|
||||
if f.Name != data.TimeSeriesTimeFieldName {
|
||||
f.Name = getName(q, f)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// this is based on the logic from the String() function in github.com/prometheus/common/model.go
|
||||
func metricNameFromLabels(f *data.Frame) string {
|
||||
labels := f.Fields[1].Labels
|
||||
func metricNameFromLabels(f *data.Field) string {
|
||||
labels := f.Labels
|
||||
metricName, hasName := labels["__name__"]
|
||||
numLabels := len(labels) - 1
|
||||
if !hasName {
|
||||
@@ -81,9 +104,9 @@ func executedQueryString(q *models.Query) string {
|
||||
return "Expr: " + q.Expr + "\n" + "Step: " + q.Step.String()
|
||||
}
|
||||
|
||||
func getName(q *models.Query, frame *data.Frame) string {
|
||||
labels := frame.Fields[1].Labels
|
||||
legend := metricNameFromLabels(frame)
|
||||
func getName(q *models.Query, field *data.Field) string {
|
||||
labels := field.Labels
|
||||
legend := metricNameFromLabels(field)
|
||||
|
||||
if q.LegendFormat == legendFormatAuto && len(labels) > 0 {
|
||||
return ""
|
||||
|
||||
Reference in New Issue
Block a user