Học RxJS
Higher-order Observables

switchMap

Chuyển sang stream mới nhất và hủy stream cũ.

Bạn gõ rx, rồi đổi thành rxjs, nhưng response của rx lại về sau và ghi đè danh sách mới. Với màn hình tìm kiếm, kết quả cũ không còn hữu ích chỉ vì request của nó chưa xong. switchMap giải quyết đúng yêu cầu này: khi nhận input mới, nó bỏ subscription của inner cũ và chuyển sang inner mới.

Phạm vi và kiến thức nền

Bài này dùng RxJS 7.8.2, TypeScript và import operator từ rxjs. Bạn nên biết pipe, subscribe và outer stream, inner stream. Ví dụ timer chạy được trong playground có RxJS; ví dụ DOM cần trình duyệt, còn endpoint HTTP là API minh họa của ứng dụng.

Mục lục

Khi nào chọn switchMap

Mình chọn switchMap khi input mới làm công việc cũ hết giá trị đối với consumer này. Tìm kiếm theo từ khóa, tải chi tiết theo route param hoặc chuyển subscription realtime sang resource vừa chọn đều thường có yêu cầu đó. Một inner có thể phát một response hoặc nhiều update; operator chuyển tiếp từng value của inner hiện tại, không chỉ lấy value cuối.

Hình dung bạn đổi kênh radio: sau khi chuyển kênh, bạn chỉ nghe kênh mới, không đợi kênh cũ phát xong. Nhưng tắt việc nghe không đồng nghĩa đài cũ ngừng phát. Tương tự, switchMap quản lý subscription của bạn; producer bên dưới có dừng công việc thật hay không là câu hỏi riêng.

Không dùng nó làm mặc định cho mọi HTTP request. Nếu mọi lệnh ghi đều phải được xử lý, bỏ theo dõi lệnh trước có thể làm mất thông tin thành công hoặc thất bại. Với lệnh ghi cần thứ tự, mình chọn concatMap; với các task độc lập cần chạy đủ, chọn mergeMap và xét giới hạn concurrency.

Cách chuyển giữa các inner

Chuyển ngay khi outer phát value

Trong RxJS 7.8.2, mỗi lần outer phát value, switchMap thực hiện theo thứ tự:

  1. Unsubscribe inner hiện tại, nếu có.
  2. Gọi project(value, index) để tạo ObservableInput mới.
  3. Subscribe kết quả projection và chuyển tiếp các value của inner này.

Nó không đợi inner mới phát response rồi mới bỏ inner cũ. Vì inner cũ bị unsubscribe trước khi projection mới chạy, nếu projection mới ném lỗi thì inner cũ cũng không được giữ lại như một phương án dự phòng.

Marble timeline
switchMap chuyển subscription sang input mới nhất
Mỗi vạch = 1 frame
switchMap chuyển subscription sang input mới nhất. outer: -a---b----|. inner A: -^-x-!. inner B: -----^-x---y|. output: ---x---x---y|

Không có hàng đợi cho input bị thay thế, cũng không có bước chạy lại inner cũ khi inner mới complete. Dòng inner A dừng tại unsubscribe, còn inner B tiếp tục ngay cả khi outer đã complete.

Không thu hồi value đã phát

Giả sử A đã phát A1, sau đó outer mới phát B. A1 vẫn là một value hợp lệ mà downstream đã nhận; operator chỉ ngăn emission tiếp theo của A đi qua subscription cũ. Nó không xóa state, thu hồi side effect hay sửa lại output trong quá khứ.

Sự kiện minh họa:

0 ms     outer A → subscribe A
100 ms   A phát A1 → output A1
150 ms   outer B → unsubscribe A → subscribe B
200 ms   A định phát A2 → không được chuyển tiếp qua subscription đã đóng
250 ms   B phát B1 → output B1

Output: A1, B1

Nếu UI cần xóa danh sách cũ ngay lúc đổi query, hãy phát trạng thái waiting hoặc loading từ inner mới. switchMap tự nó không làm UI trống đi trong thời gian chờ response.

