So, I made an observable monster.
It's much more complex than your implementation, but in its defence, it's also considerably more powerful.
It's a custom RxJS operator that should work seamlessly with the rest of the RxJS library. If you're running an HTTP call, for example, and get an error you can use retry, retryWhen, catchError, ect to handle errors and respond per-call or for the entire operator.
You can chain other operators (and therefore other observables) before and after this operator to manipulate data before it is used or after it is returned.
If you decide to cancel your function call early, you can just error the source or unsubscribe from the observable (Just like any other operator). It will clean up after itself and cancel any mid-flight operations if possible (Vanilla Promises can't be cancelled).
Anyway, yeah, it's a behemith and perhaps overkill if you're just buffering values for an async function that returns void.
Here it is:
The bufferedExhaustMap operator
This operator buffers values based on a minimum buffer length (time in milliseconds) and minimum buffer size. It always exhausts the current projected observable regardless of whether minimum time and buffer sizes have been met.
This works very much like exhaustMap, except it buffers values instead of discarding values while waiting for the current inner (projected) observable to complete.
/***
* ObservableInput works with an Array, an array-like
* object, a Promise, an iterable object, or an Observable-like object.
*/
function bufferedExhaustMap<T,R>(
project: (v:T[]) => ObservableInput<R>,
minBufferLength = 0,
minBufferCount = 1
): OperatorFunction<T,R> {
return s => new Observable(observer => {
const idle = new BehaviorSubject<boolean>(true);
const setIdle = (b) => () => idle.next(b);
const buffer = (() => {
const bufS = new Subject<(v:T[]) => T[]>();
return {
output: bufS.pipe(
scan((acc, curr) => curr(acc), [])
),
nextVal: (v:T) => bufS.next(curr => [...curr, v]),
clear: () => bufS.next(_ => [])
}
})();
const subProject = combineLatest(
idle,
idle.pipe(
filter(v => !v),
startWith(true),
switchMap(_ => timer(minBufferLength).pipe(
mapTo(true),
startWith(false)
))
),
buffer.output
).pipe(
filter(([idle, bufferedByTime, buffer]) =>
idle && bufferedByTime && buffer.length >= minBufferCount
),
tap(setIdle(false)),
tap(buffer.clear),
map(x => x[2]),
map(project),
mergeMap(projected => from(projected).pipe(
finalize(setIdle(true))
))
).subscribe(observer);
const subSource = s.subscribe({
next: buffer.nextVal,
complete: observer.complete.bind(observer),
error: e => {
subProject.unsubscribe();
observer.error(e);
}
});
return {
unsubscribe: () => {
subProject.unsubscribe();
subSource.unsubscribe();
}
}
});
}
Back to debouncedChunkedQueue
Here is how I would implement your function with the operator I built.
function debouncedChunkedQueue<T>(
fn: (items: T[]) => Promise<void>,
delay = 1000
){
const processSubject = new Subject<T>();
processSubject.pipe(
bufferedExhaustMap(fn, delay)
).subscribe();
return {
push : processSubject.next.bind(processSubject)
}
}
Though this does abstract away RxJS. Instead, I would lean into the reactive style harder and just use the operator as-is. Whatever your promise is doing to the state of your program (Setting values, transforming objects, ect), have it return that information instead and react accordingly.
If that sounds miserable to you, then you're likely better off not using this approach. That's fine too! To each their own. :)
Update # 1: A More Condenced bufferedExhaustMap
Inspired by another answer here (by kvetis), I re-implemented this operator using bufferWhen and extended it to include a minimum buffer size.
This also unsubscribes from the source rather than ignoring the source via exhaustMap while buffering. This should (in theory) be faster, though I'm not sure it would ever matter in practice.
/***
* Buffers, then projects buffered source values to an Observable which is merged in
* the output Observable only if the previous projected Observable has completed.
***/
function bufferedExhaustMap<T,R>(
project: (v:T[]) => ObservableInput<R>,
minBufferLength = 0,
minBufferCount = 1
): OperatorFunction<T,R> {
return s => defer(() => {
const idle = new BehaviorSubject(true);
const setIdle = (b) => () => idle.next(b);
const shared = s.pipe(share());
const nextBufferTime = () => forkJoin(
shared.pipe(take(minBufferCount)),
timer(minBufferLength),
idle.pipe(first(v => v))
).pipe(
tap(setIdle(false))
);
return shared.pipe(
bufferWhen(nextBufferTime),
map(project),
mergeMap(projected => from(projected).pipe(
finalize(setIdle(true))
))
);
});
}
Perhaps the most complex bit is the nextBufferTime factory function. It creates an observable that emits when the buffer should be cleared, then completes.
It does this by merging 3 observables and only completing once all three observables are complete. In this way, an inner observable completing counts as a condition being met.
- First condition: The source has emitted
minBufferCount values.
- Second condition: A timer of
minBufferLength has elapsed.
- Third condition: No projected observable is currently active (
idle == true).
const nextBufferTime = () => forkJoin(
// Take minBufferCount from source then complete
shared.pipe(take(minBufferCount)),
// Wait minBufferLength milliseconds and then complete
timer(minBufferLength),
// Wait until idle emits true, then complete
idle.pipe(first(v => v))
)