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
- Cách chuyển giữa các inner
- API và ví dụ tối thiểu
- Unsubscribe có thật sự hủy công việc
- Tìm kiếm với cancellation ngay khi đổi query
- Xử lý lỗi mà vẫn nghe input mới
- Complete và cleanup của toàn pipeline
- Các bẫy khi chọn latest
- So sánh các flattening operator
- Kiểm chứng cancellation bằng marble test
- Bài tập và checklist
- Học tiếp
- Nguồn tham khảo
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ự:
- Unsubscribe inner hiện tại, nếu có.
- Gọi
project(value, index)để tạoObservableInputmới. - 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.
B đến ở frame 5 nên inner A bị unsubscribe ngay; emission y về sau của A không thể đi xuống output.
- a
- input A
- b
- input B mới hơn
- x
- emission đầu của inner
- y
- emission tiếp theo
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, B1Nế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 |
value | Value mới nhận từ outer source |
index | Chỉ số outer emission, bắt đầu từ 0 trong mỗi subscription |
| Output | Value 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 BOuter 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 2defer 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 unavailablecatchError 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ện | Hành vi của switchMap |
|---|---|
| Outer phát value mới | Unsubscribe inner cũ rồi subscribe inner mới |
| Inner hiện tại complete | Tiế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 complete | Output complete |
| Error chưa recovery | Output error; teardown outer và inner |
| Consumer unsubscribe | Teardown 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.
switchMapgiữ 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étrefCountvà 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ọnexhaustMap, hoặc thiết kế lịch dựa trên completion. - Side effect nằm trong projection.
switchMapvẫ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,debounceTimehoặcdistinctUntilChangedcó thể chặn value trước operator. Trong ví dụ tìm kiếm, query giống nhau sautrimcố ý 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
| Operator | Input mới khi đang bận | Mình chọn khi |
|---|---|---|
switchMap | Hủy subscription cũ, xử lý input mới | Kết quả cũ hết hữu ích khi input đổi |
concatMap | Xếp hàng input mới | Cần xử lý đủ và giữ thứ tự công việc |
mergeMap | Chạy thêm inner nếu còn slot | Công việc độc lập, cần mọi kết quả |
exhaustMap | Bỏ input mới cho tới khi inner complete | Khô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
- 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 đó. - 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.
- 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. - Trong marble test, thay
switchMapbằngmergeMap. 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
Race condition
Nhận diện kết quả cũ ghi đè state mới và giới hạn của latest-only.
concatMap
Giữ mọi công việc và xử lý lần lượt theo thứ tự.
exhaustMap
Bỏ input mới thay vì hủy task đang chạy.
Subscription và teardown
Hiểu resource được dọn khi complete, error hoặc unsubscribe.
Debounce, throttle và audit
Chọn chính sách theo thời gian và vị trí trong pipeline.
Marble testing
Kiểm thử emission và subscription bằng virtual time.
Nguồn tham khảo
- RxJS 7.8.2 —
switchMap— API, thứ tự unsubscribe/project/subscribe và điều kiện output complete. - RxJS 7.8.2 —
fromFetch—AbortController,selectorvà cancellation khi đọc response body. - RxJS 7.8.2 —
takeUntil— notifier phát value mới kích hoạt completion. - RxJS 7.8.2 — Testing RxJS Code with Marble Diagrams — virtual time, notification marble và subscription marble.