API và ví dụ tối thiểu

Dạng API nên dùng trong code mới:

switchMap(project)
Thành phầnÝ nghĩa
project(value, index)Trả ObservableInput, thường là Observable hoặc Promise
valueValue mới nhận từ outer source
indexChỉ số outer emission, bắt đầu từ 0 trong mỗi subscription
OutputValue do inner hiện tại phát, không phải object Observable

switchMap không có tham số concurrency. mergeMap(project, 1) cũng không thay thế được nó: mergeMap với một slot giữ input mới trong hàng đợi, trong khi switchMap bỏ subscription cũ để xử lý input mới ngay.

import { finalize, map, of, switchMap, timer } from 'rxjs';

of('A', 'B').pipe(
  switchMap((id) => timer(100).pipe(
    map(() => `result:${id}`),
    finalize(() => console.log('finalize', id)),
  )),
).subscribe({
  next: (value) => console.log(value),
  complete: () => console.log('output complete'),
});

// Ngay khi subscribe: finalize A
// Khoảng 100 ms:      result:B
//                    output complete
//                    finalize B

Outer phát A rồi B đồng bộ, nên timer của A bị hủy trước khi kịp phát. Timer của B vẫn chạy dù of('A', 'B') đã complete. Các mốc thời gian là minh họa, không phải cam kết timer chạy chính xác tới từng mili giây.

Đừng rút ra quy tắc “switchMap chỉ phát input cuối”. Nếu thay timer(100).pipe(map(...)) bằng of(...), inner A phát và complete đồng bộ trước khi B tới; output có thể nhận kết quả của cả A và B.

RxJS 7 còn có overload resultSelector, nhưng overload đó đã deprecated. Khi cần giữ outer metadata, dùng map trong inner như { id, result } thay vì resultSelector. Và đừng gọi .subscribe() trong projection: nó trả Subscription, không phải ObservableInput, đồng thời tách công việc khỏi lifecycle mà operator đang quản lý.

Unsubscribe có thật sự hủy công việc

Cần phân biệt ba việc: không nhận kết quả cũ, dừng resource phía client, và ngăn tác động ở server. switchMap bảo đảm việc đầu tiên cho subscription do nó quản lý. Hai việc còn lại phụ thuộc producer và giao thức của hệ thống.

Promise không tự bị cancel

import { defer, of, switchMap } from 'rxjs';

function loadUser(id: number): Promise<{ id: number }> {
  console.log('start', id);
  return new Promise((resolve) => {
    setTimeout(() => {
      console.log('work finished', id);
      resolve({ id });
    }, 100);
  });
}

of(1, 2).pipe(
  switchMap((id) => defer(() => loadUser(id))),
).subscribe((user) => console.log('output', user.id));

// start 1
// start 2
// Khoảng 100 ms:
// work finished 1
// work finished 2
// output 2

defer trì hoãn lời gọi tới lúc subscribe, nhưng không thêm cơ chế cancel cho Promise. Khi chuyển sang ID 2, timer trong Promise của ID 1 vẫn chạy; chỉ kết quả của nó không được chuyển tới subscriber đã đóng. from(fetch(url)) cũng không tự bổ sung abort cho fetch.

Nếu Promise hoặc producer có side effect ghi trực tiếp vào state bên ngoài pipeline, side effect đó còn có thể xảy ra sau khi bị thay thế. Hãy đưa việc render vào consumer của output thay vì để request cũ tự cập nhật UI ngoài subscription.

HTTP cần producer có teardown

ajax của RxJS có teardown cho XHR; fromFetch dùng AbortController để abort request đang chờ khi unsubscribe. Với fromFetch, nên dùng selector nếu cần quản lý cả giai đoạn đọc response body:

import { fromFetch } from 'rxjs/fetch';

type SearchItem = { id: string; label: string };

function search$(query: string) {
  return fromFetch<SearchItem[]>(
    `/api/search?q=${encodeURIComponent(query)}`,
    {
      selector: (response) => {
        if (!response.ok) {
          throw new Error(`HTTP ${response.status}`);
        }
        return response.json() as Promise<SearchItem[]>;
      },
    },
  );
}

