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 }