| 191 | } |
| 192 | |
| 193 | func 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. |
| 225 | func (d *Distributor) queryIngesterStream(ctx context.Context, replicationSet ring.ReplicationSet, req *ingester_client.QueryRequest, partialDataEnabled bool) (*ingester_client.QueryStreamResponse, error) { |