Giả sử endpoint trả JSON array đúng schema SearchItem[]. Kiểu TypeScript ở đây không validate dữ liệu runtime; thêm validation nếu backend không bảo đảm schema. fetch cũng không tự reject vì HTTP 404 hay 500, nên cần kiểm tra response.ok.

Không có selector, fromFetch có thể phát Response và complete ngay khi nhận headers. Một bước switchMap(response => response.json()) phía sau chỉ theo dõi Promise đọc body; unsubscribe lúc đó không bảo đảm abort phần body qua lifecycle của fromFetch đã complete. selector giữ việc đọc body trong cùng subscription có thể abort.

Abort không phải rollback

Server có thể đã nhận và xử lý request trước khi client abort. Đừng dùng switchMap để bảo đảm “chỉ lệnh ghi cuối cùng được thực hiện”. Thao tác ghi cần thiết kế thứ tự, idempotency, version hoặc transaction ở cấp backend tùy yêu cầu.

Tìm kiếm với cancellation ngay khi đổi query

Ví dụ dưới đây dùng search$ ở phần HTTP. Trong trình duyệt, cần có input mang ID search; mỗi subscription vào pipeline tạo listener và execution riêng. Chỉ subscribe một lần ở owner của màn hình và unsubscribe khi owner kết thúc.

Mục tiêu là hủy request cũ ngay khi query đổi, nhưng chỉ gửi request mới nếu người dùng ngừng gõ ít nhất 300 ms. Query ngắn hơn hai ký tự đưa UI về idle, kể cả khi request trước vẫn đang chờ.

import {
  catchError,
  distinctUntilChanged,
  fromEvent,
  map,
  of,
  startWith,
  switchMap,
  timer,
} from 'rxjs';

// Dùng SearchItem và search$ từ ví dụ HTTP ở trên.
type SearchState =
  | { status: 'idle'; query: string }
  | { status: 'waiting'; query: string }
  | { status: 'loading'; query: string }
  | { status: 'success'; query: string; items: SearchItem[] }
  | { status: 'error'; query: string; message: string };

const input = document.querySelector<HTMLInputElement>('#search');
if (!input) throw new Error('Thiếu input #search');

const query$ = fromEvent(input, 'input').pipe(
  map(() => input.value.trim()),
  startWith(input.value.trim()),
  distinctUntilChanged(),
);

const state$ = query$.pipe(
  switchMap((query) => {
    if (query.length < 2) {
      return of<SearchState>({ status: 'idle', query });
    }

    return timer(300).pipe(
      switchMap(() => search$(query).pipe(
        map((items): SearchState => ({
          status: 'success', query, items,
        })),
        startWith<SearchState>({ status: 'loading', query }),
      )),
      catchError((error: unknown) => of<SearchState>({
        status: 'error',
        query,
        message: error instanceof Error ? error.message : String(error),
      })),
      startWith<SearchState>({ status: 'waiting', query }),
    );
  }),
);

const subscription = state$.subscribe({
  next: (state) => console.log('render', state),
  error: (error: unknown) => console.error('pipeline error', error),
});

// Khi rời màn hình:
// subscription.unsubscribe();

Khi query hợp lệ đến, outer switchMap hủy toàn bộ inner trước đó: timer chờ hoặc request đang chạy. Inner mới phát waiting, đợi 300 ms rồi phát loading khi subscribe request. Sau đó UI nhận success hoặc error, luôn gắn với query của inner hiện tại.

Có hai switchMap vì có hai tầng lifecycle: query chọn một session tìm kiếm, còn timer bắt đầu request trong session đó. Request vẫn thuộc inner của outer switchMap, nên cancellation đi xuyên qua cả hai tầng. Thời gian 300 ms và độ dài hai ký tự chỉ là chính sách minh họa, cần đổi theo sản phẩm.

Vì sao đặt timer bên trong switchMap

