Empresas
Empregos
  • Sobre nós
  • Soluções
    • Publicação de vagas
      Publique sua vaga e receba candidatos qualificados em 48h.
    • Avaliações de candidatos
      Mais de 500 testes técnicos e psicológicos, mais anti-fraude.
    • Headhunting
      Busca executiva personalizada do início ao fim.
    • Folha de Pagamento + EOR
      Dispersão de folha e EOR em mais de 15 países da LATAM.
  • Preços
  • Empregos

0

182
Visualizações
Rxjs: recopila todos los valores sin procesar

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), )
  • No puedo usar 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 4

Resultado 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,4
about 4 years ago · Juan Pablo Isaza
2 Respostas
Responde à pergunta

0

Tengo 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.

Operador RxJS personalizado: bufferedExhaustMap

 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)) ) ); }); }

Uso

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) ) )
about 4 years ago · Juan Pablo Isaza Relatório

0

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:

  • Cuando se recibe la primera emisión
  • Después de que se haya completado el trabajo
 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ío
  • false 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 .

about 4 years ago · Juan Pablo Isaza Relatório
Responde à pergunta
Encontrar trabalhos remotos

Descubra a nova forma de encontrar um emprego!

melhores empregos
Principais categorias de trabalho
Empresas
Postar vaga Preços Comercial
Jurídico
Termos e Condições Política de privacidade
© 2026 PeakU Inc. All Rights Reserved.
Andres GPT
Recomende algumas ofertas para mim
Preciso de ajuda