RxJS: perché il fuoco osservabile dall'interno è prima?

Sep 21 2020

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

3 olivarra1 Sep 21 2020 at 21:22

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?)

1 GogaKoreli Sep 21 2020 at 21:41

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