Files
grafana/pkg/tsdb/influxdb/healthcheck.go
T
ismail simsekandİnanç Gümüş c088d003f2 InfluxDB: Implement InfluxQL json streaming parser (#76934)
* Have the first iteration

* Prepare bench testing

* rename the test files

* Remove unnecessary test file

* Introduce influxqlStreamingParser feature flag

* Apply streaming parser feature flag

* Add new tests

* More tests

* return executedQueryString only in first frame

* add frame meta and config

* Update golden json files

* Support tags/labels

* more tests

* more tests

* Don't change original response_parser.go

* provide context

* create util package

* don't pass the row

* update converter with formatted frameName

* add executedQueryString info only to first frame

* update golden files

* rename

* update test file

* use pointer values

* update testdata

* update parsing

* update converter for null values

* prepare converter for table response

* clean up

* return timeField in fields

* handle no time column responses

* better nil field handling

* refactor the code

* add table tests

* fix config for table

* table response format

* fix value

* if there is no time column set name

* linting

* refactoring

* handle the status code

* add tracing

* Update pkg/tsdb/influxdb/influxql/converter/converter_test.go

Co-authored-by: İnanç Gümüş <m@inanc.io>

* fix import

* update test data

* sanity

* sanity

* linting

* simplicity

* return empty rsp

* rename to prevent confusion

* nullableJson field type for null values

* better handling null values

* remove duplicate test file

* fix healthcheck

* use util for pointer

* move bench test to root

* provide fake feature manager

* add more tests

* partial fix for null values in table response format

* handle partial null fields

* comments for easy testing

* move frameName allocation in readSeries

* one less append operation

* performance improvement by making string conversion once

pkg: github.com/grafana/grafana/pkg/tsdb/influxdb/influxql
             │ stream2.txt │            stream3.txt             │
             │   sec/op    │   sec/op     vs base               │
ParseJson-10   314.4m ± 1%   303.9m ± 1%  -3.34% (p=0.000 n=10)

             │ stream2.txt  │             stream3.txt              │
             │     B/op     │     B/op      vs base                │
ParseJson-10   425.2Mi ± 0%   382.7Mi ± 0%  -10.00% (p=0.000 n=10)

             │ stream2.txt │            stream3.txt             │
             │  allocs/op  │  allocs/op   vs base               │
ParseJson-10   7.224M ± 0%   6.689M ± 0%  -7.41% (p=0.000 n=10)

* add comment lines

---------

Co-authored-by: İnanç Gümüş <m@inanc.io>
2023-12-06 12:39:05 +01:00

167 lines
4.8 KiB
Go

package influxdb
import (
"context"
"errors"
"fmt"
"time"
"github.com/grafana/grafana-plugin-sdk-go/backend"
"github.com/grafana/grafana-plugin-sdk-go/backend/tracing"
"github.com/grafana/grafana/pkg/infra/log"
"github.com/grafana/grafana/pkg/services/featuremgmt"
"github.com/grafana/grafana/pkg/tsdb/influxdb/flux"
"github.com/grafana/grafana/pkg/tsdb/influxdb/fsql"
"github.com/grafana/grafana/pkg/tsdb/influxdb/influxql"
"github.com/grafana/grafana/pkg/tsdb/influxdb/models"
)
const (
refID = "healthcheck"
)
func (s *Service) CheckHealth(ctx context.Context, req *backend.CheckHealthRequest) (*backend.CheckHealthResult,
error) {
logger := logger.FromContext(ctx)
dsInfo, err := s.getDSInfo(ctx, req.PluginContext)
if err != nil {
return getHealthCheckMessage(logger, "error getting datasource info", err)
}
if dsInfo == nil {
return getHealthCheckMessage(logger, "", errors.New("invalid datasource info received"))
}
switch dsInfo.Version {
case influxVersionFlux:
return CheckFluxHealth(ctx, dsInfo, req)
case influxVersionInfluxQL:
return CheckInfluxQLHealth(ctx, dsInfo, s.features)
case influxVersionSQL:
return CheckSQLHealth(ctx, dsInfo, req)
default:
return getHealthCheckMessage(logger, "", errors.New("unknown influx version"))
}
}
func CheckFluxHealth(ctx context.Context, dsInfo *models.DatasourceInfo,
req *backend.CheckHealthRequest) (*backend.CheckHealthResult,
error) {
logger := logger.FromContext(ctx)
ds, err := flux.Query(ctx, dsInfo, backend.QueryDataRequest{
PluginContext: req.PluginContext,
Queries: []backend.DataQuery{
{
RefID: refID,
JSON: []byte(`{ "query": "buckets()" }`),
Interval: 1 * time.Minute,
MaxDataPoints: 423,
TimeRange: backend.TimeRange{
From: time.Now().AddDate(0, 0, -1),
To: time.Now(),
},
},
},
})
if err != nil {
return getHealthCheckMessage(logger, "error performing flux query", err)
}
if res, ok := ds.Responses[refID]; ok {
if res.Error != nil {
return getHealthCheckMessage(logger, "error reading buckets", res.Error)
}
if len(res.Frames) > 0 && len(res.Frames[0].Fields) > 0 {
return getHealthCheckMessage(logger, fmt.Sprintf("%d buckets found", res.Frames[0].Fields[0].Len()), nil)
}
}
return getHealthCheckMessage(logger, "", errors.New("error getting flux query buckets"))
}
func CheckInfluxQLHealth(ctx context.Context, dsInfo *models.DatasourceInfo, features featuremgmt.FeatureToggles) (*backend.CheckHealthResult, error) {
logger := logger.FromContext(ctx)
tracer := tracing.DefaultTracer()
resp, err := influxql.Query(ctx, tracer, dsInfo, &backend.QueryDataRequest{
Queries: []backend.DataQuery{
{
RefID: refID,
QueryType: "health",
JSON: []byte(`{"query": "SHOW measurements", "rawQuery": true}`),
},
},
}, features)
if err != nil {
return getHealthCheckMessage(logger, "error performing influxQL query", err)
}
if res, ok := resp.Responses[refID]; ok {
if res.Error != nil {
return getHealthCheckMessage(logger, "error reading influxDB", res.Error)
}
if len(res.Frames) == 0 {
return getHealthCheckMessage(logger, "0 measurements found", nil)
}
if len(res.Frames) > 0 && len(res.Frames[0].Fields) > 0 {
return getHealthCheckMessage(logger, fmt.Sprintf("%d measurements found", res.Frames[0].Fields[0].Len()), nil)
}
}
return getHealthCheckMessage(logger, "", errors.New("error connecting influxDB influxQL"))
}
func CheckSQLHealth(ctx context.Context, dsInfo *models.DatasourceInfo, req *backend.CheckHealthRequest) (*backend.CheckHealthResult, error) {
ds, err := fsql.Query(ctx, dsInfo, backend.QueryDataRequest{
PluginContext: req.PluginContext,
Queries: []backend.DataQuery{
{
RefID: refID,
JSON: []byte(`{ "rawSql": "select 1", "format": "table" }`),
Interval: 1 * time.Minute,
MaxDataPoints: 423,
TimeRange: backend.TimeRange{
From: time.Now().AddDate(0, 0, -1),
To: time.Now(),
},
},
},
})
if err != nil {
return getHealthCheckMessage(logger, "error performing sql query", err)
}
res := ds.Responses[refID]
if res.Error != nil {
return &backend.CheckHealthResult{
Status: backend.HealthStatusError,
Message: fmt.Sprintf("ERROR: %s", res.Error),
}, nil
}
return &backend.CheckHealthResult{
Status: backend.HealthStatusOk,
Message: "OK",
}, nil
}
func getHealthCheckMessage(logger log.Logger, message string, err error) (*backend.CheckHealthResult, error) {
if err == nil {
return &backend.CheckHealthResult{
Status: backend.HealthStatusOk,
Message: fmt.Sprintf("datasource is working. %s", message),
}, nil
}
logger.Warn("Error performing influxdb healthcheck", "err", err.Error())
errorMessage := fmt.Sprintf("%s %s", err.Error(), message)
return &backend.CheckHealthResult{
Status: backend.HealthStatusError,
Message: errorMessage,
}, nil
}