Cách viết quen thuộc là query$.pipe(debounceTime(300), distinctUntilChanged(), switchMap(search$)). Nó phù hợp nếu “mới nhất” nghĩa là query đã qua debounce, nhưng không hủy request cũ ngay lúc input thô vừa đổi.

Giả sử request cho rx đang chạy, bạn vừa gõ thêm thành rxjs. Trong 300 ms chờ debounce, switchMap chưa nhận query mới, nên response của rx vẫn có thể được render. switchMap không sai; ranh giới latest của pipeline đang đặt sau debounce.

Nếu không muốn kết quả cũ xuất hiện trong khoảng chờ đó, đặt delay trong inner như ví dụ trên. Input mới tới outer switchMap ngay, làm mất cả timer hoặc request cũ. Mình chọn cách này cho UI yêu cầu state khớp query hiện tại; còn pipeline debounce trước đơn giản hơn khi khoảng chờ đó được chấp nhận.

Query rỗng cũng phải đi tới switchMap

Đừng đặt filter(query => query.length >= 2) trước switchMap nếu xóa query phải hủy request đang chạy. Value bị filter bỏ sẽ không tới operator, nên inner cũ vẫn được theo dõi và response có thể quay lại lấp đầy ô kết quả đã xóa.

Cho query rỗng đi qua và trả of({ status: 'idle', query }) vừa hủy inner cũ, vừa phát state để UI xóa danh sách. EMPTY cũng làm chuyển subscription và không phát gì, nhưng không tự thông báo UI về trạng thái rỗng.

Loading state và finalize

finalize trong inner chạy cả khi complete, error và bị thay thế. Vì vậy, nó phù hợp để cleanup resource hoặc log lifecycle, không phải để kết luận request đã thành công.

Một bẫy là đặt tap(() => loading = true) trước switchMap rồi cho inner cũ finalize(() => loading = false). Khi input mới tới, tap bật loading trước; ngay sau đó teardown inner cũ lại tắt loading của request mới. Ví dụ trên tránh biến boolean dùng chung: inner mới tự phát state sau khi inner cũ đã được teardown, còn inner cũ không tự sửa UI trong finalize.

Xử lý lỗi mà vẫn nghe input mới

Mặc định, lỗi từ outer, exception trong projection hoặc error của inner hiện tại đều kết thúc output. Nếu bạn muốn request lỗi nhưng người dùng vẫn gõ query khác được, đặt catchError bên trong inner của từng query.

import { catchError, of, switchMap, throwError } from 'rxjs';

function task$(query: string) {
  return query === 'bad'
    ? throwError(() => new Error('Request failed'))
    : of(`ok:${query}`);
}

// Recovery từng inner: query tiếp theo vẫn được xử lý.
of('bad', 'good').pipe(
  switchMap((query) => task$(query).pipe(
    catchError(() => of(`failed:${query}`)),
  )),
).subscribe(console.log);
// failed:bad
// ok:good

// Recovery toàn pipeline: upstream lỗi đã bị teardown.
of('bad', 'good').pipe(
  switchMap((query) => task$(query)),
  catchError(() => of('search unavailable')),
).subscribe(console.log);
// search unavailable

catchError bên ngoài thay toàn bộ upstream đã lỗi bằng fallback, không tiếp tục nghe outer cũ sau khi fallback complete. Trong ví dụ tìm kiếm, recovery nằm trong session của mỗi query nên event input tiếp theo vẫn hoạt động.

Nếu một factory có thể ném lỗi đồng bộ trước khi trả Observable, bọc nó bằng defer(() => factory(query)).pipe(catchError(...)); catchError gắn vào Observable mà factory chưa trả được sẽ không bắt exception từ lời gọi factory. Error của outer vẫn là phạm vi khác, không được inner recovery xử lý.

Đặt retry trong inner khi muốn retry request cho query hiện tại. Query mới sẽ hủy cả request lẫn lịch retry của inner cũ. retry ngoài switchMap có thể subscribe lại cả outer; việc đó không có nghĩa là phát lại đúng query gần nhất, nhất là khi outer là hot stream không replay.

