Range Splitting: Process instant queries as an independent query group (#64049)

* Query splitting: enable instant queries

* Range splitting: send instant queries as another request group

* Range splitting: increase grouped splitted requests stability

We were defaulting to the `0` index as the first group for the next request batch, but there was no guarantee that the group `0` had a `.partition` entry for `requestN-1`. Now we find the first defined and use that index as the next starting group.

* Range splitting: update unit test
This commit is contained in:
Matias Chomicki
2023-03-07 07:44:13 -05:00
committed by GitHub
parent 8999de4313
commit ede3e9e5c4
3 changed files with 41 additions and 11 deletions
@@ -9,7 +9,7 @@ import * as logsTimeSplit from './logsTimeSplit';
import * as metricTimeSplit from './metricTimeSplit';
import { createLokiDatasource, getMockFrames } from './mocks';
import { runPartitionedQueries } from './querySplitting';
import { LokiQuery } from './types';
import { LokiQuery, LokiQueryType } from './types';
describe('runPartitionedQueries()', () => {
let datasource: LokiDatasource;
@@ -145,6 +145,19 @@ describe('runPartitionedQueries()', () => {
expect(datasource.runQuery).toHaveBeenCalledTimes(3);
});
});
test('Groups instant queries', async () => {
const request = getQueryOptions<LokiQuery>({
targets: [
{ expr: 'count_over_time({a="b"}[1m])', refId: 'A', queryType: LokiQueryType.Instant },
{ expr: 'count_over_time({c="d"}[1m])', refId: 'B', queryType: LokiQueryType.Instant },
],
range,
});
await expect(runPartitionedQueries(datasource, request)).toEmitValuesWith(() => {
// Instant queries are omitted from splitting
expect(datasource.runQuery).toHaveBeenCalledTimes(1);
});
});
test('Respects maxLines of logs queries', async () => {
const { logFrameA } = getMockFrames();
const request = getQueryOptions<LokiQuery>({
@@ -162,5 +175,19 @@ describe('runPartitionedQueries()', () => {
expect(datasource.runQuery).toHaveBeenCalledTimes(4);
});
});
test('Groups multiple queries into logs, queries, and instant', async () => {
const request = getQueryOptions<LokiQuery>({
targets: [
{ expr: 'count_over_time({a="b"}[1m])', refId: 'A', queryType: LokiQueryType.Instant },
{ expr: '{c="d"}', refId: 'B' },
{ expr: 'count_over_time({c="d"}[1m])', refId: 'C' },
],
range,
});
await expect(runPartitionedQueries(datasource, request)).toEmitValuesWith(() => {
// 3 days, 3 chunks, 3x Logs + 3x Metric + 1x Instant, 7 requests.
expect(datasource.runQuery).toHaveBeenCalledTimes(7);
});
});
});
});
@@ -8,7 +8,7 @@ import { LokiDatasource } from './datasource';
import { getRangeChunks as getLogsRangeChunks } from './logsTimeSplit';
import { getRangeChunks as getMetricRangeChunks } from './metricTimeSplit';
import { combineResponses, isLogsQuery } from './queryUtils';
import { LokiQuery } from './types';
import { LokiQuery, LokiQueryType } from './types';
/**
* Purposely exposing it to support doing tests without needing to update the repo.
@@ -109,7 +109,7 @@ export function runGroupedQueries(datasource: LokiDatasource, requests: LokiGrou
const nextRequest = () => {
const { nextRequestN, nextRequestGroup } = getNextRequestPointers(requests, requestGroup, requestN);
if (nextRequestN > 0) {
if (nextRequestN > 0 && nextRequestGroup >= 0) {
runNextRequest(subscriber, nextRequestN, nextRequestGroup);
return;
}
@@ -167,16 +167,18 @@ function getNextRequestPointers(requests: LokiGroupedRequest, requestGroup: numb
};
}
return {
nextRequestGroup: 0,
// Find the first group where `[requestN - 1]` is defined
nextRequestGroup: requests.findIndex((group) => group?.partition[requestN - 1] !== undefined),
nextRequestN: requestN - 1,
};
}
export function runPartitionedQueries(datasource: LokiDatasource, request: DataQueryRequest<LokiQuery>) {
const queries = request.targets.filter((query) => !query.hide);
const [logQueries, metricQueries] = partition(queries, (query) => isLogsQuery(query.expr));
const [instantQueries, normalQueries] = partition(queries, (query) => query.queryType === LokiQueryType.Instant);
const [logQueries, metricQueries] = partition(normalQueries, (query) => isLogsQuery(query.expr));
const requests = [];
const requests: LokiGroupedRequest = [];
if (logQueries.length) {
requests.push({
request: { ...request, targets: logQueries },
@@ -189,5 +191,11 @@ export function runPartitionedQueries(datasource: LokiDatasource, request: DataQ
partition: partitionTimeRange(false, request.range, request.intervalMs, metricQueries[0].resolution ?? 1),
});
}
if (instantQueries.length) {
requests.push({
request: { ...request, targets: instantQueries },
partition: [request.range],
});
}
return runGroupedQueries(datasource, requests);
}
@@ -310,11 +310,6 @@ export function requestSupportsPartitioning(allQueries: LokiQuery[]) {
.filter((query) => !query.refId.includes('do-not-chunk'))
.filter((query) => query.expr);
const instantQueries = queries.some((query) => query.queryType === LokiQueryType.Instant);
if (instantQueries) {
return false;
}
return queries.length > 0;
}