RXJS運算子-mergeMap、switchMap

May 28, 2020

Observable 與 pip

在 Rxjs中,Observable 是一個可供訂閱的物件。當一個 Observablesubscribe (訂閱)後,直到這個 Observable 發出 complete 之前,訂閱這個 Observable的 Observer 都會收到 Observable 發出的值。

因為 Observable 會在特定的條件下,不斷的發出值,因此就會形成一個資料流。可以把資料流想像成水流,而 發出值的 Observable 就是水的源頭。

而當資料流產生後,我們可能還會對資料的值一些處理,最將處理後的值發給訂閱這個 ObservableObserver。這時就可以使用 pipe。可以把 pipel 想像成一個水管。

打個比方,源 Observable 就像 水龍頭,而 pipe 像蓮蓬頭。當水龍頭打開後,我們(Observer)最後看到的水,會是蓮蓬頭導出的水。

pipe 中的 operator (運算子),input 和 output都會是 Observable

所有的運算子,其實都是一個Observable。以 map為例,當project return 值後,map 會發出project return 的值;pipe中的函式會像接水管一樣,下一個Operator訂閱上一個Operator(Observable)。

訂閱的動作,可以透過 Subscriber 完成。它是一個內部使用的類別,實現了 Observer界面,並且繼承 Subscription。內部的 Observer 界面,會被轉化為 Subscriber。 Subscriber 實現了「管理」及「取消」多個Observable 訂閱的能力 。因此,在 pipe中,當 其中的某個運算子要結束掉整個 pipe的時候,可以通過 Subscriber 的 unsubscribe() 方法,取消整個管線的訂閱。

引用官方說明:

A Pipeable Operator is a function that takes an Observable as its input and returns another Observable. It is a pure operation: the previous Observable stays unmodified.

map operator

函式: map(project: Function, thisArg: any): Observable

map operator是很好懂的一個運算子。它的 project function 會接受上一個運算子發出的值,並且經過處理後,回傳處理後的值。

範例:

// RxJS v6+
import { from } from 'rxjs';
import { map } from 'rxjs/operators';

// 发出 (1,2,3,4,5)
const source = from([1, 2, 3, 4, 5]);
// 每个数字加10
const example = source.pipe(map(val => val + 10));
// 输出: 11,12,13,14,15
const subscribe = example.subscribe(val => console.log(val));

mergeMap and switchMap

mergeMap: mergeMap(project: function: Observable, resultSelector: function: any, concurrent: number): Observable

switchMap: switchMap(project: function: Observable, resultSelector: function(outerValue, innerValue, outerIndex, innerIndex): any): Observable

很多時候,我們會需要回傳一個 Observable,並且自已處理 Observable 的發出值。最常見的範例就是在接到值後,面發起一個 http請求。此時,return的值是一個 Observable

mergeMapswitchMap 事實上也沒有什麼魔法。它就是 project function 直接回傳一個 Observable,而下一個運算子取得的值,是由這個回傳的Observable發出的。

mergeMapswitchMapmap 的差別,在於 project function。 map回傳的是Observable內的值;mergeMap/ switchMap回傳的是一個函式。

怎麼知道發出的值是從哪一個 Observable發出的?Observable何時 complete

這就是 mregeMapswitchMap 的不同處了。因為它們的project function 回傳值是 Observable,我們不會知道它會什麼時候發出值、什麼時候結束。

mergeMapswitchMap 都有內部的 Observable,下一個運算子接收的值,是由這個 Observable 發出的值。而兩者的差別只在,使用switchMap的時候,當新的observable進來了,上一個 Observable會直接 complete。而 mergeMap不會complete上一個 Observable。

當Observable被發出後,該Observable之後再發出的任何訊息,都不會再被訂閱者接收(在pipe中就是下一個運算子)。

switchMap一次只會有一個訂閱中的observable。而mergeMap會一次管理多個Observable訂閱。

因此,使用switchMap可以防止潛在可能memory leak的問題。大多數的時候,使用switchMap是比較安全的。然而,有些情況下,可能還是適合使用 mergeMap。這就要看使用情境跟經驗了。

例子

這是從rxjs官網看到,覺得比較有趣的例子。

StackBlitz

// RxJS v6+
import { interval, fromEvent, merge, empty } from 'rxjs';
import { switchMap, scan, takeWhile, startWith, mapTo } from 'rxjs/operators';

const countdownSeconds = 10;
const setHTML = id => val => (document.getElementById(id).innerHTML = val);
const pauseButton = document.getElementById('pause');
const resumeButton = document.getElementById('resume');
const interval$ = interval(1000).pipe(mapTo(-1));

const pause$ = fromEvent(pauseButton, 'click').pipe(mapTo(false));
const resume$ = fromEvent(resumeButton, 'click').pipe(mapTo(true));

const timer$ = merge(pause$, resume$)
  .pipe(
    startWith(true),
    switchMap(val => (val ? interval$ : empty())),
    tap(val => console.log(val)),
    scan((acc, curr) => (curr ? curr + acc : acc), countdownSeconds),
    takeWhile(v => v >= 0)
  )
  .subscribe(setHTML('remaining'));

首先,原 Observable是 pause和 resume兩個 event 的 merge。 因此,timer$ 要馬會接收 pause button 的 click事件,回傳 true 要馬會接收 resume button 的 click事件,回傳 false

接著,switchMap 會接收 merge 發出的 true 或 false 值,並判斷要回傳 interval$ observable 或是 empty。

interval$這個 observable會每隔一秒發出 -1 。

這裡因為使用 switchMap,所以當收到新的 click事件,前一個 回傳的 observable會被 結束掉,因此不會有潛在的 memory leak風險。

接下來,scan會將收到的值累加,並發出值給下一個運算子。

最後,takeWhile 會在 v > 0 條件成立時,發出接收到的值。

這樣整個 pipe做到的效果,就是一個 監聽 按鈕有沒有被按,並且做對應動作的倒計時器了。

在switchMap 後面 加上一個 tap/do 運算子,log一下接收的值,對switchMap會更了解。當你按下 pause時,switchMap會對內部的Obesrvable送出complete,並且發出一個新的Observable,在這裡是 empty。因此, 前一個interval$ 的 Obsevable因為已經complete,不會再發出值了,達成停止counting 的效果。

參考

https://medium.com/allen%E7%9A%84%E6%8A%80%E8%A1%93%E7%AD%86%E8%A8%98/rxjs-mergeMap-%E7%AD%86%E8%A8%98-b971778cbff4

https://rxjs-cn.github.io/learn-rxjs-operators/operators/transformation/switchMap.html

https://rxjs-cn.github.io/learn-rxjs-operators/operators/transformation/mergeMap.html



Written by Howard Chang , software engineer, programming lover, from Taiwan