Subscribe và Observer
Nhận giá trị, lỗi và tín hiệu hoàn tất từ stream.
Bạn đã có một Observable nhưng chạy chương trình lại không thấy gì. Hoặc dữ liệu đã hiện đúng, nhưng khi rời màn hình thì timer và event listener vẫn còn sống. Điểm nối giữa hai vấn đề đó là subscribe(): nó bắt đầu một execution, chuyển notification cho Observer và trả về Subscription để bạn quản lý vòng đời của execution ấy.
Phạm vi phiên bản
Bài này dùng public API ổn định của RxJS 7.x, được đối chiếu với tài liệu và mã nguồn RxJS 7.8.2. Các ví dụ dùng TypeScript; bài không giả định hành vi của RxJS 8.
Mục lục
- Mental model: subscribe mở một phiên nhận dữ liệu
- Hai cách subscribe nên dùng
- Đọc ba loại notification
- Lifecycle và teardown
- Ví dụ thực tế: banner trạng thái mạng
- Đặt side effect ở đâu
- Lựa chọn và pitfalls
- Bài tập tự kiểm tra
- Học tiếp
- Nguồn tham khảo
Mental model: subscribe mở một phiên nhận dữ liệu
Hãy xem Observable như công thức pha cà phê, còn mỗi lần subscribe() là một lần bắt đầu pha cho một người uống. Công thức mô tả việc cần làm nhưng tự nó chưa tạo ra ly cà phê. Tương tự, với Observable thông thường, subscribe() mới thiết lập execution để producer gửi dữ liệu đến consumer.
Ẩn dụ này dừng ở chỗ “mỗi lần thực hiện”. Một hot Observable có thể dùng producer đã chạy sẵn và được chia sẻ; vì vậy không nên suy ra rằng mọi lần subscribe đều tạo nguồn mới. Phần Cold và hot Observable sẽ phân tích ranh giới đó.
Observable ── subscribe(observer) ──► execution đang active
│
┌─────────────┼──────────────┐
│ │ │
next(value) error/complete unsubscribe()
│ │ │
└── Observer │ │
▼ ▼
closed ──► teardown
subscribe(...) ─────────────────────────────► trả về SubscriptionBốn ý cần giữ trong đầu:
Observablemô tả nguồn và pipeline.subscribe()bắt đầu một execution cụ thể và đăng ký consumer.Observerlà tập callback phản ứng vớinext,errorvàcomplete.Subscriptionđại diện cho execution đó;unsubscribe()dùng để dừng nó từ phía consumer.
Nếu bạn cần mô hình đầy đủ hơn về producer, Observable, Observer và Subscription, đọc Observer pattern. Trang này tập trung vào cách viết và đặt subscribe() trong code ứng dụng.
Hai cách subscribe nên dùng
Observer object cho đầy đủ ngữ cảnh
Mặc định, mình chọn Observer object. Tên callback cho biết rõ nhánh nào đang được xử lý, và bạn có thể thêm error hoặc complete mà không phải nhớ thứ tự tham số.
import { from, type Observer } from 'rxjs';
type OrderStatus = 'queued' | 'processing' | 'shipped';
const orderStatuses: OrderStatus[] = [
'queued',
'processing',
'shipped',
];
const orderStatus$ = from(orderStatuses);
const observer: Observer<OrderStatus> = {
next: (status) => console.log('next:', status),
error: (error: unknown) => console.error('error:', error),
complete: () => console.log('complete'),
};
console.log('trước subscribe');
const subscription = orderStatus$.subscribe(observer);
console.log('sau subscribe, closed =', subscription.closed);Kết quả:
trước subscribe
next: queued
next: processing
next: shipped
complete
sau subscribe, closed = trueObserver<T> đầy đủ có cả ba callback. Trong code thường ngày, RxJS cũng nhận partial observer, nên bạn chỉ cần khai báo các callback thật sự dùng:
import { of } from 'rxjs';
of(10, 20, 30).subscribe({
next: (value) => console.log(value),
});Bỏ next hoặc complete chỉ có nghĩa là consumer không phản ứng với loại notification đó. Riêng error cần được cân nhắc kỹ: nếu lỗi đi tới cuối pipeline mà không có error handler, RxJS báo nó là unhandled error trên một call stack khác.
Đừng để error biến mất khỏi thiết kế
Partial observer là hợp lệ, nhưng “không viết error handler” không phải một chiến lược xử lý lỗi. Ở boundary như UI, command hoặc integration, hãy xử lý lỗi trong pipeline hoặc cung cấp callback error có chủ đích.
Callback next cho trường hợp thật sự đơn giản
Nếu chỉ cần quan sát giá trị trong một thử nghiệm nhỏ, bạn có thể truyền trực tiếp callback next:
import { of } from 'rxjs';
of('Rx', 'JS').subscribe((part) => console.log(part));Kết quả:
Rx
JSCách này vẫn hợp lệ trong RxJS 7.8.2. Tuy nhiên, overload nhận ba callback theo vị trí là API đã bị deprecated trong RxJS 7.x:
import { of } from 'rxjs';
const source$ = of(1, 2, 3);
// Không nên dùng: callback phụ thuộc vào vị trí tham số.
source$.subscribe(
(value) => console.log(value),
(error) => console.error(error),
() => console.log('complete'),
);
// Nên dùng: ý nghĩa của từng callback được gọi tên.
source$.subscribe({
next: (value) => console.log(value),
error: (error: unknown) => console.error(error),
complete: () => console.log('complete'),
});Quy tắc thực dụng là: callback đơn chỉ dành cho next; khi cần từ hai loại notification trở lên, dùng Observer object.
Đọc ba loại notification
Một Observable execution có contract ngắn gọn:
next* (error | complete)?Nó có thể gửi nhiều next, sau đó có nhiều nhất một terminal notification là error hoặc complete. Cũng có những source sống lâu không tự gửi terminal notification, chẳng hạn DOM event.
next có thể xuất hiện nhiều lần
next(value) chuyển một giá trị sang Observer nhưng không đóng execution. Callback next là nơi consumer phản ứng với dữ liệu: render UI, chuyển dữ liệu sang adapter imperative, hoặc ghi log tại boundary.
Số lần next không cho biết source đã xong hay chưa. Một request có thể phát một giá trị rồi complete; một WebSocket có thể phát hàng nghìn giá trị mà chưa complete.
error kết thúc execution do lỗi
error(error) vừa thông báo lỗi vừa đóng execution. Sau nó sẽ không còn next và cũng không có complete trong cùng execution.
import { concat, of, throwError } from 'rxjs';
const result$ = concat(
of('dữ liệu từ cache'),
throwError(() => new Error('HTTP 503')),
of('giá trị này không bao giờ được phát'),
);
result$.subscribe({
next: (value) => console.log('next:', value),
error: (error: unknown) => {
const message = error instanceof Error ? error.message : String(error);
console.log('error:', message);
},
complete: () => console.log('complete'),
});Kết quả:
next: dữ liệu từ cache
error: HTTP 503Không có dòng thứ ba và không có complete, vì error đã kết thúc execution. Nếu muốn fallback hoặc retry, hãy biểu diễn quyết định đó bằng operator thay vì cố tiếp tục trong error handler; xem Error channel.
complete kết thúc bình thường
complete() báo producer đã gửi xong và sẽ không còn giá trị. Nó không mang payload. Nếu cần một kết quả cuối, producer phải gửi kết quả bằng next(result) trước rồi mới complete.
complete cũng không đồng nghĩa với “nghiệp vụ thành công”. Một stream có thể phát { kind: 'not-found' } rồi complete bình thường, vì “không tìm thấy” được mô hình hóa thành dữ liệu chứ không phải lỗi kỹ thuật.
Lifecycle và teardown
Observable có thể chạy đồng bộ
Ví dụ of() ở trên phát toàn bộ giá trị và complete ngay bên trong lời gọi subscribe(). Vì thế dòng sau subscribe chỉ chạy sau callback complete, và subscription.closed đã là true.
Đừng gắn “Observable” với “async”. Nguồn dữ liệu và scheduler quyết định thời điểm notification xuất hiện. of() thường đồng bộ, interval() phát theo thời gian, còn fromEvent() phụ thuộc vào sự kiện bên ngoài; cả ba vẫn dùng cùng contract.
Điều này dẫn tới một pitfall nhỏ nhưng khó chịu: đừng cố dùng biến subscription ngay bên trong callback next của một source đồng bộ để tự hủy lần đầu tiên. Callback có thể chạy trước khi phép gán biến hoàn tất. Nếu ý định là “lấy một giá trị”, dùng operator như take(1) để mô tả trực tiếp điều kiện kết thúc.
Unsubscribe không phải complete
subscribe() trả về một Subscription. Gọi unsubscribe() đóng execution từ phía consumer và kích hoạt teardown, nhưng không gọi callback complete.
import { finalize, interval } from 'rxjs';
const subscription = interval(1_000)
.pipe(
finalize(() => console.log('finalize: dọn pipeline')),
)
.subscribe({
next: (value) => console.log('next:', value),
complete: () => console.log('complete'),
});
setTimeout(() => {
subscription.unsubscribe();
console.log('closed:', subscription.closed);
}, 2_500);Kết quả thường là:
next: 0
next: 1
finalize: dọn pipeline
closed: trueThời điểm có thể xê dịch một chút theo event loop, nhưng dòng complete sẽ không xuất hiện. finalize() vẫn chạy vì finalization xảy ra khi stream complete, error hoặc bị unsubscribe.
Phân biệt trách nhiệm
Producer đặt cleanup tài nguyên ở teardown. Consumer giữ Subscription hoặc gắn subscription vào lifecycle phù hợp. Nếu side effect phải chạy trên mọi đường kết thúc của pipeline, dùng finalize() thay vì chỉ dựa vào callback complete.
Bài Subscription và teardown đi sâu hơn vào ownership, gom nhiều subscription và dọn resource.
Ví dụ thực tế: banner trạng thái mạng
Giả sử một ứng dụng web có phần tử <div id="network-status"></div>. Banner cần hiển thị trạng thái ban đầu, cập nhật khi browser chuyển online/offline và dừng lắng nghe khi màn hình bị tháo.
import {
distinctUntilChanged,
fromEvent,
map,
merge,
startWith,
} from 'rxjs';
type NetworkStatus = 'online' | 'offline';
const banner = document.querySelector<HTMLElement>('#network-status');
if (!banner) {
throw new Error('Không tìm thấy #network-status');
}
const online$ = fromEvent(window, 'online').pipe(
map((): NetworkStatus => 'online'),
);
const offline$ = fromEvent(window, 'offline').pipe(
map((): NetworkStatus => 'offline'),
);
const initialStatus: NetworkStatus = navigator.onLine ? 'online' : 'offline';
const networkStatus$ = merge(online$, offline$).pipe(
startWith(initialStatus),
distinctUntilChanged(),
);
const subscription = networkStatus$.subscribe({
next: (status) => {
banner.textContent = status === 'online' ? 'Đã kết nối' : 'Mất kết nối';
banner.dataset.status = status;
},
error: (error: unknown) => {
console.error('Không thể theo dõi trạng thái mạng:', error);
},
});
export function destroyNetworkBanner(): void {
subscription.unsubscribe();
}Dòng chảy của ví dụ:
networkStatus$mới chỉ mô tả hai event source và giá trị ban đầu.subscribe()khiến haifromEvent()đăng ký listener lênwindow.- Mỗi event trở thành một
next, và Observer cập nhật DOM. destroyNetworkBanner()đóng subscription; teardown củafromEvent()gỡ cả hai listener.- Source DOM event không tự complete, nên lifecycle của màn hình phải quyết định lúc dừng.
Trong framework, hãy ưu tiên primitive lifecycle sẵn có thay vì tự tạo hàm destroy... cho mọi component. Dù dùng cơ chế nào, câu hỏi vẫn giống nhau: ai bắt đầu subscription và ai chịu trách nhiệm kết thúc nó?
Đặt side effect ở đâu
pipe() dùng để mô tả cách dữ liệu được biến đổi; subscribe() là boundary nơi ứng dụng thực sự tiêu thụ kết quả. Vì vậy, một hàm dùng lại nên trả Observable thay vì âm thầm subscribe bên trong. Làm vậy giúp caller giữ quyền xử lý lỗi và hủy execution.
import { map, type Observable } from 'rxjs';
type ApiUser = {
id: string;
displayName: string;
};
export function selectDisplayName(
user$: Observable<ApiUser>,
): Observable<string> {
return user$.pipe(
map((user) => user.displayName.trim()),
);
}Boundary UI mới subscribe:
import { type Observable } from 'rxjs';
declare const displayName$: Observable<string>;
declare const heading: HTMLHeadingElement;
const subscription = displayName$.subscribe({
next: (name) => {
heading.textContent = name;
},
error: (error: unknown) => {
console.error('Không tải được tên hiển thị:', error);
},
});
export function destroyProfile(): void {
subscription.unsubscribe();
}Ngoại lệ là những hàm được đặt tên rõ như startTelemetry() hoặc mountWidget(): bản thân nhiệm vụ của chúng là khởi động side effect. Khi đó, hãy trả Subscription hoặc một hàm cleanup để caller vẫn quản lý được lifecycle.
Lựa chọn và pitfalls
| Tình huống | Lựa chọn mặc định | Vì sao |
|---|---|---|
| Cần xử lý giá trị, lỗi hoặc hoàn tất | subscribe({ next, error, complete }) | Tên callback rõ và không phụ thuộc vị trí tham số. |
| Chỉ in giá trị trong demo ngắn | subscribe(value => ...) | Gọn, vẫn hợp lệ trong RxJS 7.x. |
| Cần biến đổi dữ liệu | Operator trong pipe() | Giữ pipeline lazy và có thể compose; chưa cần side effect. |
| Source sống theo UI, socket hoặc timer | Lưu Subscription hoặc dùng primitive lifecycle | Có nơi rõ ràng để hủy và chạy teardown. |
| Cần cleanup trên complete, error và unsubscribe | finalize() | Bao phủ cả ba đường kết thúc. |
| Hàm thư viện muốn cung cấp dữ liệu cho caller | Trả Observable<T> | Caller sở hữu subscribe, error handling và cancellation. |
Các lỗi nên kiểm tra trước khi review xong:
- Quên rằng subscribe kích hoạt execution. Với cold Observable, subscribe hai lần có thể chạy producer hai lần, gồm cả hai HTTP request.
- Subscribe lồng nhau. Mỗi subscription con có error handling và teardown riêng, nên flow nhanh chóng khó hủy. Ưu tiên flattening operator; xem Nested subscribe.
- Cleanup chỉ trong complete handler. Unsubscribe không gọi
complete; đặt cleanup resource trong teardown hoặcfinalize(). - Gọi unsubscribe ngay sau subscribe như một nghi thức. Với công việc async, bạn có thể hủy trước emission đầu tiên. Chỉ hủy khi lifecycle hoặc điều kiện nghiệp vụ yêu cầu.
- Cho rằng error handler hồi sinh source. Error handler chỉ quan sát terminal notification. Recovery cần operator như
catchError()hoặc một subscription mới. - Để subscribe sâu trong helper. Caller mất quyền quyết định cancellation và rất khó biết side effect bắt đầu ở đâu.
Bài tập tự kiểm tra
Lấy ví dụ banner mạng và thử ba thay đổi:
- thêm
console.logtrước và sausubscribe()để xác định giá trị từstartWith()chạy đồng bộ ở đâu; - gọi
destroyNetworkBanner(), sau đó bật/tắt mạng và xác nhận banner không đổi nữa; - thêm
finalize(() => console.log('đã dọn'))vào pipeline và kiểm tra log xuất hiện khi unsubscribe, dù Observer không có callbackcomplete.
Nếu bạn giải thích được vì sao cả ba kết quả xảy ra, bạn đã nắm đúng ranh giới giữa Observer, Subscription và teardown. Bước tiếp theo là đưa các phép biến đổi ra khỏi subscribe() để pipeline dễ đọc hơn.
Học tiếp
Pipe và operator
Đưa phép biến đổi vào pipeline thay vì nhồi logic vào callback next.
Vòng đời của stream
Hiểu sâu next, error, complete và ba đường đóng subscription.
Subscription và teardown
Quản lý cancellation và cleanup tài nguyên trong code thực tế.
Error channel
Thiết kế xử lý lỗi và recovery cho pipeline.