MCPcopy Create free account
hub / github.com/cortexproject/cortex / mergeExemplarQueryResponses

Function mergeExemplarQueryResponses

pkg/distributor/query.go:193–222  ·  view source on GitHub ↗
(results []any)

Source from the content-addressed store, hash-verified

191}
192
193func mergeExemplarQueryResponses(results []any) *ingester_client.ExemplarQueryResponse {
194 var keys []string
195 exemplarResults := make(map[string]cortexpb.TimeSeries)
196 buf := make([]byte, 0, 1024)
197 for _, result := range results {
198 r := result.(*ingester_client.ExemplarQueryResponse)
199 for _, ts := range r.Timeseries {
200 lbls := string(cortexpb.FromLabelAdaptersToLabels(ts.Labels).Bytes(buf))
201 e, ok := exemplarResults[lbls]
202 if !ok {
203 exemplarResults[lbls] = ts
204 keys = append(keys, lbls)
205 } else {
206 // Merge in any missing values from another ingesters exemplars for this series.
207 e.Exemplars = mergeExemplarSets(e.Exemplars, ts.Exemplars)
208 exemplarResults[lbls] = e
209 }
210 }
211 }
212
213 // Query results from each ingester were sorted, but are not necessarily still sorted after merging.
214 sort.Strings(keys)
215
216 result := make([]cortexpb.TimeSeries, len(exemplarResults))
217 for i, k := range keys {
218 result[i] = exemplarResults[k]
219 }
220
221 return &ingester_client.ExemplarQueryResponse{Timeseries: result}
222}
223
224// queryIngesterStream queries the ingesters using the new streaming API.
225func (d *Distributor) queryIngesterStream(ctx context.Context, replicationSet ring.ReplicationSet, req *ingester_client.QueryRequest, partialDataEnabled bool) (*ingester_client.QueryStreamResponse, error) {

Callers 2

TestMergeExemplarsFunction · 0.85

Calls 3

mergeExemplarSetsFunction · 0.85
BytesMethod · 0.65

Tested by 1

TestMergeExemplarsFunction · 0.68