Complete và cleanup của toàn pipeline

Outer complete vẫn đợi inner hiện tại

Sự kiệnHành vi của switchMap
Outer phát value mớiUnsubscribe inner cũ rồi subscribe inner mới
Inner hiện tại completeTiếp tục nghe outer; không quay lại inner cũ
Outer completeĐợi inner hiện tại complete, nếu còn inner
Outer và inner hiện tại đều completeOutput complete
Error chưa recoveryOutput error; teardown outer và inner
Consumer unsubscribeTeardown outer và inner, không gọi callback complete

Vì vậy, of('A').pipe(switchMap(() => interval(1_000))) không tự complete: outer đã xong nhưng inner vẫn phát mãi. Nếu task cần hữu hạn, đặt ranh giới phù hợp như take(1) hoặc take(n) trong inner; đừng ép complete khi inner thực sự là subscription realtime cần sống cùng view.

Đặt takeUntil sau switchMap

Để hủy cả outer lẫn inner khi owner bị destroy, đặt takeUntil sau flattening operator:

import {
  interval,
  of,
  Subject,
  switchMap,
  takeUntil,
} from 'rxjs';

const destroy$ = new Subject<void>();

of('active').pipe(
  switchMap(() => interval(1_000)),
  takeUntil(destroy$),
).subscribe({
  next: (value) => console.log(value),
  complete: () => console.log('owner stopped'),
});

setTimeout(() => {
  destroy$.next();
  destroy$.complete();
}, 100);

// Khoảng 100 ms: owner stopped
// Không có emission từ interval.

Notifier phải phát value; chỉ gọi destroy$.complete() không kích hoạt takeUntil. Nếu đặt takeUntil trước switchMap, nó chỉ làm outer complete, còn inner hiện tại vẫn có thể chạy tiếp. Với outer đồng bộ trong ví dụ này, outer thậm chí đã complete trước tín hiệu destroy.

Gọi trực tiếp subscription.unsubscribe() cũng dọn cả pipeline, nhưng không phát complete tới consumer. Sự khác nhau này giải thích vì sao cleanup cần finalize hoặc teardown chứ không nên chỉ đặt trong callback complete.

Các bẫy khi chọn latest

  • Tưởng chỉ có một network request trong cả ứng dụng. switchMap giữ một inner subscription hiện tại trong mỗi execution. Hai subscriber vào cold pipeline có thể tạo hai request riêng; Promise không cancel hoặc HTTP đã được server nhận còn có thể tiếp tục ở phía dưới.
  • Inner là hot source hoặc được share. Unsubscribe khỏi inner không bắt buộc dừng producer đang phục vụ consumer khác. Với shareReplay, xét refCount và các subscriber còn sống trước khi suy luận resource đã dừng.
  • Polling nhanh hơn response. interval(1_000).pipe(switchMap(load$)) liên tục thay request nếu response luôn mất hơn một giây. Có thể không nhận kết quả nào; nếu mỗi lượt phải chạy xong thì chọn exhaustMap, hoặc thiết kế lịch dựa trên completion.
  • Side effect nằm trong projection. switchMap vẫn gọi projection cho mỗi outer emission. Không nhận output không có nghĩa side effect chưa chạy; tránh mutation state ngoài pipeline và dùng producer có teardown nếu cần cancel thật.
  • Latest chỉ tính tới outer value operator nhận. filter, debounceTime hoặc distinctUntilChanged có thể chặn value trước operator. Trong ví dụ tìm kiếm, query giống nhau sau trim cố ý không restart request; muốn refresh cùng query cần event refresh riêng.

Các bẫy này đều quay về một câu hỏi: “inner subscription cũ đóng” đã đủ cho nghiệp vụ chưa? Nếu chưa, cần thêm chính sách ở producer, lifecycle của owner hoặc backend, chứ không chỉ thay tên operator.

So sánh các flattening operator

