mirror of
https://github.com/SigNoz/signoz.git
synced 2026-08-01 18:50:35 +01:00
Remove the statementbuilder.Builders aggregate. The querier chain (querier.New, signozquerier.NewFactory, NewQuerierProviderFactories) now takes the six per-signal statement builders individually, and signoz.go::newQueryStack returns them directly via a plain multi-return. The statementbuilder package is now config-only. Remove each <signal>statementbuilder/new.go and fold its NewFactory into that package's statement_builder.go (traces' NewOperatorFactory into trace_operator_statement_builder.go), giving the layout: struct -> NewFactory -> New<X>QueryStatementBuilder.
431 lines
16 KiB
Go
431 lines
16 KiB
Go
package meterstatementbuilder
|
|
|
|
import (
|
|
"context"
|
|
"fmt"
|
|
"log/slog"
|
|
|
|
"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/statementbuilder/metricsstatementbuilder"
|
|
"github.com/SigNoz/signoz/pkg/telemetryschema/metertelemetryschema"
|
|
"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"
|
|
)
|
|
|
|
type meterQueryStatementBuilder struct {
|
|
logger *slog.Logger
|
|
metadataStore telemetrytypes.MetadataStore
|
|
fm qbtypes.FieldMapper
|
|
cb qbtypes.ConditionBuilder
|
|
metricsStatementBuilder *metricsstatementbuilder.StatementBuilder
|
|
}
|
|
|
|
var _ qbtypes.StatementBuilder[qbtypes.MetricAggregation] = (*meterQueryStatementBuilder)(nil)
|
|
|
|
// NewFactory returns a provider factory for the meter statement builder. Its New
|
|
// reuses the metrics FieldMapper/ConditionBuilder and delegates the final SELECT
|
|
// to a metrics statement builder built via the metrics factory.
|
|
func NewFactory(
|
|
metadataStore telemetrytypes.MetadataStore,
|
|
fl flagger.Flagger,
|
|
) factory.ProviderFactory[qbtypes.StatementBuilder[qbtypes.MetricAggregation], statementbuilder.Config] {
|
|
return factory.NewProviderFactory(
|
|
factory.MustNewName("meter"),
|
|
func(ctx context.Context, settings factory.ProviderSettings, cfg statementbuilder.Config) (qbtypes.StatementBuilder[qbtypes.MetricAggregation], error) {
|
|
metricsStatementBuilder, err := metricsstatementbuilder.NewFactory(metadataStore, fl).New(ctx, settings, cfg)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
fm := metricstelemetryschema.NewFieldMapper()
|
|
cb := metricstelemetryschema.NewConditionBuilder(fm)
|
|
return NewMeterQueryStatementBuilder(settings, metadataStore, fm, cb, metricsStatementBuilder), nil
|
|
},
|
|
)
|
|
}
|
|
|
|
func NewMeterQueryStatementBuilder(
|
|
settings factory.ProviderSettings,
|
|
metadataStore telemetrytypes.MetadataStore,
|
|
fieldMapper qbtypes.FieldMapper,
|
|
conditionBuilder qbtypes.ConditionBuilder,
|
|
metricsStatementBuilder *metricsstatementbuilder.StatementBuilder,
|
|
) *meterQueryStatementBuilder {
|
|
metricsSettings := factory.NewScopedProviderSettings(settings, "github.com/SigNoz/signoz/pkg/telemetryschema/metertelemetryschema")
|
|
|
|
return &meterQueryStatementBuilder{
|
|
logger: metricsSettings.Logger(),
|
|
metadataStore: metadataStore,
|
|
fm: fieldMapper,
|
|
cb: conditionBuilder,
|
|
metricsStatementBuilder: metricsStatementBuilder,
|
|
}
|
|
}
|
|
|
|
func (b *meterQueryStatementBuilder) Build(
|
|
ctx context.Context,
|
|
orgID valuer.UUID,
|
|
start uint64,
|
|
end uint64,
|
|
_ qbtypes.RequestType,
|
|
query qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation],
|
|
variables map[string]qbtypes.VariableItem,
|
|
) (*qbtypes.Statement, error) {
|
|
keySelectors := metricsstatementbuilder.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, query, keys, variables)
|
|
}
|
|
|
|
func (b *meterQueryStatementBuilder) buildPipelineStatement(
|
|
ctx context.Context,
|
|
orgID valuer.UUID,
|
|
start, end uint64,
|
|
query qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation],
|
|
keys map[string][]*telemetrytypes.TelemetryFieldKey,
|
|
variables map[string]qbtypes.VariableItem,
|
|
) (*qbtypes.Statement, error) {
|
|
var (
|
|
cteFragments []string
|
|
cteArgs [][]any
|
|
)
|
|
|
|
if qbtypes.CanShortCircuitDelta(query.Aggregations[0]) {
|
|
// spatial_aggregation_cte directly for certain delta queries
|
|
if frag, args, err := b.buildTemporalAggDeltaFastPath(ctx, orgID, start, end, query, keys, variables); 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, orgID, start, end, query, keys, variables); err != nil {
|
|
return nil, err
|
|
} else if frag != "" {
|
|
cteFragments = append(cteFragments, frag)
|
|
cteArgs = append(cteArgs, args)
|
|
}
|
|
|
|
// spatial_aggregation_cte
|
|
if frag, args, err := b.buildSpatialAggregationCTE(ctx, start, end, query, keys); err != nil {
|
|
return nil, err
|
|
} else if frag != "" {
|
|
cteFragments = append(cteFragments, frag)
|
|
cteArgs = append(cteArgs, args)
|
|
}
|
|
}
|
|
|
|
// final SELECT
|
|
return b.metricsStatementBuilder.BuildFinalSelect(cteFragments, cteArgs, query)
|
|
}
|
|
|
|
func (b *meterQueryStatementBuilder) buildTemporalAggDeltaFastPath(
|
|
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) {
|
|
var filterWhere querybuilder.PreparedWhereClause
|
|
var err 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 _, g := range query.GroupBy {
|
|
col, err := b.fm.ColumnExpressionFor(ctx, orgID, start, end, &g.TelemetryFieldKey, telemetrytypes.FieldDataTypeString, keys)
|
|
if err != nil {
|
|
return "", nil, err
|
|
}
|
|
sb.SelectMore(col)
|
|
}
|
|
|
|
tbl := metertelemetryschema.WhichSamplesTableToUse(start, end, query.Aggregations[0].Type, query.Aggregations[0].TimeAggregation, query.Aggregations[0].TableHints)
|
|
aggCol, err := metertelemetryschema.AggregationColumnForSamplesTable(start, end, query.Aggregations[0].Type, query.Aggregations[0].Temporality, query.Aggregations[0].TimeAggregation, query.Aggregations[0].TableHints)
|
|
if err != nil {
|
|
return "", nil, err
|
|
}
|
|
if query.Aggregations[0].TimeAggregation == metrictypes.TimeAggregationRate {
|
|
aggCol = fmt.Sprintf("%s/%d", aggCol, stepSec)
|
|
}
|
|
|
|
sb.SelectMore(fmt.Sprintf("%s AS value", aggCol))
|
|
sb.From(fmt.Sprintf("%s.%s AS points", metertelemetryschema.DBName, tbl))
|
|
sb.Where(
|
|
sb.In("metric_name", query.Aggregations[0].MetricName),
|
|
sb.GTE("unix_milli", start),
|
|
sb.LT("unix_milli", end),
|
|
)
|
|
if query.Filter != nil && query.Filter.Expression != "" {
|
|
filterWhere, err = querybuilder.PrepareWhereClause(query.Filter.Expression, querybuilder.FilterExprVisitorOpts{
|
|
Context: ctx,
|
|
OrgID: orgID,
|
|
Logger: b.logger,
|
|
FieldMapper: b.fm,
|
|
ConditionBuilder: b.cb,
|
|
FieldKeys: keys,
|
|
FullTextColumn: &telemetrytypes.TelemetryFieldKey{Name: "labels"},
|
|
Variables: variables,
|
|
StartNs: start,
|
|
EndNs: end,
|
|
})
|
|
if err != nil {
|
|
return "", nil, err
|
|
}
|
|
}
|
|
if !filterWhere.IsEmpty() {
|
|
sb.AddWhereClause(filterWhere.WhereClause)
|
|
}
|
|
|
|
if query.Aggregations[0].Temporality != metrictypes.Unknown {
|
|
sb.Where(sb.ILike("temporality", query.Aggregations[0].Temporality.StringValue()))
|
|
}
|
|
sb.GroupBy("ts")
|
|
sb.GroupBy(querybuilder.GroupByKeys(query.GroupBy)...)
|
|
|
|
q, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
|
|
return fmt.Sprintf("__spatial_aggregation_cte AS (%s)", q), args, nil
|
|
}
|
|
|
|
func (b *meterQueryStatementBuilder) buildTemporalAggregationCTE(
|
|
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) {
|
|
if query.Aggregations[0].Temporality == metrictypes.Delta {
|
|
return b.buildTemporalAggDelta(ctx, orgID, start, end, query, keys, variables)
|
|
}
|
|
return b.buildTemporalAggCumulativeOrUnspecified(ctx, orgID, start, end, query, keys, variables)
|
|
}
|
|
|
|
func (b *meterQueryStatementBuilder) buildTemporalAggDelta(
|
|
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) {
|
|
var filterWhere querybuilder.PreparedWhereClause
|
|
var err 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 _, g := range query.GroupBy {
|
|
col, err := b.fm.ColumnExpressionFor(ctx, orgID, start, end, &g.TelemetryFieldKey, telemetrytypes.FieldDataTypeString, keys)
|
|
if err != nil {
|
|
return "", nil, err
|
|
}
|
|
sb.SelectMore(col)
|
|
}
|
|
|
|
tbl := metertelemetryschema.WhichSamplesTableToUse(start, end, query.Aggregations[0].Type, query.Aggregations[0].TimeAggregation, query.Aggregations[0].TableHints)
|
|
aggCol, err := metertelemetryschema.AggregationColumnForSamplesTable(start, end, query.Aggregations[0].Type, query.Aggregations[0].Temporality,
|
|
query.Aggregations[0].TimeAggregation, query.Aggregations[0].TableHints)
|
|
if err != nil {
|
|
return "", nil, err
|
|
}
|
|
if query.Aggregations[0].TimeAggregation == metrictypes.TimeAggregationRate {
|
|
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", metertelemetryschema.DBName, tbl))
|
|
sb.Where(
|
|
sb.In("metric_name", query.Aggregations[0].MetricName),
|
|
sb.GTE("unix_milli", start),
|
|
sb.LT("unix_milli", end),
|
|
)
|
|
|
|
if query.Filter != nil && query.Filter.Expression != "" {
|
|
filterWhere, err = querybuilder.PrepareWhereClause(query.Filter.Expression, querybuilder.FilterExprVisitorOpts{
|
|
Context: ctx,
|
|
OrgID: orgID,
|
|
Logger: b.logger,
|
|
FieldMapper: b.fm,
|
|
ConditionBuilder: b.cb,
|
|
FieldKeys: keys,
|
|
FullTextColumn: &telemetrytypes.TelemetryFieldKey{Name: "labels"},
|
|
Variables: variables,
|
|
StartNs: start,
|
|
EndNs: end,
|
|
})
|
|
if err != nil {
|
|
return "", nil, err
|
|
}
|
|
}
|
|
if !filterWhere.IsEmpty() {
|
|
sb.AddWhereClause(filterWhere.WhereClause)
|
|
}
|
|
|
|
if query.Aggregations[0].Temporality != metrictypes.Unknown {
|
|
sb.Where(sb.ILike("temporality", query.Aggregations[0].Temporality.StringValue()))
|
|
}
|
|
|
|
sb.GroupBy("fingerprint", "ts")
|
|
sb.GroupBy(querybuilder.GroupByKeys(query.GroupBy)...)
|
|
sb.OrderBy("fingerprint", "ts")
|
|
|
|
q, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
|
|
return fmt.Sprintf("__temporal_aggregation_cte AS (%s)", q), args, nil
|
|
}
|
|
|
|
func (b *meterQueryStatementBuilder) buildTemporalAggCumulativeOrUnspecified(
|
|
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) {
|
|
var filterWhere querybuilder.PreparedWhereClause
|
|
var err 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 _, g := range query.GroupBy {
|
|
col, err := b.fm.ColumnExpressionFor(ctx, orgID, start, end, &g.TelemetryFieldKey, telemetrytypes.FieldDataTypeString, keys)
|
|
if err != nil {
|
|
return "", nil, err
|
|
}
|
|
baseSb.SelectMore(col)
|
|
}
|
|
|
|
tbl := metertelemetryschema.WhichSamplesTableToUse(start, end, query.Aggregations[0].Type, query.Aggregations[0].TimeAggregation, query.Aggregations[0].TableHints)
|
|
aggCol, err := metertelemetryschema.AggregationColumnForSamplesTable(start, end, query.Aggregations[0].Type, query.Aggregations[0].Temporality, query.Aggregations[0].TimeAggregation, query.Aggregations[0].TableHints)
|
|
if err != nil {
|
|
return "", nil, err
|
|
}
|
|
baseSb.SelectMore(fmt.Sprintf("%s AS per_series_value", aggCol))
|
|
|
|
baseSb.From(fmt.Sprintf("%s.%s AS points", metertelemetryschema.DBName, tbl))
|
|
baseSb.Where(
|
|
baseSb.In("metric_name", query.Aggregations[0].MetricName),
|
|
baseSb.GTE("unix_milli", start),
|
|
baseSb.LT("unix_milli", end),
|
|
)
|
|
if query.Filter != nil && query.Filter.Expression != "" {
|
|
filterWhere, err = querybuilder.PrepareWhereClause(query.Filter.Expression, querybuilder.FilterExprVisitorOpts{
|
|
Context: ctx,
|
|
OrgID: orgID,
|
|
Logger: b.logger,
|
|
FieldMapper: b.fm,
|
|
ConditionBuilder: b.cb,
|
|
FieldKeys: keys,
|
|
FullTextColumn: &telemetrytypes.TelemetryFieldKey{Name: "labels"},
|
|
Variables: variables,
|
|
StartNs: start,
|
|
EndNs: end,
|
|
})
|
|
if err != nil {
|
|
return "", nil, err
|
|
}
|
|
}
|
|
if !filterWhere.IsEmpty() {
|
|
baseSb.AddWhereClause(filterWhere.WhereClause)
|
|
}
|
|
|
|
if query.Aggregations[0].Temporality != metrictypes.Unknown {
|
|
baseSb.Where(baseSb.ILike("temporality", query.Aggregations[0].Temporality.StringValue()))
|
|
}
|
|
baseSb.GroupBy("fingerprint", "ts")
|
|
baseSb.GroupBy(querybuilder.GroupByKeys(query.GroupBy)...)
|
|
baseSb.OrderBy("fingerprint", "ts")
|
|
|
|
innerQuery, innerArgs := baseSb.BuildWithFlavor(sqlbuilder.ClickHouse)
|
|
|
|
switch query.Aggregations[0].TimeAggregation {
|
|
case metrictypes.TimeAggregationRate:
|
|
wrapped := sqlbuilder.NewSelectBuilder()
|
|
wrapped.Select("ts")
|
|
for _, g := range query.GroupBy {
|
|
wrapped.SelectMore(fmt.Sprintf("`%s`", g.Name))
|
|
}
|
|
wrapped.SelectMore(fmt.Sprintf("%s AS per_series_value", metricsstatementbuilder.RateTmpl))
|
|
wrapped.From(fmt.Sprintf("(%s) WINDOW rate_window AS (PARTITION BY fingerprint ORDER BY fingerprint, ts)", 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 _, g := range query.GroupBy {
|
|
wrapped.SelectMore(fmt.Sprintf("`%s`", g.Name))
|
|
}
|
|
wrapped.SelectMore(fmt.Sprintf("%s AS per_series_value", metricsstatementbuilder.IncreaseTmpl))
|
|
wrapped.From(fmt.Sprintf("(%s) WINDOW rate_window AS (PARTITION BY fingerprint ORDER BY fingerprint, ts)", 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 *meterQueryStatementBuilder) buildSpatialAggregationCTE(
|
|
_ context.Context,
|
|
_ uint64,
|
|
_ uint64,
|
|
query qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation],
|
|
_ map[string][]*telemetrytypes.TelemetryFieldKey,
|
|
) (string, []any, error) {
|
|
|
|
if query.Aggregations[0].SpaceAggregation.IsZero() {
|
|
return "", nil, errors.Newf(
|
|
errors.TypeInvalidInput,
|
|
errors.CodeInvalidInput,
|
|
"invalid space aggregation, should be one of the following: [`sum`, `avg`, `min`, `max`, `count`]",
|
|
)
|
|
}
|
|
sb := sqlbuilder.NewSelectBuilder()
|
|
|
|
sb.Select("ts")
|
|
for _, g := range query.GroupBy {
|
|
sb.SelectMore(fmt.Sprintf("`%s`", 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(querybuilder.GroupByKeys(query.GroupBy)...)
|
|
|
|
q, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
|
|
return fmt.Sprintf("__spatial_aggregation_cte AS (%s)", q), args, nil
|
|
}
|