Este es un ejercicio puramente académico para entender a los operadores. ¿Cómo puedo usar RXJS para intercalar transmisiones entrantes?
por ejemplo, ir de esto:
// RxJS v6+ import { of } from 'rxjs'; //emits any number of provided values in sequence const source = of(1, 2, 3, 4, 5); //output: 1,2,3,4,5 const subscribe = source.subscribe(val => console.log(val));a esto:
import { of } from 'rxjs'; const source1 = of(1,3,5,7,9); const source2 = of(2,4,6,8,10); // a simple merge will just append source1 and source2 // how do I obtain an output = 1,2,3,4,5,6,7,8,9,10 const subscribe = merge(source1, source2).subscribe(val => console.log(val))Restricciones: esta es una pregunta general, por lo que no quiero usar predicados basados en la separación impar/par de este ejemplo específico, sino crear una salida basada únicamente en el orden de los elementos entrantes, es decir, tomar el primero de cada flujo, luego el segundo, etc.
¿Es eso posible?
La forma estándar de hacer esto es, de hecho, la función de combinación , que puede recibir cualquier número de observables como argumentos y emite todos sus elementos en secuencia, independientemente del orden de los observables en los parámetros de la función. Pero como señaló, adjuntó los resultados, la razón de esto no es un problema con la combinación, sino que en realidad emitió los valores de forma sincrónica, por lo que todos se ejecutaron en orden en el mismo bucle de JavaScript. Puede cambiar eso mediante el uso de un observable que emite valores de forma asíncrona, o crear un observable con algunos de los programadores asíncronos disponibles.
Creé esta función asyncOf como una maqueta de un Observable con valores emitidos a lo largo del tiempo, como un cliente websocket o interacciones de usuarios.
import { merge, Observable, EMPTY } from 'rxjs'; import { toArray } from 'rxjs/operators'; const asyncOf = <T = any>(args: T[], interval = 500): Observable<T> => { let i = 0; let intervalRef; if(!args?.length) { return EMPTY; } return new Observable(observer => { intervalRef = setInterval(() => { if(args[i]) { observer.next(args[i]) } ++i; if(!args[i]) { observer.complete() if(intervalRef) { clearInterval(intervalRef) } } }, interval) }) } const source1 = asyncOf([1, 3, 5, 7, 9]); const source2 = asyncOf([2, 4, 6, 8, 10]); //output: 1,2,3,4,5,6,7,8,9,10 const subscribe = merge(source1, source2) //.pipe(toArray()) //collect all emmited values and emmit an array when the observable completes .subscribe((val) => console.log(val));En este caso, la combinación funciona correctamente, porque ambos observables se ejecutan de forma asíncrona.
Otra opción es usar uno de los programadores asíncronos (o crear uno propio) para describir cuándo se debe emitir cada elemento. Esto está más cerca de su ejemplo y es una inmersión un poco más profunda en rxjs:
import { merge, scheduled, asyncScheduler } from 'rxjs'; import { toArray } from 'rxjs/operators'; const source1 = scheduled([1, 3, 5, 7, 9], asyncScheduler); const source2 = scheduled([2, 4, 6, 8, 10], asyncScheduler); //output: 1,2,3,4,5,6,7,8,9,10 const subscribe = merge(source1, source2) //.pipe(toArray()) //collect all emmited values and emmit an array when the observable completes .subscribe((val) => console.log(val));Programado viene en RXJS 6.5+ y desaprueba el uso de la función de programador en otras funciones como de ,
Puede zip las secuencias y luego descomprimir los valores comprimidos de esta manera:
const odd$ = of(1, 3, 5, 7, 9).pipe(delay(50)); const even$ = of(2, 4, 6, 8, 10); zip([odd$, even$]) .pipe(switchMap(([odd, even]) => of(odd, even))) .subscribe(console.log);stackblitz: https://stackblitz.com/edit/rxjs-uktrkz?devtoolsheight=60&file=index.ts