Range splitting: Call subscriber.next only when there are new results to report (#64171)

This commit is contained in:
Matias Chomicki
2023-03-07 13:05:40 +01:00
committed by GitHub
parent dd12cdec4d
commit accef84ca5
@@ -90,7 +90,7 @@ function adjustTargetsFromResponseState(targets: LokiQuery[], response: DataQuer
type LokiGroupedRequest = Array<{ request: DataQueryRequest<LokiQuery>; partition: TimeRange[] }>; type LokiGroupedRequest = Array<{ request: DataQueryRequest<LokiQuery>; partition: TimeRange[] }>;
export function runGroupedQueries(datasource: LokiDatasource, requests: LokiGroupedRequest) { export function runGroupedQueries(datasource: LokiDatasource, requests: LokiGroupedRequest) {
let mergedResponse: DataQueryResponse | null; let mergedResponse: DataQueryResponse = { data: [], state: LoadingState.Streaming };
const totalRequests = Math.max(...requests.map(({ partition }) => partition.length)); const totalRequests = Math.max(...requests.map(({ partition }) => partition.length));
let shouldStop = false; let shouldStop = false;
@@ -101,23 +101,19 @@ export function runGroupedQueries(datasource: LokiDatasource, requests: LokiGrou
return; return;
} }
const done = (response: DataQueryResponse) => { const done = () => {
response.state = LoadingState.Done; mergedResponse.state = LoadingState.Done;
subscriber.next(response); subscriber.next(mergedResponse);
subscriber.complete(); subscriber.complete();
}; };
const nextRequest = () => { const nextRequest = () => {
mergedResponse = mergedResponse || { data: [] };
const { nextRequestN, nextRequestGroup } = getNextRequestPointers(requests, requestGroup, requestN); const { nextRequestN, nextRequestGroup } = getNextRequestPointers(requests, requestGroup, requestN);
if (nextRequestN > 0) { if (nextRequestN > 0) {
mergedResponse.state = LoadingState.Streaming;
subscriber.next(mergedResponse);
runNextRequest(subscriber, nextRequestN, nextRequestGroup); runNextRequest(subscriber, nextRequestN, nextRequestGroup);
return; return;
} }
done(mergedResponse); done();
}; };
const group = requests[requestGroup]; const group = requests[requestGroup];
@@ -125,7 +121,7 @@ export function runGroupedQueries(datasource: LokiDatasource, requests: LokiGrou
const range = group.partition[requestN - 1]; const range = group.partition[requestN - 1];
const targets = adjustTargetsFromResponseState(group.request.targets, mergedResponse); const targets = adjustTargetsFromResponseState(group.request.targets, mergedResponse);
if (!targets.length && mergedResponse) { if (!targets.length) {
nextRequest(); nextRequest();
return; return;
} }
@@ -140,6 +136,7 @@ export function runGroupedQueries(datasource: LokiDatasource, requests: LokiGrou
mergedResponse = combineResponses(mergedResponse, partialResponse); mergedResponse = combineResponses(mergedResponse, partialResponse);
}, },
complete: () => { complete: () => {
subscriber.next(mergedResponse);
nextRequest(); nextRequest();
}, },
error: (error) => { error: (error) => {