I have a paginated third-party resource living in a web service. What I want to do is turn that paginated resource into a stream of values, and let the client decide how many elements to use. That is, the client should not know that the original resource is paginated.
So far I got the following code:
import { from, Observable, of } from 'rxjs';
import { expand, mergeMap } from 'rxjs/operators';
interface Result {
page: number;
items: number[];
}
// Assume that this does a real HTTP request
function fetchPage(page: number = 0): Observable<Result> {
let items = [...Array(3).keys()].map((e) => e + 3 * page);
console.log('Fetching page ' + page);
return of({ page, items });
}
// Turn a paginated request into an infinite stream of 'number'
function mkInfiniteStream(): Observable<number> {
return fetchPage().pipe(
expand((res) => fetchPage(res.page + 1)),
mergeMap((res) => from(res.items))
);
}
const infinite$ = mkInfiniteStream();
This works really well: I get a lazy infinite stream of numbers, and the client can just do infinite$.pipe(take(n)) and take the first n elements, without knowing that the underlying resource is paginated.
Now, what I want to do is to share those values when dealing with multiple subscribers, that is:
infinite$.pipe(take(10)).subscribe((v) => console.log('[1] got ', v));
// Assume that later we have new subscribers
setTimeout(() => {
infinte$.pipe(take(5)).subscribe((v) => console.log('[2] got ', v));
}, 1000);
setTimeout(() => {
infinte$.pipe(take(4)).subscribe((v) => console.log('[3] got ', v));
}, 1500);
As you can see, we'll have multiple subscribers to the infinite stream, and I want to reuse the values already emitted in order to reduce the number of fetchPage calls. In this example, once a client asked for 10 items (take(10)), then any client that requests less than 10 items (ex. 5 items) should result in no calls to fetchPage since those items were already emited.
I tried the following, but I could not get the desired behavior:
const infinite$ = mkInfiniteStream().pipe(share())
Does not work since late subscribers result in multiple calls to 'fetchPage'.
const infinite$ = mkInfiniteStream().pipe(shareReplay())
Forces all values in the stream even though they are not needed (no client asked for all items yet)
Any hint would be appreciated. In case anyone wants to try the code: https://stackblitz.com/edit/n4ywfw