diff --git a/public/app/plugins/datasource/loki/querySplitting.test.ts b/public/app/plugins/datasource/loki/querySplitting.test.ts index 57f9859ea6c..5920b08fcdd 100644 --- a/public/app/plugins/datasource/loki/querySplitting.test.ts +++ b/public/app/plugins/datasource/loki/querySplitting.test.ts @@ -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) diff --git a/public/app/plugins/datasource/loki/querySplitting.ts b/public/app/plugins/datasource/loki/querySplitting.ts index 262b8708289..d5ff1c870df 100644 --- a/public/app/plugins/datasource/loki/querySplitting.ts +++ b/public/app/plugins/datasource/loki/querySplitting.ts @@ -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) { +export function runSplitQuery( + datasource: LokiDatasource, + request: DataQueryRequest, + 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, { diff --git a/public/app/plugins/datasource/loki/shardQuerySplitting.test.ts b/public/app/plugins/datasource/loki/shardQuerySplitting.test.ts index ea5e09877e0..51b8bb308ee 100644 --- a/public/app/plugins/datasource/loki/shardQuerySplitting.test.ts +++ b/public/app/plugins/datasource/loki/shardQuerySplitting.test.ts @@ -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__=""}', diff --git a/public/app/plugins/datasource/loki/shardQuerySplitting.ts b/public/app/plugins/datasource/loki/shardQuerySplitting.ts index e159263640c..9f8dc397b8d 100644 --- a/public/app/plugins/datasource/loki/shardQuerySplitting.ts +++ b/public/app/plugins/datasource/loki/shardQuerySplitting.ts @@ -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)) {