Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions cmd/pyroscope/help-all.txt.tmpl
Original file line number Diff line number Diff line change
Expand Up @@ -935,6 +935,8 @@ Usage of ./pyroscope:
Port to advertise to query-scheduler and querier (defaults to -server.http-listen-port).
-query-frontend.max-async-query-concurrency int
Maximum number of concurrent async queries per tenant. 0 to disable async queries. (default 5)
-query-frontend.query-planner-strategy string
Sets the query planner strategy, options: classic, balanced (default "classic")
-query-frontend.scheduler-worker-concurrency int
Number of concurrent workers forwarding queries to single query-scheduler. (default 5)
-query-scheduler.grpc-client-config.backoff-max-period duration
Expand Down
2 changes: 2 additions & 0 deletions cmd/pyroscope/help.txt.tmpl
Original file line number Diff line number Diff line change
Expand Up @@ -199,6 +199,8 @@ Usage of ./pyroscope:
Include profiles that were sampled out and stored with stacktraces stripped (marked __sampled__) in query results.
-query-frontend.max-async-query-concurrency int
Maximum number of concurrent async queries per tenant. 0 to disable async queries. (default 5)
-query-frontend.query-planner-strategy string
Sets the query planner strategy, options: classic, balanced (default "classic")
-query-scheduler.max-outstanding-requests-per-tenant int
Maximum number of outstanding requests per tenant per query-scheduler. In-flight requests above this limit will fail with HTTP response status code 429. (default 100)
-query-scheduler.ring.consul.hostname string
Expand Down
16 changes: 16 additions & 0 deletions pkg/frontend/frontend.go
Original file line number Diff line number Diff line change
Expand Up @@ -54,6 +54,13 @@ type Config struct {
// on the request is rejected with Unimplemented.
AsyncQueriesEnabled bool `yaml:"async_queries_enabled" category:"experimental"`

// QueryPlannerStrategy sets the query planner strategy. By default this is
// "classic" which provides the legacy query planner behavior.
//
// Optionally this can be "balanced" which uses the balanced query planner
// algorithm.
QueryPlannerStrategy string `yaml:"query_planner_strategy" doc:"hidden"`

// Used to find local IP address, that is sent to scheduler and querier-worker.
InfNames []string `yaml:"instance_interface_names" category:"advanced" doc:"default=[<private network interfaces>]"`
Addr string `yaml:"instance_addr" category:"advanced"`
Expand All @@ -79,6 +86,7 @@ func (cfg *Config) RegisterFlags(f *flag.FlagSet, logger log.Logger) {
f.BoolVar(&cfg.EnableIPv6, "query-frontend.instance-enable-ipv6", false, "Enable using a IPv6 instance address. (default false)")
f.IntVar(&cfg.Port, "query-frontend.instance-port", 0, "Port to advertise to query-scheduler and querier (defaults to -server.http-listen-port).")
f.BoolVar(&cfg.AsyncQueriesEnabled, "query-frontend.async-queries-enabled", false, "Enable the experimental asynchronous query path on SelectMergeStacktraces (default false)")
f.StringVar(&cfg.QueryPlannerStrategy, "query-frontend.query-planner-strategy", "classic", "Sets the query planner strategy, options: classic, balanced")
cfg.GRPCClientConfig.RegisterFlagsWithPrefix("query-frontend.grpc-client-config", f)
}

Expand All @@ -87,6 +95,14 @@ func (cfg *Config) Validate() error {
return fmt.Errorf("scheduler address cannot be specified when query-scheduler service discovery mode is set to '%s'", cfg.QuerySchedulerDiscovery.Mode)
}

switch cfg.QueryPlannerStrategy {
case "":
cfg.QueryPlannerStrategy = "classic"
case "classic", "balanced":
default:
return fmt.Errorf("unknown query planner strategy: %q", cfg.QueryPlannerStrategy)
}

return cfg.GRPCClientConfig.Validate()
}

