Higher-order stream
Hiểu stream phát ra các Observable con.
Bạn có một stream phát ID người dùng, và mỗi ID cần gọi một API trả về Observable. Dùng map để đổi ID thành request nghe có vẻ đúng, nhưng khi subscribe, bạn lại nhận một Observable chứ chưa nhận được dữ liệu người dùng. Đây là lúc cần hiểu higher-order stream: stream phát ra các stream con, cùng câu hỏi quan trọng hơn là ai sẽ subscribe các stream con đó.
Trước khi bắt đầu
Bạn nên biết Observable, subscribe, pipe và map. Ví dụ dùng API của RxJS 7.8.2, import từ rxjs và có thể chạy trong playground đã có RxJS. Các request trong bài là mô phỏng bằng defer, of và timer, không gọi backend thật.
Mục lục
- Từ một value đến một stream con
- map tạo Observable lồng nhưng không subscribe inner
- Flatten nối hai tầng subscription
- Value mới đến khi inner cũ chưa xong
- Lifecycle của pipeline sau khi flatten
- Vì sao không mặc định dùng subscribe lồng
- Observable lồng và Promise lồng
- Khi nào cần giữ higher-order stream
- Bài tập tự kiểm tra
- Học tiếp
- Nguồn tham khảo
Từ một value đến một stream con
Một stream thông thường có thể phát number, string hoặc object. Higher-order Observable cũng làm đúng việc đó, chỉ khác là value được phát chính là một Observable.
import { Observable, map, of } from 'rxjs';
type User = { id: number; name: string };
function loadUser$(id: number): Observable<User> {
return of({ id, name: `User ${id}` });
}
const ids$: Observable<number> = of(1, 2);
const requests$: Observable<Observable<User>> = ids$.pipe(
map((id) => loadUser$(id)),
);Đọc type từ ngoài vào trong: requests$ phát các Observable<User>, còn mỗi Observable con có thể phát User. Downstream của requests$ chưa trực tiếp nhận User.
Bạn có thể hình dung outer đưa ra từng tờ phiếu công việc, còn inner mô tả cách làm công việc trên phiếu. Nhận phiếu không đồng nghĩa với đã làm xong việc. Tuy nhiên, phép ví von này có giới hạn: inner không nhất thiết là công việc chưa bắt đầu; nó cũng có thể là hot stream đang phát dữ liệu ở nơi khác.
Outer và inner là vai trò tương đối
- Outer Observable: stream phát ra các Observable con. Trong ví dụ, đó là
requests$. - Inner Observable: từng Observable do
loadUser$(id)trả về. - Flattened output: stream phát value của inner sau khi một flattening operator quản lý subscription.
ids$ : Observable<number>
│ map(loadUser$)
▼
requests$ : Observable<Observable<User>> ← outer
│ mỗi next mang một Observable<User>
▼
inner A$ inner B$
│ │
└────┬─────┘
│ flatten theo một chiến lược
▼
users$ : Observable<User>“Outer” và “inner” không phải hai class khác nhau trong RxJS. Chúng chỉ mô tả vai trò trong một pipeline. Nếu lại bọc requests$ trong một Observable nữa, nó có thể trở thành inner ở tầng mới.
Higher-order không đồng nghĩa với bất đồng bộ
Ví dụ of(of(1, 2), of(3)) là higher-order stream dù cả outer lẫn inner đều có thể phát đồng bộ. Ngược lại, timer(1_000) phát một số sau một khoảng thời gian, nhưng vẫn là Observable<number>, không phải higher-order.
Vì vậy, hãy nhìn type của emission, không chỉ nhìn việc có API request hay timer. Một array như [1, 2] cũng chỉ là một value thông thường nếu được phát bởi of([1, 2]); nó không tự trở thành stream con.
map tạo Observable lồng nhưng không subscribe inner
map gọi projection rồi phát nguyên kết quả mà projection trả về. Nếu kết quả là Observable, map phát Observable đó; nó không tự “mở” Observable để lấy value bên trong.
Đoạn sau làm rõ khác biệt giữa tạo inner và thực thi inner cold:
import { Observable, defer, map, of } from 'rxjs';
type User = { id: number; name: string };
function loadUser$(id: number): Observable<User> {
console.log('tạo inner:', id);
return defer(() => {
console.log('subscribe inner:', id);
return of({ id, name: `User ${id}` });
});
}
const requests$ = of(1, 2).pipe(
map((id) => loadUser$(id)),
);
requests$.subscribe({
next: () => console.log('outer phát một Observable'),
complete: () => console.log('outer complete'),
});
// tạo inner: 1
// outer phát một Observable
// tạo inner: 2
// outer phát một Observable
// outer completeKhông có dòng subscribe inner: vì chưa ai subscribe các Observable con. defer đặt việc tạo dữ liệu bên trong subscription; đây là cách kiểm chứng rõ ràng hơn việc chỉ log object Observable.
Lazy execution phụ thuộc producer
map không subscribe inner, nhưng callback của bạn vẫn chạy khi outer phát value. Nếu callback gọi fetch() ngay hoặc làm side effect trước khi trả Observable, công việc đó vẫn bắt đầu. Dùng producer cold hoặc defer nếu muốn trì hoãn việc khởi tạo tới lúc inner được subscribe.
Flatten nối hai tầng subscription
Flattening operator nhận các inner, quyết định khi nào subscribe chúng, rồi chuyển tiếp value từ inner xuống một output stream. Nó quản lý quan hệ giữa subscription outer, các subscription inner và subscription downstream.
“Flatten” ở đây không phải chờ gom mọi kết quả thành một array. Inner phát value nào thì value đó có thể đi xuống ngay, tùy chiến lược đang dùng; một inner có thể phát không value, một value hoặc rất nhiều value.
Tách bước map và flatten
Dùng concatAll() để tiêu thụ từng inner theo thứ tự, đợi inner trước complete rồi mới subscribe inner sau:
import { concatAll, map, of } from 'rxjs';
of(1, 2).pipe(
map((id) => of(`user-${id}`, `detail-${id}`)),
concatAll(),
).subscribe({
next: (value) => console.log(value),
complete: () => console.log('output complete'),
});
// user-1
// detail-1
// user-2
// detail-2
// output completeType đổi theo hai bước:
Observable<number>
── map ──► Observable<Observable<string>>
── concatAll ──► Observable<string>Các inner trong ví dụ phát đồng bộ nên bạn chưa thấy thời gian chờ. Với inner bất đồng bộ, concatAll giữ các inner đến sau trong hàng đợi. Nếu inner đầu không complete, inner sau sẽ không được subscribe.
Gộp hai bước bằng higher-order mapping
Khi mỗi outer value cần tạo inner, dạng ...Map thường dễ đọc hơn vì gộp projection và flatten trong cùng một operator:
| Hai bước | Dạng gộp với cùng chiến lược |
|---|---|
map(project) rồi concatAll() | concatMap(project) |
map(project) rồi mergeAll() | mergeMap(project) |
map(project) rồi switchAll() | switchMap(project) |
map(project) rồi exhaustAll() | exhaustMap(project) |
Bảng mô tả chiến lược flatten, không cam kết projection được gọi cùng thời điểm. map(project) gọi projection cho mọi outer value ngay khi value đến. concatMap chỉ gọi projection khi đến lượt xử lý; mergeMap có giới hạn concurrency cũng chờ slot trống; exhaustMap không gọi projection cho value bị bỏ qua.
import { concatMap, of } from 'rxjs';
of(1, 2).pipe(
concatMap((id) => of(`user-${id}`, `detail-${id}`)),
).subscribe(console.log);
// user-1
// detail-1
// user-2
// detail-2Mình dùng ...Map khi input còn là ID, query hay event cần biến thành công việc. Nếu source đã phát sẵn Observable, dùng ...All để tránh viết một projection chỉ trả lại chính inner.
Value mới đến khi inner cũ chưa xong
Sau khi hiểu flatten, câu hỏi tiếp theo không còn là “làm sao lấy value?” mà là “công việc mới phải làm gì với công việc đang chạy?”. Đây là khác biệt quyết định giữa bốn chiến lược.
Ví dụ request hoàn tất ngược thứ tự
Giả sử request cho ID 1 mất 80 ms, còn ID 2 mất 20 ms. Các số này chỉ là thời gian mô phỏng để làm rõ thứ tự, không phải benchmark.
import { Observable, map, mergeMap, of, timer } from 'rxjs';
type User = { id: number; name: string };
function loadUser$(id: number): Observable<User> {
const delayMs = id === 1 ? 80 : 20;
return timer(delayMs).pipe(
map(() => ({ id, name: `User ${id}` })),
);
}
of(1, 2).pipe(
mergeMap((id) => loadUser$(id)),
).subscribe({
next: (user) => console.log('user:', user.id),
complete: () => console.log('output complete'),
});
// user: 2
// user: 1
// output completemergeMap subscribe cả hai inner mà không chờ request 1 xong. Vì inner 2 phát trước, output cũng nhận 2 trước. Outer of(1, 2) đã complete đồng bộ, nhưng output vẫn đợi các inner đang chạy.
Sự kiện mô phỏng với mergeMap:
t = 0 ms outer phát 1 → subscribe inner 1
outer phát 2 → subscribe inner 2
outer complete
t ≈ 20 ms inner 2 phát user 2 rồi complete → output nhận user 2
t ≈ 80 ms inner 1 phát user 1 rồi complete → output nhận user 1
không còn inner đang chạy → output completeNếu output được dùng để ghi vào một ô “user đang chọn”, kết quả cũ có thể ghi đè kết quả mới. Nhưng nếu đó là danh sách hai user cần tải độc lập, nhận kết quả theo thứ tự hoàn tất có thể hoàn toàn đúng.
Chọn chiến lược theo yêu cầu nghiệp vụ
| Operator | Khi outer phát value mới lúc inner cũ đang chạy | Khi phù hợp |
|---|---|---|
concatMap | Giữ value mới trong hàng đợi, chờ inner trước complete | Lưu tuần tự, cần giữ thứ tự công việc |
mergeMap | Subscribe thêm inner nếu còn slot; mặc định không giới hạn concurrency | Tác vụ độc lập, có thể chạy đồng thời |
switchMap | Unsubscribe inner cũ và subscribe inner mới | Query, route param; chỉ kết quả mới nhất còn hữu ích |
exhaustMap | Bỏ value mới, không xếp hàng, cho tới khi inner hiện tại complete | Bỏ click submit lặp trong lúc đang xử lý |
Với đúng source đồng bộ of(1, 2) và inner timer phía trên, thay operator sẽ cho:
| Operator thay vào ví dụ | ID được phát | Lý do |
|---|---|---|
concatMap | 1, rồi 2 | Inner 2 chỉ bắt đầu sau khi inner 1 complete |
mergeMap | 2, rồi 1 | Cả hai cùng hoạt động, inner 2 nhanh hơn |
switchMap | Chỉ 2 | Outer phát 2 trước khi timer của inner 1 chạy |
exhaustMap | Chỉ 1 | Outer phát 2 khi inner 1 còn đang hoạt động |
Đừng suy rộng kết quả “chỉ ID cuối” hay “chỉ ID đầu” cho mọi source. Nếu inner complete đồng bộ trước outer value tiếp theo, switchMap và exhaustMap vẫn có thể phát kết quả của mọi value.
Mình chọn theo một câu nghiệp vụ trước: đợi, chạy thêm, thay cũ hay bỏ mới. Search thường cần switchMap; thao tác ghi bắt buộc xử lý đủ thường cần concatMap hoặc mergeMap, tùy yêu cầu thứ tự. Không nên mặc định dùng switchMap chỉ vì nó giúp tránh dữ liệu cũ ở UI: unsubscribe một thao tác ghi không đảm bảo backend chưa thực hiện thao tác đó.
Concurrency không tự tạo backpressure
mergeMap(project, 2) giới hạn hai inner hoạt động cùng lúc, nhưng outer vẫn có thể phát nhanh và làm hàng đợi tăng. concatMap cũng có nguy cơ này. Nếu producer nhanh hơn consumer lâu dài, bạn cần quyết định giới hạn, gom batch, bỏ value hoặc kiểm soát tốc độ ở tầng khác.
Lifecycle của pipeline sau khi flatten
Khi dùng bốn operator phía trên, output không chỉ chuyển tiếp dữ liệu; nó còn quản lý thời điểm kết thúc và teardown. Muốn chọn đúng operator, bạn cần biết vòng đời của cả outer lẫn inner.
Complete của outer chưa phải complete của output
Với các chiến lược đang xét, output complete khi outer đã complete và không còn inner cần xử lý:
concatMap: đợi inner đang chạy và các value đã xếp hàng xử lý xong.mergeMap: đợi mọi inner đang hoạt động và hàng đợi, nếu có, xử lý xong.switchMap: đợi inner hiện tại complete; các inner đã bị thay thế không còn được theo dõi.exhaustMap: đợi inner đang hoạt động complete; các outer value bị bỏ không được xử lý lại.
Một inner interval(...) không có ranh giới kết thúc có thể làm output không bao giờ complete dù outer đã complete. Với concatMap, nó còn chặn cả hàng đợi; với exhaustMap, outer value mới sẽ tiếp tục bị bỏ khi inner còn hoạt động.
Ngược lại, inner complete không làm outer complete. Stream click vẫn có thể tiếp tục nhận click mới sau khi một request xong.
Error của inner có thể kết thúc cả pipeline
Mặc định, error từ outer hoặc từ một inner đang được subscribe đi xuống output và đóng pipeline. Với mergeMap, các inner khác đang hoạt động cũng bị unsubscribe. Nếu lỗi chỉ thuộc một request và bạn muốn giữ outer sống, đặt catchError bên trong projection:
import {
EMPTY,
catchError,
concatMap,
of,
throwError,
} from 'rxjs';
of(1, 2, 3).pipe(
concatMap((id) => {
const request$ = id === 2
? throwError(() => new Error('Request 2 thất bại'))
: of({ id });
return request$.pipe(
catchError((error: unknown) => {
const message = error instanceof Error ? error.message : String(error);
console.log('bỏ request lỗi:', message);
return EMPTY;
}),
);
}),
).subscribe({
next: (user) => console.log('user:', user.id),
complete: () => console.log('output complete'),
});
// user: 1
// bỏ request lỗi: Request 2 thất bại
// user: 3
// output completeEMPTY complete inner lỗi mà không phát value, nên concatMap có thể đi tiếp. Trong ứng dụng thật, có thể cần phát một trạng thái lỗi có kiểu rõ ràng thay vì bỏ kết quả; lựa chọn đó thuộc yêu cầu UI và nghiệp vụ.
Nếu đặt catchError(() => EMPTY) sau concatMap, lỗi request 2 sẽ làm pipeline cũ bị teardown, rồi replacement EMPTY complete. Nó không quay lại outer để tiếp tục ID 3.
Unsubscribe không phải complete
Unsubscribe subscription của output sẽ tháo các subscription outer và inner mà flattening operator quản lý. Đây là lý do một subscription ở ranh giới view thường dễ cleanup hơn nhiều subscription lồng riêng lẻ.
Nhưng unsubscribe không gọi callback complete của subscriber. Dùng finalize nếu cần cleanup hoặc quan sát kết thúc trên cả complete, error và unsubscribe. Với switchMap, finalize trong inner có thể chạy vì inner bị thay thế, không phải vì request đã thành công.
Việc resource thật có dừng hay không vẫn phụ thuộc producer. Unsubscribe khỏi timer hủy lịch phát của nó; unsubscribe khỏi Observable bọc một Promise đã chạy không tự hủy Promise. Hủy HTTP cần source có teardown tích hợp cơ chế như AbortController; và kể cả client hủy request, server vẫn có thể đã xử lý dữ liệu.
Vì sao không mặc định dùng subscribe lồng
Đoạn code sau đọc rất tự nhiên nhưng tạo hai subscription độc lập:
import { map, of, timer } from 'rxjs';
const outerSubscription = of(1, 2).subscribe((id) => {
timer(50).pipe(
map(() => ({ id })),
).subscribe((user) => console.log(user));
});
outerSubscription.unsubscribe();
// Các timer inner đã được subscribe vẫn có thể phát sau đó.RxJS không tự gắn subscription tạo trong callback next vào subscription outer. Outer complete cũng không đợi inner; error của inner không tự đi vào error handler của outer. Muốn chọn latest, tuần tự hay concurrency, bạn phải tự viết thêm logic và quản lý từng inner.
Dùng higher-order mapping giúp đưa các quyết định đó vào pipeline, rồi subscribe một lần ở ranh giới tiêu thụ. Subscribe lồng không bị cấm, nhưng mình chỉ dùng khi các tác vụ thật sự có lifecycle độc lập và việc tách đó là chủ đích, không phải cách né việc chọn operator.
Observable lồng và Promise lồng
async callback luôn trả Promise. Vì thế, map(async (...) => ...) tạo Observable<Promise<T>>, không phải Observable<T> và cũng không phải Observable<Observable<T>> theo định nghĩa chặt ở đầu bài.
import { Observable, defer, map, of, switchMap } from 'rxjs';
type User = { id: number; name: string };
function loadUser(id: number): Promise<User> {
return Promise.resolve({ id, name: `User ${id}` });
}
const ids$ = of(1, 2);
const promises$: Observable<Promise<User>> = ids$.pipe(
map(async (id) => loadUser(id)),
);
const users$: Observable<User> = ids$.pipe(
switchMap((id) => defer(() => loadUser(id))),
);
users$.subscribe((user) => console.log(user.id));
// 2Các flattening operator nhận ObservableInput, nên có thể nhận Promise trực tiếp. defer ở đây trì hoãn lời gọi loadUser tới lúc inner được subscribe; nó không bổ sung khả năng hủy cho Promise.
Vì of(1, 2) phát đồng bộ còn Promise resolve qua microtask, switchMap bỏ theo dõi Promise thứ nhất trước khi nó phát kết quả. Kết quả 1 không đi xuống UI, nhưng công việc nền của Promise thứ nhất vẫn có thể tiếp tục. Nếu cần chạy đủ mọi thao tác, hãy chọn chiến lược khác thay vì suy luận rằng “không nhận kết quả” nghĩa là “chưa thực hiện công việc”.
Khi nào cần giữ higher-order stream
Không phải cứ thấy type lồng là phải flatten ngay. Giữ higher-order stream hữu ích khi từng inner còn cần pipeline riêng:
groupBytạo stream con theo key; mỗi group có thể cần reducer hoặc quy tắc đóng riêng.windowCount,windowTimetạo stream con theo cửa sổ; bạn có thể tính toán riêng trong từng window.- Một abstraction có thể trả các inner cho tầng gọi quyết định chiến lược flatten, thay vì cố định chiến lược quá sớm.
Nhưng nếu UI chỉ cần User, hãy flatten trước ranh giới UI. Đẩy Observable<Observable<User>> ra component thường buộc component tự quản subscription và dễ tạo lại vấn đề subscribe lồng.
Không phải inner nào cũng cold
Inner từ groupBy hoặc window thường phát dữ liệu theo source và không mặc định replay value cũ. Xếp chúng vào hàng đợi rồi subscribe muộn có thể bỏ lỡ dữ liệu. Với loại inner này, hãy xét thời điểm subscribe và lifecycle của từng group/window, không áp dụng máy móc giả định “đợi rồi chạy” như với request cold.
Bài tập tự kiểm tra
- Đổi
concatAll()ở ví dụ đồng bộ thànhmergeAll()hoặcswitchAll(). Dự đoán output trước khi chạy: cả ba vẫn phát đủ bốn chuỗi vì từng inner complete đồng bộ trước outer emission tiếp theo. - Trong ví dụ timer, thử lần lượt bốn operator. Ghi lại ID được phát, thời điểm bắt đầu inner
2và thời điểm output complete; đối chiếu bảng chiến lược. - Thêm
finalize(() => console.log('kết thúc inner', id))vàoloadUser$. VớiswitchMap, giải thích vì sao inner1finalize nhưng không phát user1. - Chuyển
catchErrortrong ví dụ lỗi ra sauconcatMap. Dự đoán ID3có còn được xử lý không và giải thích dựa trên teardown, không chỉ dựa trên output.
Trước khi đọc từng operator chi tiết, hãy chọn một pipeline thật và ghi ra ba dòng: outer phát gì, inner phát gì, và value mới phải đợi / chạy thêm / thay cũ / bị bỏ. Nếu chưa viết được dòng thứ ba, chưa nên chốt flattening operator.
Học tiếp
concatMap
Đi sâu vào hàng đợi, thứ tự xử lý và inner không complete.
mergeMap
Quản lý công việc đồng thời và giới hạn concurrency.
switchMap
Chọn kết quả mới nhất và hiểu giới hạn của cancellation.
exhaustMap
Bỏ event mới khi tác vụ hiện tại còn đang hoạt động.
Race condition
Nhận diện kết quả cũ ghi đè state mới.
Subscription và teardown
Theo dõi cleanup khi complete, error hoặc unsubscribe.
Nguồn tham khảo
- RxJS 7.8.2 — Operators guide — higher-order Observable và các chiến lược flatten.
- RxJS 7.8.2 —
mergeInternals— quản lý concurrency, hàng đợi và điều kiện complete củamergeMap/concatMap. - RxJS 7.8.2 —
switchMap— unsubscribe inner cũ và đợi inner hiện tại khi outer complete. - RxJS 7.8.2 —
exhaustMap— bỏ outer value khi inner đang hoạt động, không gọi projection cho value bị bỏ. - RxJS —
defer,catchErrorvàfinalize— khởi tạo theo subscription, thay thế stream lỗi và cleanup.