Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
Split query by interval (grafana/phlare#713)
* Separate handlers in querier fronted * Clarify http responce decompression implementation details * Draft query time split * Split SelectSeries by time * Align interval to step duration * Remove unused code * Add querier.max-concurrent option * Fix connect headers
- Loading branch information
1 parent
01786d6
commit 4787ad4
Showing
34 changed files
with
869 additions
and
69 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,14 @@ | ||
package frontend | ||
|
||
import ( | ||
"context" | ||
|
||
"github.com/bufbuild/connect-go" | ||
|
||
querierv1 "github.com/grafana/phlare/api/gen/proto/go/querier/v1" | ||
"github.com/grafana/phlare/pkg/util/connectgrpc" | ||
) | ||
|
||
func (f *Frontend) Diff(ctx context.Context, c *connect.Request[querierv1.DiffRequest]) (*connect.Response[querierv1.DiffResponse], error) { | ||
return connectgrpc.RoundTripUnary[querierv1.DiffRequest, querierv1.DiffResponse](ctx, f, c) | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,14 @@ | ||
package frontend | ||
|
||
import ( | ||
"context" | ||
|
||
"github.com/bufbuild/connect-go" | ||
|
||
typesv1 "github.com/grafana/phlare/api/gen/proto/go/types/v1" | ||
"github.com/grafana/phlare/pkg/util/connectgrpc" | ||
) | ||
|
||
func (f *Frontend) LabelNames(ctx context.Context, c *connect.Request[typesv1.LabelNamesRequest]) (*connect.Response[typesv1.LabelNamesResponse], error) { | ||
return connectgrpc.RoundTripUnary[typesv1.LabelNamesRequest, typesv1.LabelNamesResponse](ctx, f, c) | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,14 @@ | ||
package frontend | ||
|
||
import ( | ||
"context" | ||
|
||
"github.com/bufbuild/connect-go" | ||
|
||
typesv1 "github.com/grafana/phlare/api/gen/proto/go/types/v1" | ||
"github.com/grafana/phlare/pkg/util/connectgrpc" | ||
) | ||
|
||
func (f *Frontend) LabelValues(ctx context.Context, c *connect.Request[typesv1.LabelValuesRequest]) (*connect.Response[typesv1.LabelValuesResponse], error) { | ||
return connectgrpc.RoundTripUnary[typesv1.LabelValuesRequest, typesv1.LabelValuesResponse](ctx, f, c) | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,14 @@ | ||
package frontend | ||
|
||
import ( | ||
"context" | ||
|
||
"github.com/bufbuild/connect-go" | ||
|
||
querierv1 "github.com/grafana/phlare/api/gen/proto/go/querier/v1" | ||
"github.com/grafana/phlare/pkg/util/connectgrpc" | ||
) | ||
|
||
func (f *Frontend) ProfileTypes(ctx context.Context, c *connect.Request[querierv1.ProfileTypesRequest]) (*connect.Response[querierv1.ProfileTypesResponse], error) { | ||
return connectgrpc.RoundTripUnary[querierv1.ProfileTypesRequest, querierv1.ProfileTypesResponse](ctx, f, c) | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,15 @@ | ||
package frontend | ||
|
||
import ( | ||
"context" | ||
|
||
"github.com/bufbuild/connect-go" | ||
|
||
profilev1 "github.com/grafana/phlare/api/gen/proto/go/google/v1" | ||
querierv1 "github.com/grafana/phlare/api/gen/proto/go/querier/v1" | ||
"github.com/grafana/phlare/pkg/util/connectgrpc" | ||
) | ||
|
||
func (f *Frontend) SelectMergeProfile(ctx context.Context, c *connect.Request[querierv1.SelectMergeProfileRequest]) (*connect.Response[profilev1.Profile], error) { | ||
return connectgrpc.RoundTripUnary[querierv1.SelectMergeProfileRequest, profilev1.Profile](ctx, f, c) | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,64 @@ | ||
package frontend | ||
|
||
import ( | ||
"context" | ||
"net/http" | ||
"time" | ||
|
||
"github.com/bufbuild/connect-go" | ||
"github.com/grafana/dskit/tenant" | ||
"golang.org/x/sync/errgroup" | ||
|
||
querierv1 "github.com/grafana/phlare/api/gen/proto/go/querier/v1" | ||
phlaremodel "github.com/grafana/phlare/pkg/model" | ||
"github.com/grafana/phlare/pkg/util/connectgrpc" | ||
"github.com/grafana/phlare/pkg/util/httpgrpc" | ||
"github.com/grafana/phlare/pkg/util/validation" | ||
) | ||
|
||
func (f *Frontend) SelectMergeStacktraces(ctx context.Context, | ||
c *connect.Request[querierv1.SelectMergeStacktracesRequest]) ( | ||
*connect.Response[querierv1.SelectMergeStacktracesResponse], error) { | ||
tenantIDs, err := tenant.TenantIDs(ctx) | ||
if err != nil { | ||
return nil, httpgrpc.Errorf(http.StatusBadRequest, err.Error()) | ||
} | ||
|
||
g, ctx := errgroup.WithContext(ctx) | ||
if maxConcurrent := validation.SmallestPositiveNonZeroIntPerTenant(tenantIDs, f.limits.MaxQueryParallelism); maxConcurrent > 0 { | ||
g.SetLimit(maxConcurrent) | ||
} | ||
|
||
m := phlaremodel.NewFlameGraphMerger() | ||
interval := validation.MaxDurationOrZeroPerTenant(tenantIDs, f.limits.QuerySplitDuration) | ||
intervals := NewTimeIntervalIterator(time.UnixMilli(c.Msg.Start), time.UnixMilli(c.Msg.End), interval) | ||
|
||
for intervals.Next() { | ||
r := intervals.At() | ||
g.Go(func() error { | ||
req := connectgrpc.CloneRequest(c, &querierv1.SelectMergeStacktracesRequest{ | ||
ProfileTypeID: c.Msg.ProfileTypeID, | ||
LabelSelector: c.Msg.LabelSelector, | ||
Start: r.Start.UnixMilli(), | ||
End: r.End.UnixMilli(), | ||
MaxNodes: c.Msg.MaxNodes, | ||
}) | ||
resp, err := connectgrpc.RoundTripUnary[ | ||
querierv1.SelectMergeStacktracesRequest, | ||
querierv1.SelectMergeStacktracesResponse](ctx, f, req) | ||
if err != nil { | ||
return err | ||
} | ||
m.MergeFlameGraph(resp.Msg.Flamegraph) | ||
return nil | ||
}) | ||
} | ||
|
||
if err = g.Wait(); err != nil { | ||
return nil, err | ||
} | ||
|
||
return connect.NewResponse(&querierv1.SelectMergeStacktracesResponse{ | ||
Flamegraph: m.FlameGraph(c.Msg.GetMaxNodes()), | ||
}), nil | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,64 @@ | ||
package frontend | ||
|
||
import ( | ||
"context" | ||
"net/http" | ||
"time" | ||
|
||
"github.com/bufbuild/connect-go" | ||
"github.com/grafana/dskit/tenant" | ||
"golang.org/x/sync/errgroup" | ||
|
||
querierv1 "github.com/grafana/phlare/api/gen/proto/go/querier/v1" | ||
phlaremodel "github.com/grafana/phlare/pkg/model" | ||
"github.com/grafana/phlare/pkg/util/connectgrpc" | ||
"github.com/grafana/phlare/pkg/util/httpgrpc" | ||
"github.com/grafana/phlare/pkg/util/validation" | ||
) | ||
|
||
func (f *Frontend) SelectSeries(ctx context.Context, | ||
c *connect.Request[querierv1.SelectSeriesRequest]) ( | ||
*connect.Response[querierv1.SelectSeriesResponse], error) { | ||
tenantIDs, err := tenant.TenantIDs(ctx) | ||
if err != nil { | ||
return nil, httpgrpc.Errorf(http.StatusBadRequest, err.Error()) | ||
} | ||
|
||
g, ctx := errgroup.WithContext(ctx) | ||
if maxConcurrent := validation.SmallestPositiveNonZeroIntPerTenant(tenantIDs, f.limits.MaxQueryParallelism); maxConcurrent > 0 { | ||
g.SetLimit(maxConcurrent) | ||
} | ||
|
||
m := phlaremodel.NewSeriesMerger(false) | ||
interval := validation.MaxDurationOrZeroPerTenant(tenantIDs, f.limits.QuerySplitDuration) | ||
intervals := NewTimeIntervalIterator(time.UnixMilli(c.Msg.Start), time.UnixMilli(c.Msg.End), interval, | ||
WithAlignment(time.Second*time.Duration(c.Msg.Step))) | ||
|
||
for intervals.Next() { | ||
r := intervals.At() | ||
g.Go(func() error { | ||
req := connectgrpc.CloneRequest(c, &querierv1.SelectSeriesRequest{ | ||
ProfileTypeID: c.Msg.ProfileTypeID, | ||
LabelSelector: c.Msg.LabelSelector, | ||
Start: r.Start.UnixMilli(), | ||
End: r.End.UnixMilli(), | ||
GroupBy: c.Msg.GroupBy, | ||
Step: c.Msg.Step, | ||
}) | ||
resp, err := connectgrpc.RoundTripUnary[ | ||
querierv1.SelectSeriesRequest, | ||
querierv1.SelectSeriesResponse](ctx, f, req) | ||
if err != nil { | ||
return err | ||
} | ||
m.MergeSeries(resp.Msg.Series) | ||
return nil | ||
}) | ||
} | ||
|
||
if err = g.Wait(); err != nil { | ||
return nil, err | ||
} | ||
|
||
return connect.NewResponse(&querierv1.SelectSeriesResponse{Series: m.Series()}), nil | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,14 @@ | ||
package frontend | ||
|
||
import ( | ||
"context" | ||
|
||
"github.com/bufbuild/connect-go" | ||
|
||
querierv1 "github.com/grafana/phlare/api/gen/proto/go/querier/v1" | ||
"github.com/grafana/phlare/pkg/util/connectgrpc" | ||
) | ||
|
||
func (f *Frontend) Series(ctx context.Context, c *connect.Request[querierv1.SeriesRequest]) (*connect.Response[querierv1.SeriesResponse], error) { | ||
return connectgrpc.RoundTripUnary[querierv1.SeriesRequest, querierv1.SeriesResponse](ctx, f, c) | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Oops, something went wrong.