Expand Down
15 changes: 12 additions & 3 deletions pkg/frontend/readpath/queryfrontend/query_frontend.go
Original file line number Diff line number Diff line change
Expand Up @@ -66,6 +66,7 @@ type QueryFrontend struct {
symbolizer Symbolizer
diagnosticsStore DiagnosticsStore
now func() time.Time
queryPlanType string

metrics *queryFrontendMetrics
}
Expand Down Expand Up @@ -132,6 +133,7 @@ func newQueryFrontendMetrics(reg prometheus.Registerer) *queryFrontendMetrics {
func NewQueryFrontend(
logger log.Logger,
limits frontend.Limits,
cfg frontend.Config,
metadataQueryClient metastorev1.MetadataQueryServiceClient,
tenantServiceClient metastorev1.TenantServiceClient,
querybackendClient QueryBackend,
Expand All @@ -148,6 +150,7 @@ func NewQueryFrontend(
symbolizer: sym,
diagnosticsStore: diagnosticsStore,
now: time.Now,
queryPlanType: cfg.QueryPlannerStrategy,
metrics: newQueryFrontendMetrics(reg),
}
return qf
Expand Down Expand Up @@ -228,7 +231,7 @@ func (q *QueryFrontend) doQuery(
endTime := time.UnixMilli(req.EndTime)
queryWindow := endTime.Sub(startTime).Round(time.Second)
traceID, _ := tracing.ExtractTraceID(ctx)
logArgs := []interface{}{
logArgs := []any{
"msg", "query weight",
"trace_id", traceID,
"tenant", strings.Join(tenants, ","),
Expand All @@ -255,8 +258,14 @@ func (q *QueryFrontend) doQuery(
blocks[i], blocks[j] = blocks[j], blocks[i]
})
xrandMutex.Unlock()
// TODO(kolesnikovae): Should be dynamic.
p := queryplan.Build(blocks, 4, 20)

var p *queryv1.QueryPlan
switch q.queryPlanType {
case "balanced":
p = queryplan.BuildBalanced(blocks, 4, 20)
default: // Always fallback to the classic planner
p = queryplan.Build(blocks, 4, 20)
}

backend := q.querybackend
if backendC != nil {
Expand Down
3 changes: 3 additions & 0 deletions pkg/frontend/readpath/queryfrontend/query_frontend_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@ import (
typesv1 "github.com/grafana/pyroscope/api/gen/proto/go/types/v1"
"github.com/grafana/pyroscope/v2/pkg/block/metadata"
"github.com/grafana/pyroscope/v2/pkg/featureflags"
"github.com/grafana/pyroscope/v2/pkg/frontend"
"github.com/grafana/pyroscope/v2/pkg/tenant"
"github.com/grafana/pyroscope/v2/pkg/test/mocks/mockfrontend"
"github.com/grafana/pyroscope/v2/pkg/test/mocks/mockmetastorev1"
Expand Down Expand Up @@ -186,6 +187,7 @@ func Test_QueryFrontend_LabelNames_WithFiltering(t *testing.T) {
qf := NewQueryFrontend(
log.NewNopLogger(),
mockLimits,
frontend.Config{},
mockMetadataClient,
nil,
mockQueryBackend,
Expand Down Expand Up @@ -329,6 +331,7 @@ func Test_QueryFrontend_Series_WithLabelNameFiltering(t *testing.T) {
qf := NewQueryFrontend(
log.NewNopLogger(),
mockLimits,
frontend.Config{},
mockMetadataClient,
nil,
mockQueryBackend,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@ import (
queryv1 "github.com/grafana/pyroscope/api/gen/proto/go/query/v1"
typesv1 "github.com/grafana/pyroscope/api/gen/proto/go/types/v1"
"github.com/grafana/pyroscope/v2/pkg/block/metadata"
"github.com/grafana/pyroscope/v2/pkg/frontend"
phlaremodel "github.com/grafana/pyroscope/v2/pkg/model"
"github.com/grafana/pyroscope/v2/pkg/pprof"
"github.com/grafana/pyroscope/v2/pkg/tenant"
Expand All @@ -41,6 +42,7 @@ func newSMPQueryFrontend(
return NewQueryFrontend(
log.NewNopLogger(),
limits,
frontend.Config{},
metaClient,
nil, // tenantServiceClient
backend,
Expand Down Expand Up @@ -682,6 +684,7 @@ func TestSelectMergeProfiles_Symbolization(t *testing.T) {
qf := NewQueryFrontend(
log.NewNopLogger(),
mockLimits,
frontend.Config{},
mockMetadataClient,
nil,
mockQueryBackend,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@ import (
querierv1 "github.com/grafana/pyroscope/api/gen/proto/go/querier/v1"
queryv1 "github.com/grafana/pyroscope/api/gen/proto/go/query/v1"
"github.com/grafana/pyroscope/v2/pkg/block/metadata"
"github.com/grafana/pyroscope/v2/pkg/frontend"
"github.com/grafana/pyroscope/v2/pkg/pprof"
"github.com/grafana/pyroscope/v2/pkg/tenant"
"github.com/grafana/pyroscope/v2/pkg/test/mocks/mockfrontend"
Expand Down Expand Up @@ -172,6 +173,7 @@ func TestSelectMergeSpanProfile_Symbolization(t *testing.T) {
qf := NewQueryFrontend(
log.NewNopLogger(),
mockLimits,
frontend.Config{},
mockMetadataClient,
nil,
mockQueryBackend,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@ import (
querierv1 "github.com/grafana/pyroscope/api/gen/proto/go/querier/v1"
queryv1 "github.com/grafana/pyroscope/api/gen/proto/go/query/v1"
"github.com/grafana/pyroscope/v2/pkg/block/metadata"
"github.com/grafana/pyroscope/v2/pkg/frontend"
phlaremodel "github.com/grafana/pyroscope/v2/pkg/model"
"github.com/grafana/pyroscope/v2/pkg/pprof"
"github.com/grafana/pyroscope/v2/pkg/tenant"
Expand Down Expand Up @@ -280,6 +281,7 @@ func TestSelectMergeStacktrace_Symbolization(t *testing.T) {
qf := NewQueryFrontend(
log.NewNopLogger(),
mockLimits,
frontend.Config{},
mockMetadataClient,
nil,
mockQueryBackend,
Expand Down Expand Up @@ -354,7 +356,7 @@ func TestSelectMergeStacktraces_DotFormat(t *testing.T) {
}},
}, nil)

qf := NewQueryFrontend(log.NewNopLogger(), mockLimits, mockMetadataClient, nil, mockQueryBackend, nil, nil, nil)
qf := NewQueryFrontend(log.NewNopLogger(), mockLimits, frontend.Config{}, mockMetadataClient, nil, mockQueryBackend, nil, nil, nil)

ctx := tenant.InjectTenantID(context.Background(), "tenant1")
start, end := smpValidTimeRange()
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ import (
"github.com/stretchr/testify/require"

querierv1 "github.com/grafana/pyroscope/api/gen/proto/go/querier/v1"
"github.com/grafana/pyroscope/v2/pkg/frontend"
"github.com/grafana/pyroscope/v2/pkg/tenant"
"github.com/grafana/pyroscope/v2/pkg/test/mocks/mockfrontend"
)
Expand All @@ -24,7 +25,7 @@ func TestSelectSeries_RejectsSubMillisecondStep(t *testing.T) {
limits.On("MaxQueryLookback", "test-tenant").Return(time.Duration(0)).Maybe()
limits.On("MaxQueryLength", "test-tenant").Return(time.Duration(0)).Maybe()

qf := NewQueryFrontend(log.NewNopLogger(), limits, nil, nil, nil, nil, nil, nil)
qf := NewQueryFrontend(log.NewNopLogger(), limits, frontend.Config{}, nil, nil, nil, nil, nil, nil)
ctx := tenant.InjectTenantID(context.Background(), "test-tenant")

_, err := qf.SelectSeries(ctx, connect.NewRequest(&querierv1.SelectSeriesRequest{
Expand All @@ -46,7 +47,7 @@ func TestSelectHeatmap_RejectsSubMillisecondStep(t *testing.T) {
for _, step := range []float64{0, 0.0001, 0.0005, 0.0009999} {
t.Run("step="+formatStep(step), func(t *testing.T) {
limits := mockfrontend.NewMockLimits(t)
qf := NewQueryFrontend(log.NewNopLogger(), limits, nil, nil, nil, nil, nil, nil)
qf := NewQueryFrontend(log.NewNopLogger(), limits, frontend.Config{}, nil, nil, nil, nil, nil, nil)
ctx := tenant.InjectTenantID(context.Background(), "test-tenant")

_, err := qf.SelectHeatmap(ctx, connect.NewRequest(&querierv1.SelectHeatmapRequest{
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@ import (
metastorev1 "github.com/grafana/pyroscope/api/gen/proto/go/metastore/v1"
querierv1 "github.com/grafana/pyroscope/api/gen/proto/go/querier/v1"
typesv1 "github.com/grafana/pyroscope/api/gen/proto/go/types/v1"
"github.com/grafana/pyroscope/v2/pkg/frontend"
phlaremodel "github.com/grafana/pyroscope/v2/pkg/model"
"github.com/grafana/pyroscope/v2/pkg/tenant"
"github.com/grafana/pyroscope/v2/pkg/test/mocks/mockfrontend"
Expand Down Expand Up @@ -136,6 +137,7 @@ func Test_QueryFrontend_Series_ProfileTypeQueryServedFromMetadata(t *testing.T)
qf := NewQueryFrontend(
log.NewNopLogger(),
mockLimits,
frontend.Config{},
mockMetadataClient,
nil,
mockQueryBackend,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@ import (
querierv1 "github.com/grafana/pyroscope/api/gen/proto/go/querier/v1"
queryv1 "github.com/grafana/pyroscope/api/gen/proto/go/query/v1"
"github.com/grafana/pyroscope/lidia"
"github.com/grafana/pyroscope/v2/pkg/frontend"
phlaremodel "github.com/grafana/pyroscope/v2/pkg/model"
"github.com/grafana/pyroscope/v2/pkg/model/symbolref"
"github.com/grafana/pyroscope/v2/pkg/tenant"
Expand Down Expand Up @@ -177,7 +178,7 @@ func TestSelectMergeStacktracesTree_SymbolRefFlagOn(t *testing.T) {
}},
}, nil).Once()

qf := NewQueryFrontend(log.NewNopLogger(), mockLimits, mockMetadataClient, nil, mockQueryBackend, mockSymbolizer, nil, nil)
qf := NewQueryFrontend(log.NewNopLogger(), mockLimits, frontend.Config{}, mockMetadataClient, nil, mockQueryBackend, mockSymbolizer, nil, nil)

ctx := tenant.InjectTenantID(context.Background(), "tenant1")
start, end := smpValidTimeRange()
Expand Down Expand Up @@ -235,7 +236,7 @@ func TestSelectMergeStacktracesTree_SymbolRefResolution(t *testing.T) {
Blocks: []*metastorev1.BlockMeta{{Id: "block_id"}},
}, nil).Once()

qf := NewQueryFrontend(log.NewNopLogger(), mockLimits, mockMetadataClient, nil, mockQueryBackend, mockSymbolizer, nil, nil)
qf := NewQueryFrontend(log.NewNopLogger(), mockLimits, frontend.Config{}, mockMetadataClient, nil, mockQueryBackend, mockSymbolizer, nil, nil)

before := testutil.ToFloat64(qf.metrics.symbolRefLocationsTotal.WithLabelValues(symbolRefLocationResolved))

Expand Down Expand Up @@ -269,7 +270,7 @@ func TestSelectMergeStacktracesTree_SymbolRefResolution(t *testing.T) {
// path (see TestRebuildInlineChainExpansionOrder for the Rebuild-side
// contract).
func TestBuildLookup_ReversesLidiaFrameOrder(t *testing.T) {
qf := NewQueryFrontend(log.NewNopLogger(), nil, nil, nil, nil, nil, nil, nil)
qf := NewQueryFrontend(log.NewNopLogger(), nil, frontend.Config{}, nil, nil, nil, nil, nil, nil)
lookup := qf.buildLookup([]binaryResolution{{
binary: symbolref.UnresolvedBinary{BuildID: "build-a", BinaryName: "libfoo.so", Addresses: []uint64{0x100}},
frames: [][]lidia.SourceInfoFrame{{
Expand Down
2 changes: 2 additions & 0 deletions pkg/pyroscope/modules_experimental.go
Original file line number Diff line number Diff line change
Expand Up @@ -87,6 +87,7 @@ func (f *Pyroscope) initQueryFrontendV2() (services.Service, error) {
f.queryFrontend = queryfrontend.NewQueryFrontend(
queryFrontendLogger,
f.Overrides,
f.Cfg.Frontend,
f.metastoreClient,
f.metastoreClient,
f.queryBackendClient,
Expand Down Expand Up @@ -145,6 +146,7 @@ func (f *Pyroscope) initQueryFrontendV12() (services.Service, error) {
f.queryFrontend = queryfrontend.NewQueryFrontend(
queryFrontendLogger,
f.Overrides,
f.Cfg.Frontend,
f.metastoreClient,
f.metastoreClient,
f.queryBackendClient,
Expand Down
Loading
Loading