Shard query splitting: run queries through time query splitting (#98126)

* Query splitting: add skipPartialUpdates option

* Shard query splitting: run queries through time splitting

* Query splitting: delegate error retry to shard splitting

* Shard query splitting: update unit tests

* Shard query splitting: test combined requests

* Formatting

* Query splitting: test new options

* Query splitting: update assertion

* Formatting
This commit is contained in:
Matias Chomicki
2025-01-15 10:26:23 +01:00
committed by GitHub
parent a6eb8abd05
commit bbade6b011
4 changed files with 99 additions and 32 deletions
@@ -64,9 +64,20 @@ describe('runSplitQuery()', () => {
});
test('Splits datasource queries', async () => {
await expect(runSplitQuery(datasource, request)).toEmitValuesWith(() => {
await expect(runSplitQuery(datasource, request)).toEmitValuesWith((emitted) => {
// 3 days, 3 chunks, 3 requests.
expect(datasource.runQuery).toHaveBeenCalledTimes(3);
// 3 sub-requests + complete
expect(emitted).toHaveLength(4);
});
});
test('Skips partial updates as an option', async () => {
await expect(runSplitQuery(datasource, request, { skipPartialUpdates: true })).toEmitValuesWith((emitted) => {
// 3 days, 3 chunks, 3 requests.
expect(datasource.runQuery).toHaveBeenCalledTimes(3);
// partial updates skipped
expect(emitted).toHaveLength(1);
});
});
@@ -80,6 +91,16 @@ describe('runSplitQuery()', () => {
});
});
test('Does not retry failed queries as an option', async () => {
jest
.mocked(datasource.runQuery)
.mockReturnValueOnce(of({ state: LoadingState.Error, errors: [{ refId: 'A', message: 'timeout' }], data: [] }));
await expect(runSplitQuery(datasource, request, { disableRetry: true })).toEmitValuesWith(() => {
// No retries
expect(datasource.runQuery).toHaveBeenCalledTimes(1);
});
});
test('Does not retry on other errors', async () => {
jest
.mocked(datasource.runQuery)
@@ -48,6 +48,17 @@ export function partitionTimeRange(
});
}
interface QuerySplittingOptions {
/**
* Tells the query splitting code to not emit partial updates. Only emit on error or when it finishes querying.
*/
skipPartialUpdates?: boolean;
/**
* Do not retry failed queries.
*/
disableRetry?: boolean;
}
/**
* Based in the state of the current response, if any, adjust target parameters such as `maxLines`.
* For `maxLines`, we will update it as `maxLines - current amount of lines`.
@@ -76,7 +87,11 @@ export function adjustTargetsFromResponseState(targets: LokiQuery[], response: D
})
.filter((target) => target.maxLines === undefined || target.maxLines > 0);
}
export function runSplitGroupedQueries(datasource: LokiDatasource, requests: LokiGroupedRequest[]) {
export function runSplitGroupedQueries(
datasource: LokiDatasource,
requests: LokiGroupedRequest[],
options: QuerySplittingOptions = {}
) {
const responseKey = requests.length ? requests[0].request.queryGroupId : uuidv4();
let mergedResponse: DataQueryResponse = { data: [], state: LoadingState.Streaming, key: responseKey };
const totalRequests = Math.max(...requests.map(({ partition }) => partition.length));
@@ -116,6 +131,9 @@ export function runSplitGroupedQueries(datasource: LokiDatasource, requests: Lok
};
const retry = (errorResponse?: DataQueryResponse) => {
if (options.disableRetry) {
return false;
}
try {
if (errorResponse && !isRetriableError(errorResponse)) {
return false;
@@ -170,13 +188,17 @@ export function runSplitGroupedQueries(datasource: LokiDatasource, requests: Lok
shouldStop = true;
}
mergedResponse = combineResponses(mergedResponse, partialResponse);
mergedResponse = updateLoadingFrame(mergedResponse, subRequest, longestPartition, requestN);
if (!options.skipPartialUpdates) {
mergedResponse = updateLoadingFrame(mergedResponse, subRequest, longestPartition, requestN);
}
},
complete: () => {
if (retrying) {
return;
}
subscriber.next(mergedResponse);
if (!options.skipPartialUpdates) {
subscriber.next(mergedResponse);
}
nextRequest();
},
error: (error) => {
@@ -268,7 +290,11 @@ function querySupportsSplitting(query: LokiQuery) {
);
}
export function runSplitQuery(datasource: LokiDatasource, request: DataQueryRequest<LokiQuery>) {
export function runSplitQuery(
datasource: LokiDatasource,
request: DataQueryRequest<LokiQuery>,
options: QuerySplittingOptions = {}
) {
const queries = request.targets.filter((query) => !query.hide).filter((query) => query.expr);
const [nonSplittingQueries, normalQueries] = partition(queries, (query) => !querySupportsSplitting(query));
const [logQueries, metricQueries] = partition(normalQueries, (query) => isLogsQuery(query.expr));
@@ -330,7 +356,7 @@ export function runSplitQuery(datasource: LokiDatasource, request: DataQueryRequ
}
const startTime = new Date();
return runSplitGroupedQueries(datasource, requests).pipe(
return runSplitGroupedQueries(datasource, requests, options).pipe(
tap((response) => {
if (response.state === LoadingState.Done) {
trackGroupedQueries(response, requests, request, startTime, {
@@ -66,6 +66,24 @@ describe('runShardSplitQuery()', () => {
});
test('Splits datasource queries', async () => {
const querySplittingRange = {
from: dateTime('2023-02-08T05:00:00.000Z'),
to: dateTime('2023-02-10T06:00:00.000Z'),
raw: {
from: dateTime('2023-02-08T05:00:00.000Z'),
to: dateTime('2023-02-10T06:00:00.000Z'),
},
};
request = createRequest([{ expr: '$SELECTOR', refId: 'A', direction: LokiQueryDirection.Scan }], {
range: querySplittingRange,
});
await expect(runShardSplitQuery(datasource, request)).toEmitValuesWith(() => {
// 5 shards, 3 groups + empty shard group, 4 requests * 3 days, 3 chunks, 3 requests = 12 requests
expect(datasource.runQuery).toHaveBeenCalledTimes(12);
});
});
test('Users query splitting for querying over a day', async () => {
await expect(runShardSplitQuery(datasource, request)).toEmitValuesWith(() => {
// 5 shards, 3 groups + empty shard group, 4 requests
expect(datasource.runQuery).toHaveBeenCalledTimes(4);
@@ -79,7 +97,8 @@ describe('runShardSplitQuery()', () => {
expect(datasource.runQuery).toHaveBeenCalledWith({
intervalMs: expect.any(Number),
range: expect.any(Object),
requestId: 'TEST_shard_0_0_2',
queryGroupId: expect.any(String),
requestId: 'TEST_shard_0_0_2_1',
targets: [
{
expr: '{a="b", __stream_shard__=~"20|10"}',
@@ -92,7 +111,8 @@ describe('runShardSplitQuery()', () => {
expect(datasource.runQuery).toHaveBeenCalledWith({
intervalMs: expect.any(Number),
range: expect.any(Object),
requestId: 'TEST_shard_0_2_2',
queryGroupId: expect.any(String),
requestId: 'TEST_shard_0_2_2_1',
targets: [
{
expr: '{a="b", __stream_shard__=~"3|2"}',
@@ -105,7 +125,8 @@ describe('runShardSplitQuery()', () => {
expect(datasource.runQuery).toHaveBeenCalledWith({
intervalMs: expect.any(Number),
range: expect.any(Object),
requestId: 'TEST_shard_0_4_1',
queryGroupId: expect.any(String),
requestId: 'TEST_shard_0_4_1_1',
targets: [
{
expr: '{a="b", __stream_shard__="1"}',
@@ -118,7 +139,8 @@ describe('runShardSplitQuery()', () => {
expect(datasource.runQuery).toHaveBeenCalledWith({
intervalMs: expect.any(Number),
range: expect.any(Object),
requestId: 'TEST_shard_0_5_1',
queryGroupId: expect.any(String),
requestId: 'TEST_shard_0_5_1_1',
targets: [
{
expr: '{a="b", __stream_shard__=""}',
@@ -147,7 +169,8 @@ describe('runShardSplitQuery()', () => {
expect(datasource.runQuery).toHaveBeenCalledWith({
intervalMs: expect.any(Number),
range: expect.any(Object),
requestId: 'TEST_shard_0_0_2',
queryGroupId: expect.any(String),
requestId: 'TEST_shard_0_0_2_1',
targets: [
{
expr: '{service_name="test", filter="true", __stream_shard__=~"20|10"}',
@@ -208,17 +231,6 @@ describe('runShardSplitQuery()', () => {
});
test('Adjusts the group size based on errors and execution time', async () => {
const request = createRequest([{ expr: '$SELECTOR', refId: 'A', direction: LokiQueryDirection.Scan }], {
range: {
from: dateTime('2024-11-13T05:00:00.000Z'),
to: dateTime('2024-11-14T06:00:00.000Z'),
raw: {
from: dateTime('2024-11-13T05:00:00.000Z'),
to: dateTime('2024-11-14T06:00:00.000Z'),
},
},
});
jest
.mocked(datasource.languageProvider.fetchLabelValues)
.mockResolvedValue(['1', '10', '2', '20', '3', '4', '5', '6', '7', '8', '9']);
@@ -373,7 +385,8 @@ describe('runShardSplitQuery()', () => {
expect(datasource.runQuery).toHaveBeenCalledWith({
intervalMs: expect.any(Number),
range: expect.any(Object),
requestId: 'TEST_shard_0_0_3',
queryGroupId: expect.any(String),
requestId: 'TEST_shard_0_0_3_1',
targets: [
{
expr: '{a="b", __stream_shard__=~"20|10|9"}',
@@ -387,7 +400,8 @@ describe('runShardSplitQuery()', () => {
expect(datasource.runQuery).toHaveBeenCalledWith({
intervalMs: expect.any(Number),
range: expect.any(Object),
requestId: 'TEST_shard_0_3_4',
queryGroupId: expect.any(String),
requestId: 'TEST_shard_0_3_4_1',
targets: [
{
expr: '{a="b", __stream_shard__=~"8|7|6|5"}',
@@ -401,7 +415,8 @@ describe('runShardSplitQuery()', () => {
expect(datasource.runQuery).toHaveBeenCalledWith({
intervalMs: expect.any(Number),
range: expect.any(Object),
requestId: 'TEST_shard_0_3_2',
queryGroupId: expect.any(String),
requestId: 'TEST_shard_0_3_2_1',
targets: [
{
expr: '{a="b", __stream_shard__=~"8|7"}',
@@ -415,7 +430,8 @@ describe('runShardSplitQuery()', () => {
expect(datasource.runQuery).toHaveBeenCalledWith({
intervalMs: expect.any(Number),
range: expect.any(Object),
requestId: 'TEST_shard_0_5_3',
queryGroupId: expect.any(String),
requestId: 'TEST_shard_0_5_3_1',
targets: [
{
expr: '{a="b", __stream_shard__=~"6|5|4"}',
@@ -429,7 +445,8 @@ describe('runShardSplitQuery()', () => {
expect(datasource.runQuery).toHaveBeenCalledWith({
intervalMs: expect.any(Number),
range: expect.any(Object),
requestId: 'TEST_shard_0_8_2',
queryGroupId: expect.any(String),
requestId: 'TEST_shard_0_8_2_1',
targets: [
{
expr: '{a="b", __stream_shard__=~"3|2"}',
@@ -443,7 +460,8 @@ describe('runShardSplitQuery()', () => {
expect(datasource.runQuery).toHaveBeenCalledWith({
intervalMs: expect.any(Number),
range: expect.any(Object),
requestId: 'TEST_shard_0_10_1',
queryGroupId: expect.any(String),
requestId: 'TEST_shard_0_10_1_1',
targets: [
{
expr: '{a="b", __stream_shard__="1"}',
@@ -457,7 +475,8 @@ describe('runShardSplitQuery()', () => {
expect(datasource.runQuery).toHaveBeenCalledWith({
intervalMs: expect.any(Number),
range: expect.any(Object),
requestId: 'TEST_shard_0_11_1',
queryGroupId: expect.any(String),
requestId: 'TEST_shard_0_11_1_1',
targets: [
{
expr: '{a="b", __stream_shard__=""}',
@@ -160,9 +160,10 @@ function splitQueriesByStreamShard(
debug(shardsToQuery.length ? `Querying ${shardsToQuery.join(', ')}` : 'Running regular query');
const queryRunner =
shardsToQuery.length > 0 ? datasource.runQuery.bind(datasource) : runSplitQuery.bind(null, datasource);
subquerySubscription = queryRunner(subRequest).subscribe({
subquerySubscription = runSplitQuery(datasource, subRequest, {
skipPartialUpdates: true,
disableRetry: true,
}).subscribe({
next: (partialResponse: DataQueryResponse) => {
if ((partialResponse.errors ?? []).length > 0 || partialResponse.error != null) {
if (retry(partialResponse)) {