mirror of
https://github.com/SigNoz/signoz.git
synced 2026-09-19 01:40: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
420 lines
16 KiB
Go
420 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
|
|
storage qbtypes.Storage
|
|
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 storage 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
|
|
}
|
|
return NewMeterQueryStatementBuilder(settings, metadataStore, metricstelemetryschema.NewStorage(), metricsStatementBuilder), nil
|
|
},
|
|
)
|
|
}
|
|
|
|
func NewMeterQueryStatementBuilder(
|
|
settings factory.ProviderSettings,
|
|
metadataStore telemetrytypes.MetadataStore,
|
|
storage qbtypes.Storage,
|
|
metricsStatementBuilder *metricsstatementbuilder.StatementBuilder,
|
|
) *meterQueryStatementBuilder {
|
|
metricsSettings := factory.NewScopedProviderSettings(settings, "github.com/SigNoz/signoz/pkg/telemetryschema/metertelemetryschema")
|
|
|
|
return &meterQueryStatementBuilder{
|
|
logger: metricsSettings.Logger(),
|
|
metadataStore: metadataStore,
|
|
storage: storage,
|
|
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, qbtypes.RequestTypeTimeSeries, 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,
|
|
))
|
|
info := querybuilder.NewQueryInfo(ctx, orgID, nil, telemetrytypes.SignalMetrics, &telemetrytypes.MetricContext{MetricName: query.Aggregations[0].MetricName}, start, end)
|
|
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, metricsstatementbuilder.GroupByColumnAlias(i, g.Name))))
|
|
}
|
|
|
|
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,
|
|
Query: info,
|
|
Storage: b.storage,
|
|
Logger: b.logger,
|
|
FieldKeys: keys,
|
|
FullTextColumn: &telemetrytypes.TelemetryFieldKey{Name: "labels"},
|
|
Variables: variables,
|
|
})
|
|
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(metricsstatementbuilder.GroupByAliases(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,
|
|
))
|
|
|
|
info := querybuilder.NewQueryInfo(ctx, orgID, nil, telemetrytypes.SignalMetrics, &telemetrytypes.MetricContext{MetricName: query.Aggregations[0].MetricName}, start, end)
|
|
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, metricsstatementbuilder.GroupByColumnAlias(i, g.Name))))
|
|
}
|
|
|
|
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,
|
|
Query: info,
|
|
Storage: b.storage,
|
|
Logger: b.logger,
|
|
FieldKeys: keys,
|
|
FullTextColumn: &telemetrytypes.TelemetryFieldKey{Name: "labels"},
|
|
Variables: variables,
|
|
})
|
|
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(metricsstatementbuilder.GroupByAliases(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,
|
|
))
|
|
info := querybuilder.NewQueryInfo(ctx, orgID, nil, telemetrytypes.SignalMetrics, &telemetrytypes.MetricContext{MetricName: query.Aggregations[0].MetricName}, start, end)
|
|
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
|
|
}
|
|
baseSb.SelectMore(sqlbuilder.Escape(fmt.Sprintf("%s AS %s", col, metricsstatementbuilder.GroupByColumnAlias(i, g.Name))))
|
|
}
|
|
|
|
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,
|
|
Query: info,
|
|
Storage: b.storage,
|
|
Logger: b.logger,
|
|
FieldKeys: keys,
|
|
FullTextColumn: &telemetrytypes.TelemetryFieldKey{Name: "labels"},
|
|
Variables: variables,
|
|
})
|
|
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(metricsstatementbuilder.GroupByAliases(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 i, g := range query.GroupBy {
|
|
wrapped.SelectMore(sqlbuilder.Escape(metricsstatementbuilder.GroupByColumnAlias(i, 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 i, g := range query.GroupBy {
|
|
wrapped.SelectMore(sqlbuilder.Escape(metricsstatementbuilder.GroupByColumnAlias(i, 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 i, g := range query.GroupBy {
|
|
sb.SelectMore(sqlbuilder.Escape(metricsstatementbuilder.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(metricsstatementbuilder.GroupByAliases(query.GroupBy)...)
|
|
|
|
q, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
|
|
return fmt.Sprintf("__spatial_aggregation_cte AS (%s)", q), args, nil
|
|
}
|