forked from thanos-io/thanos
-
Notifications
You must be signed in to change notification settings - Fork 2
/
downsampled.go
110 lines (98 loc) · 2.99 KB
/
downsampled.go
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
// Copyright (c) The Thanos Authors.
// Licensed under the Apache License 2.0.
package queryfrontend
import (
"context"
"github.com/prometheus/client_golang/prometheus"
"github.com/prometheus/client_golang/prometheus/promauto"
"github.com/thanos-io/thanos/internal/cortex/querier/queryrange"
"github.com/thanos-io/thanos/pkg/compact/downsample"
)
// DownsampledMiddleware creates a new Middleware that requests downsampled data
// should response to original request with auto max_source_resolution not contain data points.
func DownsampledMiddleware(merger queryrange.Merger, registerer prometheus.Registerer) queryrange.Middleware {
return queryrange.MiddlewareFunc(func(next queryrange.Handler) queryrange.Handler {
return downsampled{
next: next,
merger: merger,
additionalQueriesCount: promauto.With(registerer).NewCounter(prometheus.CounterOpts{
Namespace: "thanos",
Name: "frontend_downsampled_extra_queries_total",
Help: "Total number of additional queries for downsampled data",
}),
}
})
}
type downsampled struct {
next queryrange.Handler
merger queryrange.Merger
// Metrics.
additionalQueriesCount prometheus.Counter
}
var resolutions = []int64{downsample.ResLevel1, downsample.ResLevel2}
func (d downsampled) Do(ctx context.Context, req queryrange.Request) (queryrange.Response, error) {
tqrr, ok := req.(*ThanosQueryRangeRequest)
if !ok || !tqrr.AutoDownsampling {
return d.next.Do(ctx, req)
}
var (
resps = make([]queryrange.Response, 0)
resp queryrange.Response
err error
i int
)
forLoop:
for i < len(resolutions) {
if i > 0 {
d.additionalQueriesCount.Inc()
}
r := *tqrr
resp, err = d.next.Do(ctx, &r)
if err != nil {
return nil, err
}
resps = append(resps, resp)
// Set MaxSourceResolution for next request, if any.
for i < len(resolutions) {
if tqrr.MaxSourceResolution < resolutions[i] {
tqrr.AutoDownsampling = false
tqrr.MaxSourceResolution = resolutions[i]
break
}
i++
}
m := minResponseTime(resp)
switch m {
case tqrr.Start: // Response not impacted by retention policy.
break forLoop
case -1: // Empty response, retry with higher MaxSourceResolution.
continue
default: // Data partially present, query for empty part with higher MaxSourceResolution.
tqrr.End = m - tqrr.Step
}
if tqrr.Start > tqrr.End {
break forLoop
}
}
response, err := d.merger.MergeResponse(req, resps...)
if err != nil {
return nil, err
}
return response, nil
}
// minResponseTime returns earliest timestamp in r.Data.Result.
// -1 is returned if r contains no data points.
// Each SampleStream within r.Data.Result must be sorted by timestamp.
func minResponseTime(r queryrange.Response) int64 {
var res = r.(*queryrange.PrometheusResponse).Data.Result
if len(res) == 0 || len(res[0].Samples) == 0 {
return -1
}
var minTs = res[0].Samples[0].TimestampMs
for _, sampleStream := range res[1:] {
if ts := sampleStream.Samples[0].TimestampMs; ts < minTs {
minTs = ts
}
}
return minTs
}