mirror of
https://github.com/SigNoz/signoz.git
synced 2026-09-20 18:30:41 +01:00
The semantic convention family as first call citizen revealed that the
current state of the query builder needs a bit refactoring for long term
maintenance.
The `FieldMapper` and `ConditionBuilder` are now one abstraction
`Storage`.
A storage now answers
- what the compiler cannot know i.e one read per field key (the bare
SQL, the membership present or absent, what an absent row reads, and
whether the read keeps its type or filters only).
- the fallback for a key metadata does not report
- its traits
- and one Condition compilation part.
And we introduce a new type to use in the system, `Resolved`
```
// Resolved is what resolution produces for one key: its meanings, and how
// they came to be. It is the only thing the compilers receive. Compile it
// with the operator and value it was resolved with.
type Resolved struct {
Key *telemetrytypes.TelemetryFieldKey
Fields []*telemetrytypes.LogicalField
// FromFallback: the fields came from the storage's fallback, not from
// metadata matches.
FromFallback bool
// Ambiguous: the matches held several interpretations.
Ambiguous bool
// Skipped: the storage contributes nothing for this key.
Skipped bool
Warnings []string
}
```
The prepared SQL has no changes, where it changed, it specifically made
the expression better by removing the redundant part.
- The prepared SQL remains identical with this refactoring
- No changes to integration tests
Assisted-by: Claude Fable 5.1
1034 lines
40 KiB
Go
1034 lines
40 KiB
Go
package metricsstatementbuilder
|
|
|
|
import (
|
|
"context"
|
|
"fmt"
|
|
"log/slog"
|
|
"math"
|
|
"slices"
|
|
"strconv"
|
|
"strings"
|
|
"time"
|
|
|
|
"github.com/SigNoz/signoz/pkg/clickhousesql"
|
|
"github.com/SigNoz/signoz/pkg/errors"
|
|
"github.com/SigNoz/signoz/pkg/factory"
|
|
"github.com/SigNoz/signoz/pkg/flagger"
|
|
"github.com/SigNoz/signoz/pkg/querybuilder"
|
|
"github.com/SigNoz/signoz/pkg/statementbuilder"
|
|
"github.com/SigNoz/signoz/pkg/telemetryschema/metricstelemetryschema"
|
|
"github.com/SigNoz/signoz/pkg/types/metrictypes"
|
|
qbtypes "github.com/SigNoz/signoz/pkg/types/querybuildertypes/querybuildertypesv5"
|
|
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
|
|
"github.com/SigNoz/signoz/pkg/valuer"
|
|
"github.com/huandu/go-sqlbuilder"
|
|
)
|
|
|
|
const (
|
|
RateTmpl = `multiIf(row_number() OVER rate_window = 1, nan, (per_series_value - lagInFrame(per_series_value, 1) OVER rate_window) < 0, per_series_value / (ts - lagInFrame(ts, 1) OVER rate_window), (per_series_value - lagInFrame(per_series_value, 1) OVER rate_window) / (ts - lagInFrame(ts, 1) OVER rate_window))`
|
|
|
|
IncreaseTmpl = `multiIf(row_number() OVER rate_window = 1, nan, (per_series_value - lagInFrame(per_series_value, 1) OVER rate_window) < 0, per_series_value, per_series_value - lagInFrame(per_series_value, 1) OVER rate_window)`
|
|
|
|
RateMultiTemporalityTmpl = `IF(LOWER(temporality) LIKE LOWER('delta'), %s, multiIf(row_number() OVER rate_window = 1, nan, (%s - lagInFrame(%s, 1) OVER rate_window) < 0, %s / (ts - lagInFrame(ts, 1) OVER rate_window), (%s - lagInFrame(%s, 1) OVER rate_window) / (ts - lagInFrame(ts, 1) OVER rate_window))) AS per_series_value`
|
|
|
|
IncreaseMultiTemporality = `IF(LOWER(temporality) LIKE LOWER('delta'), %s, multiIf(row_number() OVER rate_window = 1, nan, (%s - lagInFrame(%s, 1) OVER rate_window) < 0, %s, (%s - lagInFrame(%s, 1) OVER rate_window))) AS per_series_value`
|
|
|
|
OthersMultiTemporality = `IF(LOWER(temporality) LIKE LOWER('delta'), %s, %s) AS per_series_value`
|
|
)
|
|
|
|
type StatementBuilder struct {
|
|
logger *slog.Logger
|
|
metadataStore telemetrytypes.MetadataStore
|
|
storage qbtypes.Storage
|
|
flagger flagger.Flagger
|
|
}
|
|
|
|
var _ qbtypes.StatementBuilder[qbtypes.MetricAggregation] = (*StatementBuilder)(nil)
|
|
|
|
// NewFactory returns a provider factory for the metrics statement builder. Its
|
|
// New internalizes the storage and yields the concrete
|
|
// *StatementBuilder so the meter builder can reuse it.
|
|
func NewFactory(
|
|
metadataStore telemetrytypes.MetadataStore,
|
|
fl flagger.Flagger,
|
|
) factory.ProviderFactory[*StatementBuilder, statementbuilder.Config] {
|
|
return factory.NewProviderFactory(
|
|
factory.MustNewName("metrics"),
|
|
func(_ context.Context, settings factory.ProviderSettings, _ statementbuilder.Config) (*StatementBuilder, error) {
|
|
return NewMetricQueryStatementBuilder(settings, metadataStore, metricstelemetryschema.NewStorage(), fl), nil
|
|
},
|
|
)
|
|
}
|
|
|
|
func NewMetricQueryStatementBuilder(
|
|
settings factory.ProviderSettings,
|
|
metadataStore telemetrytypes.MetadataStore,
|
|
storage qbtypes.Storage,
|
|
flagger flagger.Flagger,
|
|
) *StatementBuilder {
|
|
metricsSettings := factory.NewScopedProviderSettings(settings, "github.com/SigNoz/signoz/pkg/telemetryschema/metricstelemetryschema")
|
|
return &StatementBuilder{
|
|
logger: metricsSettings.Logger(),
|
|
metadataStore: metadataStore,
|
|
storage: storage,
|
|
flagger: flagger,
|
|
}
|
|
}
|
|
|
|
func GetKeySelectors(query qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation]) []*telemetrytypes.FieldKeySelector {
|
|
var keySelectors []*telemetrytypes.FieldKeySelector
|
|
if query.Filter != nil && query.Filter.Expression != "" {
|
|
whereClauseSelectors := querybuilder.QueryStringToKeysSelectors(query.Filter.Expression)
|
|
keySelectors = append(keySelectors, whereClauseSelectors...)
|
|
}
|
|
|
|
for idx := range query.GroupBy {
|
|
groupBy := query.GroupBy[idx]
|
|
selectors := querybuilder.QueryStringToKeysSelectors(groupBy.Name)
|
|
keySelectors = append(keySelectors, selectors...)
|
|
}
|
|
|
|
for idx := range query.Order {
|
|
keySelectors = append(keySelectors, &telemetrytypes.FieldKeySelector{
|
|
Name: query.Order[idx].Key.Name,
|
|
Signal: telemetrytypes.SignalMetrics,
|
|
FieldContext: query.Order[idx].Key.FieldContext,
|
|
FieldDataType: query.Order[idx].Key.FieldDataType,
|
|
})
|
|
}
|
|
|
|
for idx := range keySelectors {
|
|
keySelectors[idx].Signal = telemetrytypes.SignalMetrics
|
|
keySelectors[idx].SelectorMatchType = telemetrytypes.FieldSelectorMatchTypeExact
|
|
keySelectors[idx].MetricContext = &telemetrytypes.MetricContext{
|
|
MetricName: query.Aggregations[0].MetricName,
|
|
}
|
|
keySelectors[idx].Source = query.Source
|
|
}
|
|
return keySelectors
|
|
}
|
|
|
|
func (b *StatementBuilder) Build(
|
|
ctx context.Context,
|
|
orgID valuer.UUID,
|
|
start uint64,
|
|
end uint64,
|
|
requestType qbtypes.RequestType,
|
|
query qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation],
|
|
variables map[string]qbtypes.VariableItem,
|
|
) (*qbtypes.Statement, error) {
|
|
keySelectors := GetKeySelectors(query)
|
|
keys, _, err := b.metadataStore.GetKeysMulti(ctx, orgID, keySelectors)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
start, end = querybuilder.AdjustedMetricTimeRange(start, end, uint64(query.StepInterval.Seconds()), query)
|
|
|
|
return b.buildPipelineStatement(ctx, orgID, start, end, requestType, query, keys, variables)
|
|
}
|
|
|
|
func (b *StatementBuilder) buildPipelineStatement(
|
|
ctx context.Context,
|
|
orgID valuer.UUID,
|
|
start, end uint64,
|
|
requestType qbtypes.RequestType,
|
|
query qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation],
|
|
keys map[string][]*telemetrytypes.TelemetryFieldKey,
|
|
variables map[string]qbtypes.VariableItem,
|
|
) (*qbtypes.Statement, error) {
|
|
var (
|
|
cteFragments []string
|
|
cteArgs [][]any
|
|
)
|
|
|
|
cteQuery := query
|
|
if query.Aggregations[0].Type == metrictypes.HistogramType {
|
|
query.GroupBy = slices.DeleteFunc(slices.Clone(query.GroupBy), isHistogramBucket)
|
|
cteQuery = rewriteQueryForHistogramCTE(requestType, query)
|
|
}
|
|
|
|
agg := cteQuery.Aggregations[0]
|
|
|
|
// A reduced metric reads the raw buffer for recent short windows, and
|
|
// samples_v4/agg (unioned with the reduced tables) otherwise. The buffer is
|
|
// shaped exactly like samples_v4 / time_series_v4, so once the table names are
|
|
// chosen the rest of the pipeline is unchanged.
|
|
useBuffer := agg.Reduced &&
|
|
end-start < metricstelemetryschema.OneDayInMilliseconds &&
|
|
start >= uint64(time.Now().UnixMilli())-metricstelemetryschema.OneDayInMilliseconds
|
|
|
|
samplesTable, _ := metricstelemetryschema.WhichSamplesTableToUse(start, end, agg.Type, agg.TimeAggregation, useBuffer, agg.TableHints)
|
|
tsStart, tsEnd, _, tsTable := metricstelemetryschema.WhichTSTableToUse(start, end, useBuffer, agg.TableHints)
|
|
|
|
var timeSeriesCTE string
|
|
var timeSeriesCTEArgs []any
|
|
var filterWarnings []string
|
|
var err error
|
|
|
|
if timeSeriesCTE, timeSeriesCTEArgs, filterWarnings, err = b.buildTimeSeriesCTE(ctx, orgID, tsStart, tsEnd, cteQuery, keys, variables, tsTable); err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
if qbtypes.CanShortCircuitDelta(agg) {
|
|
// spatial_aggregation_cte directly for certain delta queries
|
|
if frag, args, err := b.buildTemporalAggDeltaFastPath(start, end, cteQuery, samplesTable, timeSeriesCTE, timeSeriesCTEArgs); err != nil {
|
|
return nil, err
|
|
} else if frag != "" {
|
|
cteFragments = append(cteFragments, frag)
|
|
cteArgs = append(cteArgs, args)
|
|
}
|
|
} else {
|
|
// temporal_aggregation_cte
|
|
if frag, args, err := b.buildTemporalAggregationCTE(ctx, start, end, cteQuery, keys, samplesTable, timeSeriesCTE, timeSeriesCTEArgs); err != nil {
|
|
return nil, err
|
|
} else if frag != "" {
|
|
cteFragments = append(cteFragments, frag)
|
|
cteArgs = append(cteArgs, args)
|
|
}
|
|
|
|
// spatial_aggregation_cte
|
|
if frag, args := b.buildSpatialAggregationCTE(ctx, start, end, cteQuery, keys); frag != "" {
|
|
cteFragments = append(cteFragments, frag)
|
|
cteArgs = append(cteArgs, args)
|
|
}
|
|
}
|
|
|
|
var reducedFragments []string
|
|
var reducedArgs [][]any
|
|
if agg.Reduced && !useBuffer {
|
|
var tsCTE string
|
|
var tsArgs []any
|
|
// time series rows are written on hour boundaries
|
|
tsStart := start - (start % metricstelemetryschema.OneHourInMilliseconds)
|
|
if tsCTE, tsArgs, err = b.buildReducedTimeSeriesCTE(ctx, orgID, tsStart, end, cteQuery, keys, variables); err != nil {
|
|
return nil, err
|
|
}
|
|
if qbtypes.CanShortCircuitReduced(agg) {
|
|
// spatial_aggregation_cte directly, no per-series level
|
|
if spatialFrag, spatialArgs, ok := b.buildReducedSpatialAggFastPath(start, end, cteQuery, tsCTE, tsArgs); ok {
|
|
reducedFragments = []string{spatialFrag}
|
|
reducedArgs = [][]any{spatialArgs}
|
|
}
|
|
} else if temporalFrag, temporalArgs, ok := b.buildReducedTemporalAggregationCTE(start, end, cteQuery, tsCTE, tsArgs); ok {
|
|
spatialFrag, spatialArgs := b.buildReducedSpatialAggregationCTE(cteQuery)
|
|
reducedFragments = []string{temporalFrag, spatialFrag}
|
|
reducedArgs = [][]any{temporalArgs, spatialArgs}
|
|
}
|
|
}
|
|
|
|
mainStmt, err := b.BuildFinalSelect(cteFragments, cteArgs, requestType, query)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
mainStmt.Warnings = append(mainStmt.Warnings, filterWarnings...)
|
|
if reducedFragments == nil {
|
|
return mainStmt, nil
|
|
}
|
|
reducedStmt, err := b.BuildFinalSelect(reducedFragments, reducedArgs, requestType, query)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
return unionStatements(mainStmt, reducedStmt, query)
|
|
}
|
|
|
|
func rewriteQueryForHistogramCTE(requestType qbtypes.RequestType, query qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation]) qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation] {
|
|
query.GroupBy = append(slices.Clone(query.GroupBy), qbtypes.GroupByKey{
|
|
TelemetryFieldKey: telemetrytypes.TelemetryFieldKey{Name: histogramBucketKey},
|
|
})
|
|
|
|
query.Aggregations = slices.Clone(query.Aggregations)
|
|
if query.Aggregations[0].SpaceAggregation.IsPercentile() && requestType != qbtypes.RequestTypeHeatmap {
|
|
query.Aggregations[0].TimeAggregation = metrictypes.TimeAggregationRate
|
|
} else {
|
|
query.Aggregations[0].TimeAggregation = metrictypes.TimeAggregationIncrease
|
|
}
|
|
query.Aggregations[0].SpaceAggregation = metrictypes.SpaceAggregationSum
|
|
|
|
return query
|
|
}
|
|
|
|
func unionStatements(main, reduced *qbtypes.Statement, query qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation]) (*qbtypes.Statement, error) {
|
|
orderBy := "ts"
|
|
for i, g := range query.GroupBy {
|
|
orderBy = GroupByColumnAlias(i, g.Name) + ", " + orderBy
|
|
}
|
|
q := fmt.Sprintf(
|
|
"SELECT * FROM (%s) UNION ALL SELECT * FROM (%s) ORDER BY %s SETTINGS do_not_merge_across_partitions_select_final = 1, optimize_move_to_prewhere_if_final = 1",
|
|
main.Query, reduced.Query, orderBy,
|
|
)
|
|
args := append(append([]any{}, main.Args...), reduced.Args...)
|
|
warnings := append(append([]string{}, main.Warnings...), reduced.Warnings...)
|
|
return &qbtypes.Statement{Query: q, Args: args, Warnings: warnings}, nil
|
|
}
|
|
|
|
func (b *StatementBuilder) buildReducedTimeSeriesCTE(
|
|
ctx context.Context,
|
|
orgID valuer.UUID,
|
|
start, end uint64,
|
|
query qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation],
|
|
keys map[string][]*telemetrytypes.TelemetryFieldKey,
|
|
variables map[string]qbtypes.VariableItem,
|
|
) (string, []any, error) {
|
|
sb := sqlbuilder.NewSelectBuilder()
|
|
info := querybuilder.NewQueryInfo(ctx, orgID, b.flagger, telemetrytypes.SignalMetrics, &telemetrytypes.MetricContext{MetricName: query.Aggregations[0].MetricName}, start, end)
|
|
|
|
var preparedWhereClause querybuilder.PreparedWhereClause
|
|
var err error
|
|
if query.Filter != nil && query.Filter.Expression != "" {
|
|
preparedWhereClause, err = querybuilder.PrepareWhereClause(query.Filter.Expression, querybuilder.FilterExprVisitorOpts{
|
|
Context: ctx,
|
|
Query: info,
|
|
Storage: b.storage,
|
|
Logger: b.logger,
|
|
FieldKeys: keys,
|
|
FullTextColumn: &telemetrytypes.TelemetryFieldKey{Name: "labels"},
|
|
Variables: variables,
|
|
})
|
|
if err != nil {
|
|
return "", nil, err
|
|
}
|
|
}
|
|
|
|
sb.From(fmt.Sprintf("%s.%s", metricstelemetryschema.DBName, metricstelemetryschema.TimeseriesV4ReducedLocalTableName))
|
|
sb.Select("fingerprint")
|
|
for i, g := range query.GroupBy {
|
|
col, err := querybuilder.ResolveColumn(ctx, info, b.storage, &g.TelemetryFieldKey, telemetrytypes.FieldDataTypeString, keys)
|
|
if err != nil {
|
|
return "", nil, err
|
|
}
|
|
sb.SelectMore(sqlbuilder.Escape(fmt.Sprintf("%s AS %s", col, GroupByColumnAlias(i, g.Name))))
|
|
}
|
|
sb.Where(
|
|
sb.In("metric_name", query.Aggregations[0].MetricName),
|
|
sb.GTE("unix_milli", start),
|
|
sb.LTE("unix_milli", end),
|
|
)
|
|
|
|
if !preparedWhereClause.IsEmpty() {
|
|
sb.AddWhereClause(preparedWhereClause.WhereClause)
|
|
}
|
|
sb.GroupBy("fingerprint")
|
|
sb.GroupBy(GroupByAliases(query.GroupBy)...)
|
|
|
|
q, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
|
|
// the caller joins this into another builder, which compiles it again
|
|
return fmt.Sprintf("(%s) AS filtered_time_series", sqlbuilder.Escape(q)), args, nil
|
|
}
|
|
|
|
// buildReducedSpatialAggFastPath is the reduced analog of
|
|
// buildTemporalAggDeltaFastPath: for combinations where the temporal and
|
|
// spatial aggregations collapse (CanShortCircuitReduced), it emits the
|
|
// spatial_aggregation_cte in one level with no per-series grouping, so shards
|
|
// send one state per (step, group) instead of per (series, step, group).
|
|
// FINAL still dedups recomputed 60s buckets at scan time.
|
|
func (b *StatementBuilder) buildReducedSpatialAggFastPath(
|
|
start, end uint64,
|
|
query qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation],
|
|
timeSeriesCTE string,
|
|
timeSeriesCTEArgs []any,
|
|
) (string, []any, bool) {
|
|
agg := query.Aggregations[0]
|
|
stepSec := int64(query.StepInterval.Seconds())
|
|
|
|
value, _, ok := metricstelemetryschema.ReducedValueColumn(agg.Type, agg.SpaceAggregation)
|
|
if !ok {
|
|
return "", nil, false
|
|
}
|
|
|
|
sb := sqlbuilder.NewSelectBuilder()
|
|
sb.Select(fmt.Sprintf("toStartOfInterval(toDateTime(intDiv(unix_milli, 1000)), toIntervalSecond(%d)) AS ts", stepSec))
|
|
for i, g := range query.GroupBy {
|
|
sb.SelectMore(sqlbuilder.Escape(GroupByColumnAlias(i, g.Name)))
|
|
}
|
|
sb.SelectMore(fmt.Sprintf("%s AS value", metricstelemetryschema.ReducedTimeAggregationColumn(agg.TimeAggregation, stepSec, value)))
|
|
sb.From(fmt.Sprintf("%s.%s AS points FINAL", metricstelemetryschema.DBName, metricstelemetryschema.WhichReducedSamplesTableToUse(agg.Type)))
|
|
sb.JoinWithOption(sqlbuilder.InnerJoin, timeSeriesCTE, "points.reduced_fingerprint = filtered_time_series.fingerprint")
|
|
sb.Where(
|
|
sb.In("metric_name", agg.MetricName),
|
|
sb.GTE("unix_milli", start),
|
|
sb.LT("unix_milli", end),
|
|
)
|
|
sb.GroupBy("ts")
|
|
sb.GroupBy(GroupByAliases(query.GroupBy)...)
|
|
|
|
q, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse, timeSeriesCTEArgs...)
|
|
return fmt.Sprintf("__spatial_aggregation_cte AS (%s)", q), args, true
|
|
}
|
|
|
|
func (b *StatementBuilder) buildReducedTemporalAggregationCTE(
|
|
start, end uint64,
|
|
query qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation],
|
|
timeSeriesCTE string,
|
|
timeSeriesCTEArgs []any,
|
|
) (string, []any, bool) {
|
|
agg := query.Aggregations[0]
|
|
stepSec := int64(query.StepInterval.Seconds())
|
|
|
|
value, weight, ok := metricstelemetryschema.ReducedValueColumn(agg.Type, agg.SpaceAggregation)
|
|
if !ok {
|
|
return "", nil, false
|
|
}
|
|
|
|
// TODO(srikanthccv): add _5m/_30m tables similar to samples_v4
|
|
// and wire them up in querier before GA
|
|
sb := sqlbuilder.NewSelectBuilder()
|
|
sb.Select("points.reduced_fingerprint AS fingerprint")
|
|
sb.SelectMore(fmt.Sprintf("toStartOfInterval(toDateTime(intDiv(unix_milli, 1000)), toIntervalSecond(%d)) AS ts", stepSec))
|
|
for i, g := range query.GroupBy {
|
|
sb.SelectMore(sqlbuilder.Escape(GroupByColumnAlias(i, g.Name)))
|
|
}
|
|
sb.SelectMore(fmt.Sprintf("%s AS per_series_value", metricstelemetryschema.ReducedTimeAggregationColumn(agg.TimeAggregation, stepSec, value)))
|
|
if weight != "" {
|
|
// count_series is a series count, not additive over time, so the avg
|
|
// denominator is reduced with avg
|
|
sb.SelectMore(fmt.Sprintf("avg(%s) AS per_series_weight", weight))
|
|
}
|
|
sb.From(fmt.Sprintf("%s.%s AS points FINAL", metricstelemetryschema.DBName, metricstelemetryschema.WhichReducedSamplesTableToUse(agg.Type)))
|
|
sb.JoinWithOption(sqlbuilder.InnerJoin, timeSeriesCTE, "points.reduced_fingerprint = filtered_time_series.fingerprint")
|
|
sb.Where(
|
|
sb.In("metric_name", agg.MetricName),
|
|
sb.GTE("unix_milli", start),
|
|
sb.LT("unix_milli", end),
|
|
)
|
|
sb.GroupBy("fingerprint", "ts")
|
|
sb.GroupBy(GroupByAliases(query.GroupBy)...)
|
|
|
|
q, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse, timeSeriesCTEArgs...)
|
|
return fmt.Sprintf("__temporal_aggregation_cte AS (%s)", q), args, true
|
|
}
|
|
|
|
func (b *StatementBuilder) buildReducedSpatialAggregationCTE(
|
|
query qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation],
|
|
) (string, []any) {
|
|
spatial := "sum(per_series_value)"
|
|
switch query.Aggregations[0].SpaceAggregation {
|
|
case metrictypes.SpaceAggregationAvg:
|
|
spatial = "sum(per_series_value) / sum(per_series_weight)"
|
|
case metrictypes.SpaceAggregationMin:
|
|
spatial = "min(per_series_value)"
|
|
case metrictypes.SpaceAggregationMax:
|
|
spatial = "max(per_series_value)"
|
|
}
|
|
|
|
sb := sqlbuilder.NewSelectBuilder()
|
|
sb.Select("ts")
|
|
for i, g := range query.GroupBy {
|
|
sb.SelectMore(sqlbuilder.Escape(GroupByColumnAlias(i, g.Name)))
|
|
}
|
|
sb.SelectMore(spatial + " AS value")
|
|
sb.From("__temporal_aggregation_cte")
|
|
sb.GroupBy("ts")
|
|
sb.GroupBy(GroupByAliases(query.GroupBy)...)
|
|
|
|
q, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
|
|
return fmt.Sprintf("__spatial_aggregation_cte AS (%s)", q), args
|
|
}
|
|
|
|
func (b *StatementBuilder) buildTemporalAggDeltaFastPath(
|
|
start, end uint64,
|
|
query qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation],
|
|
samplesTable string,
|
|
timeSeriesCTE string,
|
|
timeSeriesCTEArgs []any,
|
|
) (string, []any, error) {
|
|
stepSec := int64(query.StepInterval.Seconds())
|
|
|
|
sb := sqlbuilder.NewSelectBuilder()
|
|
|
|
sb.SelectMore(fmt.Sprintf(
|
|
"toStartOfInterval(toDateTime(intDiv(unix_milli, 1000)), toIntervalSecond(%d)) AS ts",
|
|
stepSec,
|
|
))
|
|
for i, g := range query.GroupBy {
|
|
sb.SelectMore(sqlbuilder.Escape(GroupByColumnAlias(i, g.Name)))
|
|
}
|
|
|
|
var aggCol string
|
|
if query.Aggregations[0].SpaceAggregation.IsPercentile() &&
|
|
query.Aggregations[0].Type == metrictypes.ExpHistogramType {
|
|
// merging sketches already spans every series in the step, so neither a
|
|
// samples-table value column nor the rate divisor applies
|
|
aggCol = fmt.Sprintf("quantilesDDMerge(0.01, %f)(sketch)[1]", query.Aggregations[0].SpaceAggregation.Percentile())
|
|
} else {
|
|
col, err := metricstelemetryschema.AggregationColumnForSamplesTable(
|
|
samplesTable, query.Aggregations[0].Temporality, query.Aggregations[0].TimeAggregation,
|
|
)
|
|
if err != nil {
|
|
return "", nil, err
|
|
}
|
|
aggCol = col
|
|
if query.Aggregations[0].TimeAggregation == metrictypes.TimeAggregationRate {
|
|
// TODO(srikanthccv): should it be step interval or use [start_time_unix_nano](https://github.com/open-telemetry/opentelemetry-proto/blob/d3fb76d70deb0874692bd0ebe03148580d85f3bb/opentelemetry/proto/metrics/v1/metrics.proto#L400C11-L400C31)?
|
|
aggCol = fmt.Sprintf("%s/%d", aggCol, stepSec)
|
|
}
|
|
}
|
|
|
|
sb.SelectMore(fmt.Sprintf("%s AS value", aggCol))
|
|
|
|
sb.From(fmt.Sprintf("%s.%s AS points", metricstelemetryschema.DBName, samplesTable))
|
|
sb.JoinWithOption(sqlbuilder.InnerJoin, timeSeriesCTE, "points.fingerprint = filtered_time_series.fingerprint")
|
|
sb.Where(
|
|
sb.In("metric_name", query.Aggregations[0].MetricName),
|
|
sb.GTE("unix_milli", start),
|
|
sb.LT("unix_milli", end),
|
|
)
|
|
sb.GroupBy("ts")
|
|
sb.GroupBy(GroupByAliases(query.GroupBy)...)
|
|
|
|
q, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse, timeSeriesCTEArgs...)
|
|
return fmt.Sprintf("__spatial_aggregation_cte AS (%s)", q), args, nil
|
|
}
|
|
|
|
func (b *StatementBuilder) buildTimeSeriesCTE(
|
|
ctx context.Context,
|
|
orgID valuer.UUID,
|
|
start, end uint64,
|
|
query qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation],
|
|
keys map[string][]*telemetrytypes.TelemetryFieldKey,
|
|
variables map[string]qbtypes.VariableItem,
|
|
tsTable string,
|
|
) (string, []any, []string, error) {
|
|
sb := sqlbuilder.NewSelectBuilder()
|
|
info := querybuilder.NewQueryInfo(ctx, orgID, b.flagger, telemetrytypes.SignalMetrics, &telemetrytypes.MetricContext{MetricName: query.Aggregations[0].MetricName}, start, end)
|
|
|
|
var preparedWhereClause querybuilder.PreparedWhereClause
|
|
var err error
|
|
|
|
if query.Filter != nil && query.Filter.Expression != "" {
|
|
preparedWhereClause, err = querybuilder.PrepareWhereClause(query.Filter.Expression, querybuilder.FilterExprVisitorOpts{
|
|
Context: ctx,
|
|
Query: info,
|
|
Storage: b.storage,
|
|
Logger: b.logger,
|
|
FieldKeys: keys,
|
|
FullTextColumn: &telemetrytypes.TelemetryFieldKey{Name: "labels"},
|
|
Variables: variables,
|
|
})
|
|
if err != nil {
|
|
return "", nil, nil, err
|
|
}
|
|
}
|
|
|
|
sb.From(fmt.Sprintf("%s.%s", metricstelemetryschema.DBName, tsTable))
|
|
|
|
sb.Select("fingerprint")
|
|
for i, g := range query.GroupBy {
|
|
col, err := querybuilder.ResolveColumn(ctx, info, b.storage, &g.TelemetryFieldKey, telemetrytypes.FieldDataTypeString, keys)
|
|
if err != nil {
|
|
return "", nil, nil, err
|
|
}
|
|
sb.SelectMore(sqlbuilder.Escape(fmt.Sprintf("%s AS %s", col, GroupByColumnAlias(i, g.Name))))
|
|
}
|
|
|
|
sb.Where(
|
|
sb.In("metric_name", query.Aggregations[0].MetricName),
|
|
sb.GTE("unix_milli", start),
|
|
sb.LTE("unix_milli", end),
|
|
)
|
|
|
|
if query.Aggregations[0].Temporality != metrictypes.Multiple && query.Aggregations[0].Temporality != metrictypes.Unknown {
|
|
sb.Where(sb.ILike("temporality", query.Aggregations[0].Temporality.StringValue()))
|
|
}
|
|
|
|
// the buffer holds both raw rows and the reduced catalog rows; the raw read
|
|
// only wants the original series
|
|
if tsTable == metricstelemetryschema.TimeseriesV4BufferLocalTableName {
|
|
sb.Where(sb.EQ("is_reduced", false))
|
|
}
|
|
|
|
if !preparedWhereClause.IsEmpty() {
|
|
sb.AddWhereClause(preparedWhereClause.WhereClause)
|
|
}
|
|
|
|
sb.GroupBy("fingerprint")
|
|
sb.GroupBy(GroupByAliases(query.GroupBy)...)
|
|
|
|
q, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
|
|
// the caller joins this into another builder, which compiles it again
|
|
return fmt.Sprintf("(%s) AS filtered_time_series", sqlbuilder.Escape(q)), args, preparedWhereClause.Warnings, nil
|
|
}
|
|
|
|
func (b *StatementBuilder) buildTemporalAggregationCTE(
|
|
ctx context.Context,
|
|
start, end uint64,
|
|
query qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation],
|
|
_ map[string][]*telemetrytypes.TelemetryFieldKey,
|
|
samplesTable string,
|
|
timeSeriesCTE string,
|
|
timeSeriesCTEArgs []any,
|
|
) (string, []any, error) {
|
|
if query.Aggregations[0].Temporality == metrictypes.Delta {
|
|
return b.buildTemporalAggDelta(ctx, start, end, query, samplesTable, timeSeriesCTE, timeSeriesCTEArgs)
|
|
} else if query.Aggregations[0].Temporality != metrictypes.Multiple {
|
|
return b.buildTemporalAggCumulativeOrUnspecified(ctx, start, end, query, samplesTable, timeSeriesCTE, timeSeriesCTEArgs)
|
|
}
|
|
return b.buildTemporalAggForMultipleTemporalities(ctx, start, end, query, samplesTable, timeSeriesCTE, timeSeriesCTEArgs)
|
|
}
|
|
|
|
func (b *StatementBuilder) buildTemporalAggDelta(
|
|
_ context.Context,
|
|
start, end uint64,
|
|
query qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation],
|
|
samplesTable string,
|
|
timeSeriesCTE string,
|
|
timeSeriesCTEArgs []any,
|
|
) (string, []any, error) {
|
|
stepSec := int64(query.StepInterval.Seconds())
|
|
|
|
sb := sqlbuilder.NewSelectBuilder()
|
|
|
|
sb.Select("fingerprint")
|
|
sb.SelectMore(fmt.Sprintf(
|
|
"toStartOfInterval(toDateTime(intDiv(unix_milli, 1000)), toIntervalSecond(%d)) AS ts",
|
|
stepSec,
|
|
))
|
|
for i, g := range query.GroupBy {
|
|
sb.SelectMore(sqlbuilder.Escape(GroupByColumnAlias(i, g.Name)))
|
|
}
|
|
|
|
aggCol, err := metricstelemetryschema.AggregationColumnForSamplesTable(samplesTable, query.Aggregations[0].Temporality, query.Aggregations[0].TimeAggregation)
|
|
if err != nil {
|
|
return "", nil, err
|
|
}
|
|
if query.Aggregations[0].TimeAggregation == metrictypes.TimeAggregationRate {
|
|
// TODO(srikanthccv): should it be step interval or use [start_time_unix_nano](https://github.com/open-telemetry/opentelemetry-proto/blob/d3fb76d70deb0874692bd0ebe03148580d85f3bb/opentelemetry/proto/metrics/v1/metrics.proto#L400C11-L400C31)?
|
|
aggCol = fmt.Sprintf("%s/%d", aggCol, stepSec)
|
|
}
|
|
|
|
sb.SelectMore(fmt.Sprintf("%s AS per_series_value", aggCol))
|
|
|
|
sb.From(fmt.Sprintf("%s.%s AS points", metricstelemetryschema.DBName, samplesTable))
|
|
sb.JoinWithOption(sqlbuilder.InnerJoin, timeSeriesCTE, "points.fingerprint = filtered_time_series.fingerprint")
|
|
sb.Where(
|
|
sb.In("metric_name", query.Aggregations[0].MetricName),
|
|
sb.GTE("unix_milli", start),
|
|
sb.LT("unix_milli", end),
|
|
)
|
|
sb.GroupBy("fingerprint", "ts")
|
|
sb.GroupBy(GroupByAliases(query.GroupBy)...)
|
|
sb.OrderBy("fingerprint", "ts")
|
|
|
|
q, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse, timeSeriesCTEArgs...)
|
|
return fmt.Sprintf("__temporal_aggregation_cte AS (%s)", q), args, nil
|
|
}
|
|
|
|
func (b *StatementBuilder) buildTemporalAggCumulativeOrUnspecified(
|
|
_ context.Context,
|
|
start, end uint64,
|
|
query qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation],
|
|
samplesTable string,
|
|
timeSeriesCTE string,
|
|
timeSeriesCTEArgs []any,
|
|
) (string, []any, error) {
|
|
stepSec := int64(query.StepInterval.Seconds())
|
|
|
|
baseSb := sqlbuilder.NewSelectBuilder()
|
|
baseSb.Select("fingerprint")
|
|
baseSb.SelectMore(fmt.Sprintf(
|
|
"toStartOfInterval(toDateTime(intDiv(unix_milli, 1000)), toIntervalSecond(%d)) AS ts",
|
|
stepSec,
|
|
))
|
|
for i, g := range query.GroupBy {
|
|
baseSb.SelectMore(sqlbuilder.Escape(GroupByColumnAlias(i, g.Name)))
|
|
}
|
|
|
|
aggCol, err := metricstelemetryschema.AggregationColumnForSamplesTable(samplesTable, query.Aggregations[0].Temporality, query.Aggregations[0].TimeAggregation)
|
|
if err != nil {
|
|
return "", nil, err
|
|
}
|
|
baseSb.SelectMore(fmt.Sprintf("%s AS per_series_value", aggCol))
|
|
|
|
baseSb.From(fmt.Sprintf("%s.%s AS points", metricstelemetryschema.DBName, samplesTable))
|
|
baseSb.JoinWithOption(sqlbuilder.InnerJoin, timeSeriesCTE, "points.fingerprint = filtered_time_series.fingerprint")
|
|
baseSb.Where(
|
|
baseSb.In("metric_name", query.Aggregations[0].MetricName),
|
|
baseSb.GTE("unix_milli", start),
|
|
baseSb.LT("unix_milli", end),
|
|
)
|
|
baseSb.GroupBy("fingerprint", "ts")
|
|
baseSb.GroupBy(GroupByAliases(query.GroupBy)...)
|
|
baseSb.OrderBy("fingerprint", "ts")
|
|
|
|
innerQuery, innerArgs := baseSb.BuildWithFlavor(sqlbuilder.ClickHouse, timeSeriesCTEArgs...)
|
|
|
|
switch query.Aggregations[0].TimeAggregation {
|
|
case metrictypes.TimeAggregationRate:
|
|
wrapped := sqlbuilder.NewSelectBuilder()
|
|
wrapped.Select("ts")
|
|
for i, g := range query.GroupBy {
|
|
wrapped.SelectMore(sqlbuilder.Escape(GroupByColumnAlias(i, g.Name)))
|
|
}
|
|
wrapped.SelectMore(fmt.Sprintf("%s AS per_series_value", RateTmpl))
|
|
wrapped.From(fmt.Sprintf("(%s) WINDOW rate_window AS (PARTITION BY fingerprint ORDER BY fingerprint, ts)", sqlbuilder.Escape(innerQuery)))
|
|
q, args := wrapped.BuildWithFlavor(sqlbuilder.ClickHouse, innerArgs...)
|
|
return fmt.Sprintf("__temporal_aggregation_cte AS (%s)", q), args, nil
|
|
|
|
case metrictypes.TimeAggregationIncrease:
|
|
wrapped := sqlbuilder.NewSelectBuilder()
|
|
wrapped.Select("ts")
|
|
for i, g := range query.GroupBy {
|
|
wrapped.SelectMore(sqlbuilder.Escape(GroupByColumnAlias(i, g.Name)))
|
|
}
|
|
wrapped.SelectMore(fmt.Sprintf("%s AS per_series_value", IncreaseTmpl))
|
|
wrapped.From(fmt.Sprintf("(%s) WINDOW rate_window AS (PARTITION BY fingerprint ORDER BY fingerprint, ts)", sqlbuilder.Escape(innerQuery)))
|
|
q, args := wrapped.BuildWithFlavor(sqlbuilder.ClickHouse, innerArgs...)
|
|
return fmt.Sprintf("__temporal_aggregation_cte AS (%s)", q), args, nil
|
|
default:
|
|
return fmt.Sprintf("__temporal_aggregation_cte AS (%s)", innerQuery), innerArgs, nil
|
|
}
|
|
}
|
|
|
|
func (b *StatementBuilder) buildTemporalAggForMultipleTemporalities(
|
|
_ context.Context,
|
|
start, end uint64,
|
|
query qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation],
|
|
samplesTable string,
|
|
timeSeriesCTE string,
|
|
timeSeriesCTEArgs []any,
|
|
) (string, []any, error) {
|
|
stepSec := int64(query.StepInterval.Seconds())
|
|
sb := sqlbuilder.NewSelectBuilder()
|
|
|
|
sb.SelectMore(fmt.Sprintf(
|
|
"toStartOfInterval(toDateTime(intDiv(unix_milli, 1000)), toIntervalSecond(%d)) AS ts",
|
|
stepSec,
|
|
))
|
|
for i, g := range query.GroupBy {
|
|
sb.SelectMore(sqlbuilder.Escape(GroupByColumnAlias(i, g.Name)))
|
|
}
|
|
|
|
aggForDeltaTemporality, err := metricstelemetryschema.AggregationColumnForSamplesTable(samplesTable, metrictypes.Delta, query.Aggregations[0].TimeAggregation)
|
|
if err != nil {
|
|
return "", nil, err
|
|
}
|
|
aggForCumulativeTemporality, err := metricstelemetryschema.AggregationColumnForSamplesTable(samplesTable, metrictypes.Cumulative, query.Aggregations[0].TimeAggregation)
|
|
if err != nil {
|
|
return "", nil, err
|
|
}
|
|
if query.Aggregations[0].TimeAggregation == metrictypes.TimeAggregationRate {
|
|
aggForDeltaTemporality = fmt.Sprintf("%s/%d", aggForDeltaTemporality, stepSec)
|
|
}
|
|
|
|
switch query.Aggregations[0].TimeAggregation {
|
|
case metrictypes.TimeAggregationRate:
|
|
rateExpr := fmt.Sprintf(RateMultiTemporalityTmpl,
|
|
aggForDeltaTemporality,
|
|
aggForCumulativeTemporality, aggForCumulativeTemporality, aggForCumulativeTemporality,
|
|
aggForCumulativeTemporality, aggForCumulativeTemporality,
|
|
)
|
|
sb.SelectMore(rateExpr)
|
|
case metrictypes.TimeAggregationIncrease:
|
|
increaseExpr := fmt.Sprintf(IncreaseMultiTemporality,
|
|
aggForDeltaTemporality,
|
|
aggForCumulativeTemporality, aggForCumulativeTemporality, aggForCumulativeTemporality,
|
|
aggForCumulativeTemporality, aggForCumulativeTemporality,
|
|
)
|
|
sb.SelectMore(increaseExpr)
|
|
default:
|
|
expr := fmt.Sprintf(OthersMultiTemporality, aggForDeltaTemporality, aggForCumulativeTemporality)
|
|
sb.SelectMore(expr)
|
|
}
|
|
|
|
sb.From(fmt.Sprintf("%s.%s AS points", metricstelemetryschema.DBName, samplesTable))
|
|
sb.JoinWithOption(sqlbuilder.InnerJoin, timeSeriesCTE, "points.fingerprint = filtered_time_series.fingerprint")
|
|
sb.Where(
|
|
sb.In("metric_name", query.Aggregations[0].MetricName),
|
|
sb.GTE("unix_milli", start),
|
|
sb.LT("unix_milli", end),
|
|
)
|
|
sb.GroupBy("fingerprint", "ts", "temporality")
|
|
sb.GroupBy(GroupByAliases(query.GroupBy)...)
|
|
queryWithoutWindow, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse, timeSeriesCTEArgs...)
|
|
queryWithWindowAndOrder := queryWithoutWindow + " WINDOW rate_window AS (PARTITION BY fingerprint ORDER BY fingerprint ASC, ts ASC) ORDER BY ts"
|
|
return fmt.Sprintf("__temporal_aggregation_cte AS (%s)", queryWithWindowAndOrder), args, nil
|
|
}
|
|
|
|
func (b *StatementBuilder) buildSpatialAggregationCTE(
|
|
_ context.Context,
|
|
_ uint64,
|
|
_ uint64,
|
|
query qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation],
|
|
_ map[string][]*telemetrytypes.TelemetryFieldKey,
|
|
) (string, []any) {
|
|
sb := sqlbuilder.NewSelectBuilder()
|
|
|
|
sb.Select("ts")
|
|
for i, g := range query.GroupBy {
|
|
sb.SelectMore(sqlbuilder.Escape(GroupByColumnAlias(i, g.Name)))
|
|
}
|
|
sb.SelectMore(fmt.Sprintf("%s(per_series_value) AS value", query.Aggregations[0].SpaceAggregation.StringValue()))
|
|
sb.From("__temporal_aggregation_cte")
|
|
sb.Where(sb.EQ("isNaN(per_series_value)", 0))
|
|
if query.Aggregations[0].ValueFilter != nil {
|
|
sb.Where(sb.EQ("per_series_value", query.Aggregations[0].ValueFilter.Value))
|
|
}
|
|
sb.GroupBy("ts")
|
|
sb.GroupBy(GroupByAliases(query.GroupBy)...)
|
|
|
|
q, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
|
|
return fmt.Sprintf("__spatial_aggregation_cte AS (%s)", q), args
|
|
}
|
|
|
|
func (b *StatementBuilder) BuildFinalSelect(
|
|
cteFragments []string,
|
|
cteArgs [][]any,
|
|
requestType qbtypes.RequestType,
|
|
query qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation],
|
|
) (*qbtypes.Statement, error) {
|
|
combined := querybuilder.CombineCTEs(cteFragments)
|
|
|
|
var args []any
|
|
for _, a := range cteArgs {
|
|
args = append(args, a...)
|
|
}
|
|
|
|
if requestType == qbtypes.RequestTypeHeatmap {
|
|
return buildHeatmapFinalSelect(combined, args, query)
|
|
}
|
|
return buildAggregationFinalSelect(combined, args, query)
|
|
}
|
|
|
|
// buildAggregationFinalSelect reads __spatial_aggregation_cte as one value per
|
|
// (group, timestamp), which is what every request type but heatmap wants.
|
|
func buildAggregationFinalSelect(
|
|
combined string,
|
|
args []any,
|
|
query qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation],
|
|
) (*qbtypes.Statement, error) {
|
|
metricType := query.Aggregations[0].Type
|
|
spaceAgg := query.Aggregations[0].SpaceAggregation
|
|
|
|
sb := sqlbuilder.NewSelectBuilder()
|
|
|
|
if metricType == metrictypes.HistogramType && spaceAgg.IsPercentile() {
|
|
quantile := query.Aggregations[0].SpaceAggregation.Percentile()
|
|
sb.Select("ts")
|
|
for i, g := range query.GroupBy {
|
|
sb.SelectMore(sqlbuilder.Escape(GroupByColumnAlias(i, g.Name)))
|
|
}
|
|
sb.SelectMore(fmt.Sprintf(
|
|
"histogramQuantile(arrayMap(x -> toFloat64(x), groupArray(le)), groupArray(value), %.3f) AS value",
|
|
quantile,
|
|
))
|
|
sb.From("__spatial_aggregation_cte")
|
|
sb.GroupBy(GroupByAliases(query.GroupBy)...)
|
|
sb.GroupBy("ts")
|
|
if query.Having != nil && query.Having.Expression != "" {
|
|
rewriter := querybuilder.NewHavingExpressionRewriter()
|
|
rewrittenExpr, err := rewriter.RewriteForMetrics(query.Having.Expression, query.Aggregations)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
sb.Having(rewrittenExpr)
|
|
}
|
|
} else if metricType == metrictypes.HistogramType && spaceAgg == metrictypes.SpaceAggregationCount && query.Aggregations[0].ComparisonSpaceAggregationParam != nil {
|
|
sb.Select("ts")
|
|
|
|
for i, g := range query.GroupBy {
|
|
sb.SelectMore(sqlbuilder.Escape(GroupByColumnAlias(i, g.Name)))
|
|
}
|
|
|
|
aggQuery, err := metricstelemetryschema.AggregationQueryForHistogramCountWithParams(query.Aggregations[0].ComparisonSpaceAggregationParam)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
sb.SelectMore(aggQuery)
|
|
|
|
sb.From("__spatial_aggregation_cte")
|
|
|
|
sb.GroupBy(GroupByAliases(query.GroupBy)...)
|
|
sb.GroupBy("ts")
|
|
|
|
if query.Having != nil && query.Having.Expression != "" {
|
|
rewriter := querybuilder.NewHavingExpressionRewriter()
|
|
rewrittenExpr, err := rewriter.RewriteForMetrics(query.Having.Expression, query.Aggregations)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
sb.Having(rewrittenExpr)
|
|
}
|
|
} else {
|
|
// for count aggregation on histograms with no params, the exact result of spatial aggregation can be sent forward
|
|
sb.Select("*")
|
|
sb.From("__spatial_aggregation_cte")
|
|
if query.Having != nil && query.Having.Expression != "" {
|
|
rewriter := querybuilder.NewHavingExpressionRewriter()
|
|
rewrittenExpr, err := rewriter.RewriteForMetrics(query.Having.Expression, query.Aggregations)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
sb.Where(rewrittenExpr)
|
|
}
|
|
}
|
|
sb.OrderBy(GroupByAliases(query.GroupBy)...)
|
|
sb.OrderBy("ts")
|
|
if metricType == metrictypes.HistogramType && spaceAgg == metrictypes.SpaceAggregationCount && query.Aggregations[0].ComparisonSpaceAggregationParam == nil {
|
|
sb.OrderBy("toFloat64(le)")
|
|
}
|
|
|
|
q, a := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
|
|
return &qbtypes.Statement{Query: combined + q, Args: append(args, a...)}, nil
|
|
}
|
|
|
|
const (
|
|
histogramBucketKey = "le"
|
|
|
|
heatmapValueAlias = "__result_0"
|
|
heatmapWindow = "__heatmap_window"
|
|
)
|
|
|
|
func isHistogramBucket(k qbtypes.GroupByKey) bool { return k.Name == histogramBucketKey }
|
|
|
|
// buildHeatmapFinalSelect turns __spatial_aggregation_cte into one row per
|
|
// heatmap cell: (ts, group labels..., bucket upper bound, count).
|
|
func buildHeatmapFinalSelect(
|
|
combined string,
|
|
args []any,
|
|
query qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation],
|
|
) (*qbtypes.Statement, error) {
|
|
if query.Aggregations[0].Type == metrictypes.HistogramType {
|
|
return buildHistogramHeatmapFinalSelect(combined, args, query)
|
|
}
|
|
return buildValueHeatmapFinalSelect(combined, args, query)
|
|
}
|
|
|
|
// buildHistogramHeatmapFinalSelect differences the cumulative per-`le` counts in
|
|
// __spatial_aggregation_cte into a count per bucket. A bucket runs from the `le`
|
|
// below it up to its own, so the `le=+Inf` row reaches the reader as the
|
|
// overflow and the lowest `le` as a bucket open below, which is where a
|
|
// negative observation would have been counted.
|
|
func buildHistogramHeatmapFinalSelect(
|
|
combined string,
|
|
args []any,
|
|
query qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation],
|
|
) (*qbtypes.Statement, error) {
|
|
groupAliases := GroupByAliases(query.GroupBy)
|
|
partitionBy := append(append([]string{}, groupAliases...), "ts")
|
|
|
|
sb := sqlbuilder.NewSelectBuilder()
|
|
sb.Select("ts")
|
|
sb.SelectMore(groupAliases...)
|
|
sb.SelectMore(fmt.Sprintf(
|
|
"lagInFrame(toFloat64(%s), 1, toFloat64('-Inf')) OVER %s AS %s",
|
|
histogramBucketKey, heatmapWindow, qbtypes.HeatmapBucketMinColumn,
|
|
))
|
|
sb.SelectMore(fmt.Sprintf("toFloat64(%s) AS %s", histogramBucketKey, qbtypes.HeatmapBucketMaxColumn))
|
|
// a partial scrape can break monotonicity across `le`, and a negative cell
|
|
// count has no meaning
|
|
sb.SelectMore(fmt.Sprintf(
|
|
"greatest(value - lagInFrame(value, 1, 0) OVER %s, 0) AS %s",
|
|
heatmapWindow, heatmapValueAlias,
|
|
))
|
|
// sqlbuilder has no WINDOW clause; appending it to FROM lands it between FROM
|
|
// and ORDER BY, since these statements carry no WHERE or GROUP BY
|
|
sb.From(fmt.Sprintf(
|
|
"__spatial_aggregation_cte WINDOW %s AS (PARTITION BY %s ORDER BY toFloat64(%s))",
|
|
heatmapWindow, strings.Join(partitionBy, ", "), histogramBucketKey,
|
|
))
|
|
sb.OrderBy(groupAliases...)
|
|
sb.OrderBy("ts", fmt.Sprintf("toFloat64(%s)", histogramBucketKey))
|
|
|
|
q, a := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
|
|
return &qbtypes.Statement{Query: combined + q, Args: append(args, a...)}, nil
|
|
}
|
|
|
|
// buildValueHeatmapFinalSelect places each spatially aggregated value in a
|
|
// bucket of the requested axis. __spatial_aggregation_cte holds one row per
|
|
// (group, timestamp), so every cell counts exactly one.
|
|
func buildValueHeatmapFinalSelect(
|
|
combined string,
|
|
args []any,
|
|
query qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation],
|
|
) (*qbtypes.Statement, error) {
|
|
bucketMin, bucketMax, err := renderHeatmapBucketExprs(*query.Aggregations[0].HeatmapBucketing)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
groupAliases := GroupByAliases(query.GroupBy)
|
|
|
|
sb := sqlbuilder.NewSelectBuilder()
|
|
sb.Select("ts")
|
|
sb.SelectMore(groupAliases...)
|
|
sb.SelectMore(fmt.Sprintf("%s AS %s", bucketMin, qbtypes.HeatmapBucketMinColumn))
|
|
sb.SelectMore(fmt.Sprintf("%s AS %s", bucketMax, qbtypes.HeatmapBucketMaxColumn))
|
|
sb.SelectMore(fmt.Sprintf("toFloat64(1) AS %s", heatmapValueAlias))
|
|
sb.From("__spatial_aggregation_cte")
|
|
sb.OrderBy(groupAliases...)
|
|
sb.OrderBy("ts", qbtypes.HeatmapBucketMaxColumn)
|
|
|
|
q, a := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
|
|
return &qbtypes.Statement{Query: combined + q, Args: append(args, a...)}, nil
|
|
}
|
|
|
|
// renderHeatmapBucketExprs renders the bucket (min, max] that `value` falls in.
|
|
// Only the bucket under everything the axis covers is open below, and only the
|
|
// one over it is open above.
|
|
func renderHeatmapBucketExprs(bucketing qbtypes.HeatmapBucketing) (minExpr, maxExpr string, err error) {
|
|
switch bucketing.Kind {
|
|
case qbtypes.BucketsKindLinear:
|
|
return renderLinearBucketExprs(bucketing)
|
|
case qbtypes.BucketsKindLog:
|
|
return renderLogBucketExprs()
|
|
default:
|
|
return "", "", errors.NewInvalidInputf(errors.CodeInvalidInput,
|
|
"unsupported bucketsScaling %q for heatmap requests", bucketing.Kind.StringValue())
|
|
}
|
|
}
|
|
|
|
func renderLinearBucketExprs(bucketing qbtypes.HeatmapBucketing) (string, string, error) {
|
|
maxValue := formatFloat(bucketing.MaxValue)
|
|
numBuckets := strconv.Itoa(bucketing.NumBuckets)
|
|
index := fmt.Sprintf("least(greatest(ceil(value * %s / %s), 1), %s)", numBuckets, maxValue, numBuckets)
|
|
|
|
minExpr := fmt.Sprintf(
|
|
"multiIf(value <= 0, toFloat64('-Inf'), value > %s, toFloat64(%s), (%s - 1) * %s / %s)",
|
|
maxValue, maxValue, index, maxValue, numBuckets,
|
|
)
|
|
maxExpr := fmt.Sprintf(
|
|
"multiIf(value <= 0, toFloat64(0), value > %s, toFloat64('+Inf'), %s * %s / %s)",
|
|
maxValue, index, maxValue, numBuckets,
|
|
)
|
|
return minExpr, maxExpr, nil
|
|
}
|
|
|
|
// ClickHouse buckets at MaxLogScale whatever HeatmapBucketing.LogScale asks for;
|
|
// postprocessing folds the axis down afterwards.
|
|
func renderLogBucketExprs() (string, string, error) {
|
|
bucketsPerDoubling := formatFloat(math.Exp2(qbtypes.MaxLogScale))
|
|
lowest := formatFloat(qbtypes.MinLogUpperBound)
|
|
highest := formatFloat(qbtypes.MaxLogUpperBound)
|
|
|
|
minExpr := fmt.Sprintf(
|
|
"multiIf(value <= 0, toFloat64('-Inf'), value <= %s, toFloat64(0), value > %s, toFloat64(%s), pow(2, (ceil(log2(value) * %s) - 1) / %s))",
|
|
lowest, highest, highest, bucketsPerDoubling, bucketsPerDoubling,
|
|
)
|
|
maxExpr := fmt.Sprintf(
|
|
"multiIf(value <= 0, toFloat64(0), value <= %s, %s, value > %s, toFloat64('+Inf'), pow(2, ceil(log2(value) * %s) / %s))",
|
|
lowest, lowest, highest, bucketsPerDoubling, bucketsPerDoubling,
|
|
)
|
|
return minExpr, maxExpr, nil
|
|
}
|
|
|
|
// formatFloat renders a float64 as the shortest literal that reads back as the
|
|
// same value, so an upper bound computed from it is identical on every row.
|
|
func formatFloat(v float64) string {
|
|
return strconv.FormatFloat(v, 'g', -1, 64)
|
|
}
|
|
|
|
func GroupByColumnAlias(i int, name string) string {
|
|
if name == histogramBucketKey {
|
|
return clickhousesql.Identifier(histogramBucketKey)
|
|
}
|
|
return clickhousesql.Identifier(fmt.Sprintf("__GROUP_BY_KEY_%d_%s", i, name))
|
|
}
|
|
|
|
func GroupByAliases(groupBy []qbtypes.GroupByKey) []string {
|
|
aliases := make([]string, 0, len(groupBy))
|
|
for i := range groupBy {
|
|
aliases = append(aliases, sqlbuilder.Escape(GroupByColumnAlias(i, groupBy[i].Name)))
|
|
}
|
|
return aliases
|
|
}
|