OperatorInput mới khi đang bậnMình chọn khi
switchMapHủy subscription cũ, xử lý input mớiKết quả cũ hết hữu ích khi input đổi
concatMapXếp hàng input mớiCần xử lý đủ và giữ thứ tự công việc
mergeMapChạy thêm inner nếu còn slotCông việc độc lập, cần mọi kết quả
exhaustMapBỏ input mới cho tới khi inner completeKhông muốn click submit lặp tạo thêm task

Với autosave, đừng chọn switchMap chỉ vì “nội dung cuối cùng là mới nhất”. Client bỏ response cũ không ngăn lệnh save cũ tới server muộn và ghi đè bản mới. Nếu chỉ quan tâm snapshot cuối nhưng backend có thể nhận request chồng lấp, cần version hoặc cơ chế concurrency phù hợp; nếu cần tuần tự hóa thao tác ở client, xét concatMap.

Kiểm chứng cancellation bằng marble test

Test không chỉ nên kiểm tra output. Nếu yêu cầu là inner cũ bị unsubscribe, hãy assert subscription window của nó. Ví dụ dùng TestScheduler.run, với mỗi dấu - tương ứng một frame virtual time:

import { deepStrictEqual } from 'node:assert/strict';
import { switchMap } from 'rxjs';
import { TestScheduler } from 'rxjs/testing';

const scheduler = new TestScheduler((actual, expected) => {
  deepStrictEqual(actual, expected);
});

scheduler.run(({ cold, expectObservable, expectSubscriptions }) => {
  const outer$ = cold('-a---b----|');
  const innerA$ = cold('--x---y|');
  const innerB$ = cold('--x---y|');

  const result$ = outer$.pipe(
    switchMap((value) => value === 'a' ? innerA$ : innerB$),
  );

  expectObservable(result$).toBe('---x---x---y|');
  expectSubscriptions(outer$.subscriptions).toBe('^---------!');
  expectSubscriptions(innerA$.subscriptions).toBe('-^---!');
  expectSubscriptions(innerB$.subscriptions).toBe('-----^------!');
});

A bắt đầu ở frame 1, phát x ở frame 3, rồi bị unsubscribe ở frame 5 khi B tới. y của A dự kiến ở frame 7 nhưng không được chuyển tiếp. B bắt đầu ở frame 5, phát x ở 7, y ở 11 và complete ở 12; outer đã complete ở frame 10 nhưng output vẫn đợi B.

^ và ! thuộc subscription marble, lần lượt là subscribe và unsubscribe; | thuộc notification marble, là complete. Test này xác nhận lifecycle RxJS, không chứng minh HTTP bị abort thật. Muốn kiểm tra abort, thêm integration test với producer hoặc mock fetch có quan sát AbortSignal.

Bài tập và checklist

  1. Thay timer của ví dụ tối thiểu bằng of(...). Dự đoán vì sao A vẫn phát kết quả dù outer phát B ngay sau đó.
  2. Trong ví dụ Promise, log từ công việc của ID 1 vẫn xuất hiện nhưng output chỉ có ID 2. Chỉ rõ đâu là producer, đâu là consumer đã bị unsubscribe.
  3. Với UI tìm kiếm, đang chờ request thì xóa input. Kiểm tra state về idle, request cũ không render và subscription còn nghe query mới.
  4. Trong marble test, thay switchMap bằng mergeMap. Dự đoán thêm emission của A và sửa subscription expectation trước khi chạy.

Trước khi áp dụng vào pipeline thật, trả lời bốn câu:

  • Input mới có thật sự khiến công việc cũ không còn cần thiết không?
  • Producer có teardown để dừng công việc, hay chỉ bỏ kết quả?
  • Ranh giới latest nằm trước hay sau debounce, filter và deduplication?
  • Error recovery và cleanup có nằm đúng cấp của inner và owner không?

Nếu câu đầu tiên là “mọi thao tác đều phải chạy”, đổi chiến lược trước khi tối ưu cancellation. Nếu câu đầu tiên là “chỉ kết quả hiện tại cần render”, hãy viết một test đổi input khi inner cũ chưa xong và assert cả output lẫn subscription window.

Học tiếp

Nguồn tham khảo

On this page