Higher-order Observables
Làm chủ Observable lồng nhau và các chiến lược flattening.
Bạn gõ thêm một ký tự trong ô tìm kiếm khi request trước vẫn chưa xong. Hoặc bạn bấm nút lưu hai lần liên tiếp. Cả hai đều là “event tạo ra công việc async”, nhưng cách xử lý đúng lại khác nhau: tìm kiếm thường chỉ cần kết quả mới nhất, còn thao tác ghi dữ liệu có thể cần giữ đủ mọi công việc và đúng thứ tự.
Higher-order Observables giúp mô tả công việc đó bằng stream. Điều cần chọn không chỉ là operator có chữ Map, mà là khi event mới đến, bạn muốn làm gì với công việc đang chạy? Trang này cho bạn mental model và cách chọn chiến lược; các bài con sẽ đi sâu vào từng operator.
Trước khi bắt đầu
Bạn nên biết pipe(), map(), subscribe() và sự khác nhau giữa complete với unsubscribe. Các ví dụ dùng public API của RxJS 7.x, đối chiếu với mã nguồn 7.8.2. Ví dụ timer là mô phỏng, không phải HTTP request thật.
Mục lục
- Bạn sẽ học gì trong nhóm này
- Outer stream và inner stream
- Bốn chiến lược flattening
- Một ví dụ để so sánh cả bốn
- Lifecycle của pipeline lồng nhau
- Những bẫy dễ gặp
- Lộ trình học đề xuất
- Bài tập tự kiểm tra
- Học tiếp
- Nguồn tham khảo
Bạn sẽ học gì trong nhóm này
Sau nhóm bài này, bạn nên đọc được một pipeline và trả lời bốn câu hỏi:
- Event nào thuộc outer stream, và nó tạo ra inner stream nào?
- Nếu hai công việc chồng lấp, pipeline xếp hàng, chạy đồng thời, hủy cái cũ hay bỏ cái mới?
- Kết quả có thể đến khác thứ tự event ban đầu không, và UI có chấp nhận điều đó không?
- Khi error, complete hoặc unsubscribe xảy ra, inner nào còn sống và resource nào được dọn?
Mình khuyên bạn chọn operator từ yêu cầu nghiệp vụ trước, rồi mới viết pipeline. “Cứ dùng switchMap cho HTTP” không phải quy tắc an toàn: đọc kết quả tìm kiếm và ghi một đơn hàng không có cùng cancellation policy.
Outer stream và inner stream
Hình dung mỗi event là một phiếu công việc. Phiếu nói cần làm gì, còn inner Observable mô tả công việc đó sẽ phát kết quả thế nào. Flattening operator quyết định khi nào nhận phiếu, khi nào bắt đầu xử lý và kết quả nào được chuyển tiếp.
Ẩn dụ này chỉ giúp nhìn ra hai tầng. Observable không phải một hàng đợi bền vững: dữ liệu chờ trong RAM không tự sống qua restart và không có đảm bảo exactly-once.
| Tầng | Phát ra gì? | Ví dụ tìm kiếm |
|---|---|---|
| Outer stream | Event hoặc dữ liệu đầu vào | Chuỗi query từ ô input |
| Inner stream | Kết quả của một công việc | Observable trả kết quả cho một query |
| Higher-order stream | Các inner Observable | Observable<Observable<SearchResult>> |
| Flattened stream | Giá trị do inner phát | Observable<SearchResult> |
Vì sao map chưa đủ
map() chỉ biến đổi từng giá trị. Nếu callback trả về Observable, kết quả vẫn là một stream phát ra Observable; map() không tự subscribe inner.
import { map, of } from 'rxjs';
import type { Observable } from 'rxjs';
const ids$ = of(1, 2);
const users$$: Observable<Observable<string>> = ids$.pipe(
map((id) => of(`user-${id}`)),
);
users$$.subscribe((inner$) => {
console.log('có phương thức subscribe:', typeof inner$.subscribe);
});Output là hai dòng có phương thức subscribe: function, không phải user-1, user-2. Dấu $$ ở đây chỉ là quy ước đặt tên để nhắc rằng có hai tầng; RxJS không đọc tên biến để quyết định hành vi.
Flattening nối hai tầng thế nào
Flattening operator subscribe inner theo một policy rồi chuyển các notification của inner xuống consumer. Nếu muốn xử lý từng ID tuần tự, bạn có thể viết:
import { concatMap, of } from 'rxjs';
const users$ = of(1, 2).pipe(
concatMap((id) => of(`user-${id}`)),
);
users$.subscribe({
next: (user) => console.log(user),
complete: () => console.log('complete'),
});user-1
user-2
completeouter: ID ──► project(ID) ──► inner Observable
│
flattening operator
quản lý subscribe/teardown
│
▼
kết quả ──► consumerVề mental model, concatMap(project) kết hợp bước tạo inner với bước flatten tuần tự. Nếu đã có sẵn stream phát ra các Observable, bạn dùng nhóm concatAll, mergeAll, switchAll, exhaustAll tương ứng. Chẳng hạn map(project) rồi switchAll() mô tả hai bước mà switchMap(project) gộp lại.
Đừng suy ra rằng mọi tổ hợp map(project) rồi operator All đều gọi project ở cùng thời điểm với operator Map. Trong RxJS 7.8.2, concatMap và mergeMap có concurrency hữu hạn có thể trì hoãn gọi project cho giá trị đang chờ; exhaustMap không gọi project cho event bị bỏ qua. Điều này quan trọng nếu callback tạo Promise hoặc side effect eager.
Bốn chiến lược flattening
Tất cả bốn operator đều tạo và subscribe inner. Điểm khác nhau xuất hiện khi outer phát thêm một giá trị trong lúc inner trước chưa complete.
| Operator | Khi giá trị mới đến mà đang bận | Số inner active | Thứ tự kết quả | Trường hợp thường phù hợp |
|---|---|---|---|---|
concatMap | Xếp giá trị vào buffer để xử lý sau | Tối đa 1 | Theo thứ tự công việc được nhận | Ghi dữ liệu tuần tự, chuỗi thao tác phụ thuộc |
mergeMap | Bắt đầu inner mới nếu còn slot; nếu hết slot thì buffer | Tối đa concurrent, mặc định Infinity | Có thể xen kẽ, không đảm bảo thứ tự outer | Công việc độc lập, tải nhiều tài nguyên |
switchMap | Unsubscribe inner cũ, subscribe inner mới | Tối đa 1 | Chỉ chuyển tiếp từ inner hiện đang được chọn | Tìm kiếm, đọc dữ liệu theo lựa chọn hiện tại |
exhaustMap | Bỏ qua giá trị mới, không xếp hàng | Tối đa 1 | Chỉ có kết quả từ các công việc được nhận | Chặn submit lặp khi một submit đang chạy |
“Active” trong bảng là subscription mà operator quản lý. Nếu producer không hỗ trợ cancellation, công việc bên ngoài của một inner đã bị unsubscribe vẫn có thể tiếp tục.
concatMap giữ đủ và giữ thứ tự
Nếu mỗi bản cập nhật cần được xử lý sau bản trước, mình chọn concatMap. Nó chờ inner hiện tại complete, không chỉ chờ inner phát một giá trị, rồi mới xử lý giá trị tiếp theo trong buffer.
Đổi lại, công việc chậm sẽ chặn toàn bộ phần sau. Nếu outer phát nhanh hơn tốc độ xử lý, buffer có thể tăng liên tục; nếu inner không bao giờ complete, phần còn lại sẽ chờ mãi. concatMap phù hợp với một nguồn công việc có tốc độ hoặc số lượng được kiểm soát, không phải một durable job queue.
mergeMap cho phép công việc chồng lấp
Nếu các công việc độc lập và thứ tự kết quả không quan trọng, mergeMap cho phép nhiều inner cùng active. Ví dụ, tải thông tin cho nhiều ID có thể dùng mergeMap(loadById, 4) để giới hạn tối đa bốn inner subscription tại một thời điểm; số 4 là cấu hình minh họa, cần chọn theo khả năng của hệ thống.
“Chạy đồng thời” ở đây không có nghĩa RxJS tạo thread mới. Các request hoặc timer có thể chồng lấp thời gian chờ, nhưng code JavaScript đồng bộ vẫn chạy theo execution model của runtime.
Giới hạn concurrency cũng không giới hạn độ dài buffer: các giá trị đến khi đã đủ slot vẫn được giữ lại. Trong RxJS 7.x, mergeMap(project, 1) có cùng chiến lược tuần tự với concatMap(project).
switchMap chỉ theo inner mới nhất
Với kết quả tìm kiếm theo query hiện tại, mình chọn switchMap: query mới khiến subscription của query cũ bị hủy. Điều này tránh việc một response cũ đến muộn ghi đè response mới qua chính pipeline đó.
Nhưng switchMap không xóa những giá trị đã phát trước khi chuyển inner. Nếu inner cũ đã cập nhật UI, UI vẫn giữ trạng thái đó cho đến khi code của bạn cập nhật tiếp. Operator cũng không đảm bảo chỉ phát một giá trị: inner hiện tại có thể phát nhiều lần.
Unsubscribe không phải rollback
switchMap hủy subscription cũ, không đảm bảo server ngừng xử lý hay hoàn tác thao tác ghi. Với from(existingPromise), Promise nền vẫn chạy. Đừng dùng nó cho một chuỗi lệnh ghi mà mọi lệnh đều phải được xử lý chỉ vì muốn tránh kết quả cũ trên UI.
exhaustMap giữ công việc đang chạy
Với nút submit mà click lặp trong lúc đang gửi phải bị bỏ qua, exhaustMap diễn đạt policy đó trực tiếp. Nó nhận event khi đang rảnh; từ lúc inner bắt đầu đến khi inner complete, các event mới bị bỏ qua.
Nó không nhớ “click cuối cùng” để chạy sau. Nếu inner là một stream sống vô hạn, operator sẽ tiếp tục bỏ qua mọi event tiếp theo. Với thao tác quan trọng, nên kết hợp trạng thái nút trên UI và idempotency phía server: exhaustMap chỉ chặn lặp trong phạm vi subscription hiện tại, không chặn request từ tab hoặc client khác.
Một ví dụ để so sánh cả bốn
Giả sử outer phát A, B, C đồng bộ, rồi complete. Mỗi công việc phát đúng một kết quả và complete sau một khoảng chờ mô phỏng: A chờ 300 ms, B chờ 100 ms, C chờ 200 ms. Khoảng chờ tính từ lúc inner được subscribe, không phải từ lúc outer phát nhãn.
import {
concatMap,
exhaustMap,
map,
mergeMap,
of,
switchMap,
timer,
} from 'rxjs';
import type { OperatorFunction } from 'rxjs';
type Job = 'A' | 'B' | 'C';
const duration: Record<Job, number> = {
A: 300,
B: 100,
C: 200,
};
const work = (job: Job) => timer(duration[job]).pipe(
map(() => job),
);
const strategies: Record<string, OperatorFunction<Job, Job>> = {
concatMap: concatMap(work),
mergeMap: mergeMap(work),
switchMap: switchMap(work),
exhaustMap: exhaustMap(work),
};
const jobs: Job[] = ['A', 'B', 'C'];
for (const [name, strategy] of Object.entries(strategies)) {
of(...jobs).pipe(strategy).subscribe({
next: (job) => console.log(`${name}: ${job}`),
complete: () => console.log(`${name}: complete`),
});
}Bốn subscription độc lập nên log giữa các chiến lược có thể xen kẽ. Đọc riêng log của từng chiến lược sẽ thấy:
| Chiến lược | Inner được bắt đầu | Kết quả | Thời điểm mô phỏng của kết quả |
|---|---|---|---|
concatMap | A, rồi B, rồi C | A → B → C | 300 → 400 → 600 ms |
mergeMap | A, B, C ngay trong lượt phát đồng bộ | B → C → A | 100 → 200 → 300 ms |
switchMap | A rồi hủy A; B rồi hủy B; giữ C | C | 200 ms |
exhaustMap | Chỉ A; bỏ B và C | A | 300 ms |
Đây là thời gian lý tưởng để suy luận, không phải cam kết chính xác của wall clock. Callback timer có thể chạy muộn khi runtime bận. Muốn kiểm tra timing có tính quyết định, dùng TestScheduler.
Outer tại t≈0: A → B → C → complete
concatMap: [ A 300 ms ][ B 100 ms ][ C 200 ms ] → complete
mergeMap: [ A 300 ms ] → complete
[ B 100 ms ]
[ C 200 ms ]
switchMap: hủy A, hủy B, [ C 200 ms ] → complete
exhaustMap: [ A 300 ms ], bỏ B và C → completeVì inner là timer async nên chúng vẫn active khi B và C đến. Nếu đổi work thành of(job) phát và complete ngay, cả bốn đều có thể cho A → B → C: lúc event tiếp theo đến, inner trước đã xong. Đó là lý do ví dụ thuần đồng bộ thường không thể hiện khác biệt giữa các flattening policy.
Lifecycle của pipeline lồng nhau
Chọn được operator vẫn chưa đủ. Một pipeline có outer, inner và consumer; phải nhìn cả ba khi quyết định cách kết thúc.
Outer complete không có nghĩa inner dừng ngay
Trong ví dụ trên, of() complete ngay sau C, nhưng flattened stream vẫn chờ công việc được nhận:
concatMapchờ inner active và xử lý hết buffer;mergeMapchờ mọi inner active, đồng thời xử lý hết buffer nếu có giới hạn concurrency;switchMapchờ inner hiện tại;exhaustMapchờ inner đã nhận, không xử lý lại những event bị bỏ qua.
Nếu outer không complete, inner complete cũng không tự làm cả pipeline complete: pipeline còn chờ event outer tiếp theo. Ngược lại, nếu outer đã complete nhưng inner cần chờ không bao giờ complete, flattened stream cũng không complete.
Error chưa được xử lý từ outer, inner hoặc callback project sẽ kết thúc flattened stream bằng error và teardown các subscription được nó quản lý. Consumer gọi unsubscribe() cũng hủy outer cùng các inner active, nhưng không gọi callback complete của consumer. Những giá trị còn trong buffer không được chạy tiếp sau cancellation đó.
Đặt error boundary theo phạm vi cần phục hồi
Nếu một request lỗi nhưng các event sau vẫn phải được xử lý, đặt catchError bên trong inner. Ví dụ này dùng hai query đồng bộ để nhìn rõ việc query sau vẫn được nhận khi query đầu lỗi:
import {
catchError,
defer,
map,
of,
switchMap,
throwError,
} from 'rxjs';
const search = (query: string) => defer(() => {
if (query === 'fail') {
return throwError(() => new Error('Lỗi mô phỏng'));
}
return of([`Kết quả cho ${query}`]);
});
of('fail', 'rxjs').pipe(
switchMap((query) => search(query).pipe(
map((items) => ({ kind: 'success' as const, query, items })),
catchError((error: unknown) => of({
kind: 'error' as const,
query,
message: String(error),
})),
)),
).subscribe((result) => console.log(result.kind, result.query));error fail
success rxjsNếu chuyển catchError ra sau switchMap và trả một of(...) fallback, error của query đầu sẽ đóng upstream; fallback phát rồi complete, không giữ nguyên subscription để nhận query sau. Vị trí bên ngoài vẫn đúng nếu bạn chủ đích kết thúc cả workflow và thay nó bằng một stream khác. Xem catchError để chọn recovery policy.
Đặt cancellation ở đúng tầng
Nếu dùng takeUntil(destroy$) để dừng toàn bộ pipeline khi màn hình bị tháo, mặc định nên đặt nó sau flattening operator:
queries$ → switchMap(search) → takeUntil(destroy$) → consumerKhi destroy$ phát, takeUntil complete output và unsubscribe upstream, bao gồm cả inner đang active. Nếu đặt nó trước switchMap, nó chỉ làm outer complete; switchMap vẫn có thể chờ inner hiện tại hoàn thành. concatMap thậm chí có thể tiếp tục xử lý buffer đã nhận.
Kiểm tra teardown của producer
Dừng pipeline chỉ dừng được resource khi source có teardown tương ứng. Timer của RxJS có thể bị hủy; Promise thường không thể. Request muốn abort thật sự cần adapter hỗ trợ cancellation. Và kể cả abort phía client, bạn vẫn không được giả định server đã rollback.
Những bẫy dễ gặp
| Bẫy | Vì sao sai | Cách xử lý |
|---|---|---|
subscribe() bên trong subscribe() rồi chỉ hủy subscription ngoài | Inner tự tạo không tự thuộc lifecycle của outer | Trả Observable từ callback và dùng flattening operator; giữ một subscription ở boundary |
Dùng mergeMap cho UI chỉ cần lựa chọn hiện tại | Response cũ có thể đến sau và ghi đè state mới | Dùng switchMap cho luồng đọc phù hợp, hoặc kiểm tra request/version ID trước khi cập nhật |
Dùng switchMap cho mọi thao tác ghi | Một phần workflow có thể bị bỏ dở ở phía client, trong khi server đã xử lý | Xác định rõ yêu cầu giữ đủ, tuần tự, đồng thời hay bỏ lặp trước khi chọn operator |
Xem mergeMap(..., 4) là backpressure đầy đủ | Giới hạn inner active nhưng không giới hạn buffer, không làm producer outer chậm lại | Kiểm soát tốc độ đầu vào, lượng dữ liệu hoặc dùng cơ chế queue phù hợp |
| Tạo tất cả Promise trước khi flatten | Promise đã bắt đầu trước lúc operator quyết định subscribe | Tạo công việc trong project, hoặc dùng defer khi cần trì hoãn đến lúc subscribe |
Inner không complete với concatMap hoặc exhaustMap | Không có thời điểm để bắt đầu công việc chờ hoặc nhận event mới | Chọn nguồn hữu hạn, take(...) hoặc timeout phù hợp nghiệp vụ |
Gắn finalize rồi hiểu mọi lần gọi là thành công | finalize chạy cả khi error hoặc unsubscribe | Dùng nó cho cleanup; ghi nhận thành công qua kết quả hoặc protocol của tác vụ |
Đừng chỉ quan sát output khi debug. Hãy log lúc inner bắt đầu, lúc phát kết quả và lúc finalize chạy; một inner không phát kết quả có thể chưa từng được nhận, đang xếp hàng, bị hủy hoặc bị lỗi. Các trường hợp đó cần cách sửa khác nhau.
Lộ trình học đề xuất
Đọc theo thứ tự dưới đây để mỗi bài trả lời một quyết định cụ thể:
| Bài | Câu hỏi cần trả lời |
|---|---|
| Higher-order stream | Outer và inner là gì; khi nào map tạo thêm một tầng Observable? |
| concatMap | Làm sao giữ thứ tự, và vì sao inner phải complete để queue đi tiếp? |
| mergeMap | Khi nào cho phép chồng lấp và kiểm soát concurrency thế nào? |
| switchMap | Kết quả nào bị bỏ và cancellation thực sự dừng được gì? |
| exhaustMap | Khi nào nên bỏ event mới thay vì hủy công việc cũ hoặc xếp hàng? |
| Race condition | Response đến sai thứ tự có làm state sai không, và cần bảo vệ ở tầng nào? |
Bài tập tự kiểm tra
Dùng ví dụ A, B, C phía trên, dự đoán trước rồi mới chạy:
- Đổi
mergeMap(work)thànhmergeMap(work, 2)và cho C chờ 150 ms thay vì 200 ms. B xong ở 100 ms, C mới được bắt đầu; output lý tưởng là B → C → A ở 100 → 250 → 300 ms. Log thời điểm bắt đầu để xác nhận không quá hai inner active. - Đổi
workthànhof(job). Giải thích vì saoexhaustMapkhông bỏ B, C vàswitchMapkhông làm mất A, B. - Đổi
workthànhinterval(100). Với outer vẫn làof('A', 'B', 'C'), operator nào còn phát kết quả, operator nào chặn công việc sau, và vì sao pipeline không complete? Chủ động unsubscribe sau một khoảng thử nghiệm. - Chuyển
catchErrortrong ví dụ query ra ngoàiswitchMap. Xác nhận vì saosuccess rxjsbiến mất khi fallback là mộtof(...)hữu hạn. - Thêm log
finalizevào inner A, B, C. Phân biệt inner đã complete, inner bịswitchMaphủy và inner chưa từng đượcexhaustMapnhận.
Viết lại một pipeline bạn đang dùng bằng một câu nghiệp vụ: “khi event mới đến, công việc cũ phải…”. Nếu câu đó chưa rõ, đừng chốt operator chỉ dựa vào tên.
Học tiếp
Bắt đầu với bài Higher-order stream nếu bạn vẫn lẫn giá trị đầu vào với Observable trả về từ callback. Nếu đã rõ hai tầng, chọn bài operator khớp với policy bạn cần và kiểm tra lifecycle bằng marble test.
Higher-order stream
Nắm rõ hai tầng trước khi chọn flattening policy.
concatMap
Xử lý đủ công việc theo thứ tự.
mergeMap
Cho phép công việc độc lập chồng lấp có kiểm soát.
switchMap
Theo inner mới nhất và hiểu giới hạn cancellation.
exhaustMap
Bỏ event lặp trong lúc công việc hiện tại chưa xong.
Race condition
Bảo vệ state trước kết quả đến sai thứ tự.
Subscription và teardown
Kiểm tra resource có thực sự được dọn khi hủy hay không.
Marble testing
Kiểm thử output và subscription timing bằng virtual time.