datasources: querier: configurable concurrent-query-limit (#114585)
This commit is contained in:
@@ -338,8 +338,8 @@ func prepareQuery(
|
||||
}, nil
|
||||
}
|
||||
|
||||
func handlePreparedQuery(ctx context.Context, pq *preparedQuery) (*backend.QueryDataResponse, error) {
|
||||
resp, err := service.QueryData(ctx, pq.logger, pq.cache, pq.exprSvc, pq.mReq, pq.builder, pq.headers)
|
||||
func handlePreparedQuery(ctx context.Context, pq *preparedQuery, concurrentQueryLimit int) (*backend.QueryDataResponse, error) {
|
||||
resp, err := service.QueryData(ctx, pq.logger, pq.cache, pq.exprSvc, pq.mReq, pq.builder, pq.headers, concurrentQueryLimit)
|
||||
pq.reportMetrics()
|
||||
return resp, err
|
||||
}
|
||||
@@ -357,7 +357,7 @@ func handleQuery(
|
||||
responder.Error(err)
|
||||
return nil, err
|
||||
}
|
||||
return handlePreparedQuery(ctx, pq)
|
||||
return handlePreparedQuery(ctx, pq, b.concurrentQueryLimit)
|
||||
}
|
||||
|
||||
type responderWrapper struct {
|
||||
|
||||
@@ -3,10 +3,11 @@ package query
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"runtime"
|
||||
|
||||
"github.com/prometheus/client_golang/prometheus"
|
||||
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||||
"k8s.io/apimachinery/pkg/runtime"
|
||||
apiruntime "k8s.io/apimachinery/pkg/runtime"
|
||||
"k8s.io/apimachinery/pkg/runtime/schema"
|
||||
"k8s.io/apiserver/pkg/authorization/authorizer"
|
||||
"k8s.io/apiserver/pkg/registry/rest"
|
||||
@@ -62,6 +63,7 @@ func NewQueryAPIBuilder(
|
||||
tracer tracing.Tracer,
|
||||
legacyDatasourceLookup service.LegacyDataSourceLookup,
|
||||
connections DataSourceConnectionProvider,
|
||||
concurrentQueryLimit int,
|
||||
) (*QueryAPIBuilder, error) {
|
||||
// Include well typed query definitions
|
||||
var queryTypes *query.QueryTypeDefinitionList
|
||||
@@ -80,7 +82,7 @@ func NewQueryAPIBuilder(
|
||||
}
|
||||
|
||||
return &QueryAPIBuilder{
|
||||
concurrentQueryLimit: 4,
|
||||
concurrentQueryLimit: concurrentQueryLimit,
|
||||
log: log.New("query_apiserver"),
|
||||
instanceProvider: instanceProvider,
|
||||
authorizer: ar,
|
||||
@@ -142,6 +144,7 @@ func RegisterAPIService(
|
||||
tracer,
|
||||
legacyDatasourceLookup,
|
||||
&connectionsProvider{dsService: dataSourcesService, registry: reg},
|
||||
cfg.SectionWithEnvOverrides("query").Key("concurrent_query_limit").MustInt(runtime.NumCPU()),
|
||||
)
|
||||
apiregistration.RegisterAPI(builder)
|
||||
return builder, err
|
||||
@@ -151,7 +154,7 @@ func (b *QueryAPIBuilder) GetGroupVersion() schema.GroupVersion {
|
||||
return query.SchemeGroupVersion
|
||||
}
|
||||
|
||||
func addKnownTypes(scheme *runtime.Scheme, gv schema.GroupVersion) {
|
||||
func addKnownTypes(scheme *apiruntime.Scheme, gv schema.GroupVersion) {
|
||||
scheme.AddKnownTypes(gv,
|
||||
&query.DataSourceApiServer{},
|
||||
&query.DataSourceApiServerList{},
|
||||
@@ -165,7 +168,7 @@ func addKnownTypes(scheme *runtime.Scheme, gv schema.GroupVersion) {
|
||||
)
|
||||
}
|
||||
|
||||
func (b *QueryAPIBuilder) InstallSchema(scheme *runtime.Scheme) error {
|
||||
func (b *QueryAPIBuilder) InstallSchema(scheme *apiruntime.Scheme) error {
|
||||
addKnownTypes(scheme, query.SchemeGroupVersion)
|
||||
metav1.AddToGroupVersion(scheme, query.SchemeGroupVersion)
|
||||
return scheme.SetVersionPriority(query.SchemeGroupVersion)
|
||||
|
||||
@@ -226,7 +226,7 @@ func buildErrorResponses(err error, queries []*simplejson.Json) splitResponse {
|
||||
return splitResponse{er, http.Header{}}
|
||||
}
|
||||
|
||||
func QueryData(ctx context.Context, log log.Logger, dscache datasources.CacheService, exprService *expr.Service, reqDTO dtos.MetricRequest, qsDatasourceClientBuilder dsquerierclient.QSDatasourceClientBuilder, headers map[string]string) (*backend.QueryDataResponse, error) {
|
||||
func QueryData(ctx context.Context, log log.Logger, dscache datasources.CacheService, exprService *expr.Service, reqDTO dtos.MetricRequest, qsDatasourceClientBuilder dsquerierclient.QSDatasourceClientBuilder, headers map[string]string, concurrentQueryLimit int) (*backend.QueryDataResponse, error) {
|
||||
s := &ServiceImpl{
|
||||
log: log,
|
||||
dataSourceCache: dscache,
|
||||
@@ -234,7 +234,7 @@ func QueryData(ctx context.Context, log log.Logger, dscache datasources.CacheSer
|
||||
dataSourceRequestValidator: validations.ProvideValidator(),
|
||||
qsDatasourceClientBuilder: qsDatasourceClientBuilder,
|
||||
headers: headers,
|
||||
concurrentQueryLimit: 16, // TODO: make it configurable
|
||||
concurrentQueryLimit: concurrentQueryLimit,
|
||||
}
|
||||
|
||||
user, err := identity.GetRequester(ctx)
|
||||
|
||||
Reference in New Issue
Block a user