diff --git a/pkg/registry/apis/query/query.go b/pkg/registry/apis/query/query.go index 987251bcdb0..99323eb2721 100644 --- a/pkg/registry/apis/query/query.go +++ b/pkg/registry/apis/query/query.go @@ -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 { diff --git a/pkg/registry/apis/query/register.go b/pkg/registry/apis/query/register.go index 221b589b511..e7ab7fb9b07 100644 --- a/pkg/registry/apis/query/register.go +++ b/pkg/registry/apis/query/register.go @@ -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) diff --git a/pkg/services/query/query.go b/pkg/services/query/query.go index 34abc9a2dbf..dd646fa2acc 100644 --- a/pkg/services/query/query.go +++ b/pkg/services/query/query.go @@ -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)