RxJS: perché il fuoco osservabile dall'interno è prima?
Sto usando RxJS 6 e ho il seguente problema di esempio :
Vogliamo bufferizzare gli elementi per uno specifico bufferTimema se non succede nulla per un periodo di tempo maggiore di quello bufferTimeche vogliamo che il primo elemento si attivi immediatamente.
Sequenza:
[------bufferTime------]
Input over time:
[1, 2, 3, -------------|---4, 5, 6 ----------------]
Output over time:
[1]-----------------[2,3]---[4]------------------[5,6]
Questo è il codice che mi porta lì:
source$.pipe( buffer(source$.pipe(
throttleTime(bufferTime, asyncScheduler, {leading: true, trailing: true}),
delay(10) // <-- This here bugs me like crazy though!
)
)
La mia domanda riguarda l' delayoperatore. Quando lo ometto, il buffer si attiva con un elenco vuoto perché $source.pipe(throttleTime(...))è più veloce del passaggio del buffer.
Senza delay
[------bufferTime------]
Input over time:
[1, 2, 3, -------------|---4, 5, 6 ----------------]
Output over time:
[]------------------[1,2,3]--[]------------------[4,5,6]
C'è un modo per sbarazzarsi di delay?
Risposte
Il problema è che l'ordine in cui gli osservatori si iscrivono a uno stream viene preservato quando arriva un nuovo valore: il primo che si è iscritto riceverà per primo il valore, poi il successivo e così via.
Nella v6 , buffersi iscriverà prima al flusso interno, che si iscriverà per source$primo, e poi (buffer) si iscriverà al source$secondo posto. Per chiarire, quando arriva un nuovo valore, il primo a ricevere il nuovo valore sarà l'osservabile interno, perché sottoscritto prima alla fonte buffer.
Nella v7 è stato sostituito, quindi dovrebbe funzionare senza alcun lavoro aggiuntivo
Come soluzione più pulita per v6, c'è un operatore che può aiutare a: subscribeOn. Con questo puoi specificare uno scheduler da utilizzare quando ti iscrivi all'operatore, nel tuo caso utilizzando un micro-task con asapScheduleresso dovrebbe funzionare:
source$.pipe( buffer(source$.pipe(
subscribeOn(asapScheduler),
throttleTime(bufferTime, asyncScheduler, {leading: true, trailing: true}),
)
)
In questo modo la sottoscrizione interna al sorgente verrà posticipata dopo che il buffer si è iscritto ad essa.
Tieni presente che se stai creando un operatore personalizzato con questa implementazione, potresti voler inviare in multicast la sorgente $ con l'overload publish(multicasted$ => ...), altrimenti questo frammento farà 2 iscrizioni a source$(potenzialmente inviando una richiesta due volte?)
Prova il codice seguente, puoi giocarci nello StackBlitz :
const source$: Observable<string> = interval(300).pipe( map(x => "a" + x), share() ); source$.pipe(
exhaustMap(value => {
return concat(
of([value]),
source$.pipe(
bufferTime(1000),
take(1)
)
);
})
).subscribe(x => console.log(x, new Date()));
L'output è:
["a0"] 2020-09-21T14:39:49.895Z
["a1", "a2", "a3"] 2020-09-21T14:39:50.899Z
["a4"] 2020-09-21T14:39:51.096Z
["a5", "a6", "a7"] 2020-09-21T14:40:30.967Z