Quiero usar Subject tiene una cola de tareas asíncrona con simultaneidad = 1.
La recopilación de nuevas tareas es más rápida que la realización de esas tareas, por lo que, como optimización del rendimiento, me gustaría procesar todas las tareas no procesadas hasta ahora.
El siguiente código repasa tarea por tarea:
this.tasks .pipe( mergeMap(async task => { await this.processTasks([task]) }, 1), )Me gustaría convertir el código anterior en algo similar a:
this.tasks .pipe( mergeMap___AllUnprocessedUntilNow(async tasks => { await this.processTasks(tasks) }, 1), )bufferTime o bufferCount porque introducirán una latencia adicional en cada nueva tarea.Resultado real de la ejecución:
* adding task 1 * procssing task 1 (takes alot of time to perform it) * meanwhile, adding more tasks: 2,3,4 * processed task 1 * procssing task 2 * processed task 2 * procssing task 3 * processed task 3 * procssing task 4 * processed task 4Resultado esperado de la ejecución:
* adding task 1 * procssing task 1 (takes alot of time to perform it) * meanwhile, adding more tasks: 2,3,4 * processed task 1 * procssing task 2,3,4 * processed task 2,3,4Tengo un operador personalizado que hace esto. Atención, hay algunas partes que hacen que esto sea más complicado de lo que piensas al principio. Además de eso, esto está un poco sobrediseñado, pero parece que puede usarlo con los valores predeterminados para obtener lo que desea.
function bufferedExhaustMap<T, R>( project: (v: T[]) => ObservableInput<R>, minBufferLength = 0, minBufferCount = 1, concurrent = 1 ): OperatorFunction<T, R> { function diffCounter(){ const incOrDec = new Subject<boolean>(); return { onAvailable: incOrDec.pipe( scan((acc, curr) => (curr ? ++acc : --acc), 0), startWith(0), shareReplay(1), first(count => count < concurrent), mapTo(true) ), start: () => incOrDec.next(true), end: () => incOrDec.next(false) }; } return source => defer(() => { const projectCount = diffCounter(); const shared = source.pipe(share()); const nextBufferTime = () => forkJoin([ shared.pipe(take(minBufferCount), delay(0)), timer(minBufferLength), projectCount.onAvailable ]); return shared.pipe( bufferWhen(nextBufferTime), delayWhen(() => projectCount.onAvailable), tap(projectCount.start), map(project), mergeMap(projected => from(projected).pipe( finalize(projectCount.end)) ) ); }); }Tenga en cuenta que eliminé algunas promesas superfluas/mezclas de RxJS, ya que es un olor a código.
this.tasks.pipe( mergeMap___AllUnprocessedUntilNow(async tasks => { await this.processTasks(tasks) }, 1), )se convierte
this.tasks.pipe( mergeMap___AllUnprocessedUntilNow( tasks => this.processTasks(tasks) , 1 ) ) (recuerde que todos los operadores RxJS de alto orden convertirán los iterables y las promesas en observables para usted. De lo contrario, puede usar el operador from en su lugar).
Dado que bufferedExhaustMap tiene una concurrencia predeterminada de 1, puede escribir esto con el nuevo operador de la siguiente manera:
this.tasks.pipe( bufferedExhaustMap( tasks => this.processTasks(tasks) ) )Puede lograr esto con una combinación de buffer y concatMap , junto con un poco de ayuda de filter y tap .
La idea es activar la liberación del búfer en el momento correcto. En tu caso hay dos ocasiones para soltar:
const release$ = new Subject<void>(); let releaseOnNextEmit = true; const work$ = item$.pipe( tap(() => { if(releaseOnNextEmit) setTimeout(() => release$.next(), 0); }), buffer(release$), tap(tasks => releaseOnNextEmit = tasks.length === 0), filter(tasks => tasks.length > 0), concatMap(tasks => processTasks(tasks)), tap(() => release$.next()) ); Puede ver que usamos un Asunto para desencadenar la liberación del búfer y un indicador releaseOnNextEmit para indicar si recibir una emisión debe desencadenar la liberación del búfer.
Después de que se emita el búfer, establecemos la bandera en consecuencia:
true cuando el búfer está vacíofalse cuando NO está vacío El filter se usa para evitar pasar matrices vacías a concatMap ( el búfer emitirá una matriz vacía si el búfer está vacío cuando se emite su disparador ).
Usamos concatMap para ejecutar el trabajo real, luego simplemente liberamos el búfer después.
Aquí hay una demostración de StackBlitz en funcionamiento.
Con todos los tap allí, esto no es muy fácil de seguir, por lo que puede ser beneficioso ponerlo en un operador personalizado:
function bufferConcatMap<T, R>(project: (a: T[]) => ObservableInput<R>) { const release$ = new Subject<void>(); let releaseOnNextEmit = true; return (source$: Observable<T>) => source$.pipe( tap(() => { if(releaseOnNextEmit) setTimeout(() => release$.next(), 0); }), buffer(release$), tap(tasks => releaseOnNextEmit = tasks.length === 0), filter(tasks => tasks.length > 0), concatMap(project), tap(() => release$.next()) ); } const work$ = item$.pipe( bufferConcatMap(tasks => processTasks(tasks)) );Aquí hay otro StackBlitz .