Observable & Subscription
Đi sâu vào cách Observable sản xuất dữ liệu và Subscription quản lý tài nguyên.
Bạn nhìn thấy cùng một kiểu Observable<T>, nhưng bên dưới nó có thể là một HTTP request chạy riêng cho từng subscriber, một WebSocket dùng chung, hoặc một event listener phải được gỡ khi màn hình biến mất. Nếu chỉ nhìn các giá trị được phát ra, bạn rất dễ tạo request trùng, bỏ lỡ emission hoặc để timer tiếp tục chạy sau khi UI đã bị tháo.
Nhóm bài này giúp bạn đọc toàn bộ vòng đời của một execution: dữ liệu được tạo ở đâu, subscribe() thực sự làm gì, các subscriber có dùng chung producer hay không, và tài nguyên được dọn vào lúc nào.
Phạm vi phiên bản
Nội dung dùng public API của RxJS 7.x và đối chiếu với RxJS 7.8.2. Đây là trang định hướng cho cả nhóm; mỗi bài con sẽ đi sâu vào một quyết định cụ thể.
Mục lục
- Bạn sẽ học gì trong nhóm này
- Bức tranh tổng thể
- Một lần subscribe diễn ra thế nào
- Hai trục cần tách riêng
- Cleanup là một phần của contract
- Ví dụ xuyên suốt với một timer
- Lộ trình học đề xuất
- Checklist khi đọc một Observable
- 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 trả lời được các câu hỏi sau mà không cần đoán từ tên biến có dấu $:
- Source bắt đầu ngay khi được tạo, hay chỉ bắt đầu khi có
subscribe()? - Mỗi subscription tạo producer riêng, hay nhiều subscription quan sát cùng producer?
- Subscriber đến muộn có nhận lại dữ liệu cũ không?
unsubscribe()chỉ ngừng nhận notification, hay còn dừng được công việc bên ngoài?- Ai sở hữu timer, listener, socket hoặc request, và ai chịu trách nhiệm đóng chúng?
- Khi nào nên dùng creation function có sẵn, khi nào mới cần
new Observable()?
Đây là những câu hỏi về semantics và ownership, không chỉ về syntax. Hai pipeline cùng trả Observable<User> vẫn có thể khác hoàn toàn: một pipeline gửi request mới cho mỗi subscriber, pipeline kia dùng lại một request hoặc một giá trị đã cache.
Bức tranh tổng thể
Bốn vai trò xuất hiện xuyên suốt nhóm này:
| Vai trò | Trách nhiệm | Ví dụ |
|---|---|---|
| Producer | Thật sự tạo dữ liệu hoặc side effect | Timer, HTTP request, DOM, WebSocket |
| Observable | Mô tả cách kết nối producer với consumer | interval(1000), fromEvent(...), pipeline qua pipe() |
| Observer | Phản ứng với next, error, complete | Object truyền vào subscribe() |
| Subscription | Đại diện cho một lần đăng ký cụ thể và quản lý việc hủy | Giá trị trả về từ subscribe() |
next(value)
┌──────────┐ ┌────────────┐ ─────────────► ┌──────────┐
│ Producer │ ───────► │ Observable │ error(err) │ Observer │
└────┬─────┘ └──────┬─────┘ ─────────────► └──────────┘
│ │ complete()
│ │
└──── cleanup ◄─────────┴──── Subscription
│
unsubscribe()Sơ đồ này là điểm xuất phát, không phải toàn bộ câu chuyện. Một Observable có thể tạo producer mới, kết nối vào producer đã tồn tại hoặc đặt một cơ chế chia sẻ ở giữa. Vì vậy, muốn hiểu stream, bạn phải theo đường đi từ subscribe() ngược lên nơi producer được tạo.
Observable mô tả, subscribe khởi chạy
Observable thường là một mô tả lazy: khai báo pipeline chưa đồng nghĩa công việc đã chạy. subscribe() thiết lập một execution cụ thể và trả về Subscription đại diện cho execution đó.
import { defer, of } from 'rxjs';
const value$ = defer(() => {
console.log('producer: bắt đầu');
return of(42);
});
console.log('đã khai báo');
const subscription = value$.subscribe({
next: (value) => console.log('next:', value),
complete: () => console.log('complete'),
});
console.log('closed:', subscription.closed);Kết quả:
đã khai báo
producer: bắt đầu
next: 42
complete
closed: trueVí dụ này còn cho thấy Observable không mặc định là async. of(42) phát đồng bộ, nên execution đã complete trước khi subscribe() trả về. Nguồn dữ liệu và scheduler quyết định timing; kiểu Observable<T> tự nó không nói điều đó.
Lazy không phải quy tắc cho mọi producer
Promise bắt đầu khi được tạo, DOM có thể phát event dù không có subscriber, còn WebSocket có thể được mở bởi code khác. Observable chỉ trì hoãn được công việc nếu chính phần tạo công việc cũng nằm trong lúc subscribe, chẳng hạn defer(() => from(fetch(...))).
Producer mới là nơi công việc xảy ra
Khi debug stream, đừng dừng ở operator cuối cùng. Hãy tìm thứ thật sự tạo dữ liệu:
of()vàfrom(array)tự duyệt dữ liệu có sẵn;from(Promise)quan sát một Promise, nhưng không khiến Promise trở thành lazy;fromEvent()gắn handler vào một producer bên ngoài như DOM;interval()tạo lịch phát theo từng subscription;Subjectnhận notification từ code khác rồi phân phối cho subscriber;share()tạo một boundary để nhiều subscriber dùng chung một source subscription.
Producer quyết định ba thứ quan trọng: khi nào bắt đầu, phát cho ai, và dừng bằng cách nào. Observable cùng các operator đóng gói các quyết định đó để consumer có một API thống nhất.
Một lần subscribe diễn ra thế nào
Với một Observable thông thường, có thể đọc subscribe() theo chuỗi sau:
1. Consumer gọi subscribe(observer)
│
▼
2. Observable thiết lập execution và kết nối producer
│
▼
3. Producer gửi next* rồi có thể error hoặc complete
│
▼
4. Consumer hoặc operator có thể unsubscribe sớm
│
▼
5. Subscription đóng và teardown chạyMột execution tuân theo contract:
next* (error | complete)?Nó có thể gửi nhiều next, rồi kết thúc bằng tối đa một error hoặc complete. Một nguồn sống lâu cũng có thể không tự gửi terminal notification nào. unsubscribe() đứng ngoài grammar này vì nó là hành động của consumer, không phải notification từ producer.
Điểm dễ nhầm nhất là unsubscribe() không gọi callback complete. Cả complete, error và unsubscribe đều đóng subscription và kích hoạt teardown đã đăng ký, nhưng ý nghĩa của chúng khác nhau:
| Đường kết thúc | Ai khởi phát? | Observer nhận gì? | Teardown chạy? |
|---|---|---|---|
complete() | Producer | complete | Có |
error(err) | Producer hoặc pipeline | error | Có |
unsubscribe() | Consumer hoặc operator upstream/downstream | Không có terminal notification cho consumer hủy | Có |
Hai trục cần tách riêng
Cold/hot và unicast/multicast thường xuất hiện cùng nhau nên dễ bị dùng như từ đồng nghĩa. Tách chúng thành hai câu hỏi sẽ giúp bạn review code chính xác hơn.
Trục cold và hot
Trục này hỏi: producer được tạo trong lúc subscribe, hay đã tồn tại bên ngoài subscription?
- Cold Observable tạo execution hoặc producer riêng cho từng subscription.
of(),range()vàinterval()là các ví dụ cold điển hình. - Hot Observable kết nối subscriber vào producer đã tồn tại hoặc đang chạy. DOM event, một WebSocket dùng chung và
Subjectthường mang hành vi hot.
COLD HOT
subscribe A ─► producer A ─► A ┌─► A
producer ┤
subscribe B ─► producer B ─► B └─► BCold cho mỗi consumer sự độc lập, nhưng side effect có thể chạy lặp. Hot cho phép chia sẻ một nguồn sống, nhưng subscriber đến muộn có thể bỏ lỡ dữ liệu và lifecycle của producer cần một owner rõ ràng.
Trục unicast và multicast
Trục này hỏi: một execution phân phối notification cho một hay nhiều consumer?
- Unicast là quan hệ một producer với một consumer trong execution đó.
- Multicast cho nhiều consumer quan sát cùng một producer.
Plain cold Observable thường unicast vì mỗi lần subscribe tạo execution riêng. share() có thể đặt một Subject nội bộ vào pipeline để các subscriber đang chồng lấp dùng chung source subscription. Tuy nhiên, multicast không tự động có nghĩa là cache: subscriber đến muộn chỉ nhận dữ liệu cũ nếu bạn chọn thêm replay semantics.
| Câu hỏi | Khái niệm cần dùng |
|---|---|
| Producer bắt đầu ở đâu và lúc nào? | Cold / hot |
| Một producer đang phục vụ bao nhiêu consumer? | Unicast / multicast |
| Subscriber đến muộn có nhận emission cũ không? | Replay / cache policy |
| Khi subscriber cuối rời đi thì source có dừng không? | RefCount và connection lifecycle |
Cách phân tích thực dụng
Đặt log ở nơi producer bắt đầu và nơi teardown chạy. Giá trị đầu ra giống nhau không chứng minh hai subscriber dùng chung execution; số lần producer được tạo mới là bằng chứng.
Cleanup là một phần của contract
Một subscription đóng không có nghĩa mọi công việc bên ngoài tự biến mất. RxJS chỉ có thể cleanup thứ mà source hoặc operator đã đăng ký cách dọn tương ứng.
interval()biết cách hủy timer nội bộ;fromEvent()biết cách gỡ listener mà nó đã thêm;- custom Observable phải trả teardown cho resource mà nó tạo;
from(existingPromise)ngừng chuyển kết quả đến subscriber đã đóng, nhưng không thể tự hủy Promise;- một request chỉ bị abort thật sự khi source tích hợp cơ chế như
AbortController.
import { Observable } from 'rxjs';
const heartbeat$ = new Observable<number>((subscriber) => {
let beat = 0;
const timerId = setInterval(() => {
subscriber.next(beat++);
}, 1_000);
return () => {
clearInterval(timerId);
console.log('teardown: đã dừng timer');
};
});
const subscription = heartbeat$.subscribe({
next: (beat) => console.log('beat:', beat),
complete: () => console.log('complete'),
});
setTimeout(() => subscription.unsubscribe(), 2_500);Sau khoảng 2,5 giây, teardown dừng timer nhưng callback complete không chạy. Nếu bỏ clearInterval, subscriber đã đóng sẽ không nhận thêm notification, nhưng timer vẫn đánh thức event loop. “Không còn thấy log” chưa đủ để kết luận resource đã được giải phóng.
Mặc định, hãy giữ quy tắc ownership sau:
- code tạo resource phải định nghĩa teardown;
- code gọi
subscribe()phải biết subscription sống đến lúc nào; - code cần chạy cleanup ở cấp pipeline trên cả complete, error và unsubscribe nên dùng
finalize(); - framework có lifecycle primitive riêng thì ưu tiên primitive đó thay vì tự quản lý mọi subscription bằng tay.
Ví dụ xuyên suốt với một timer
Một timer nhỏ đủ để nhìn thấy producer ownership, số execution và teardown mà không bị nhiễu bởi network hay framework.
Hai subscription độc lập
import { defer, finalize, interval, take } from 'rxjs';
const ticks$ = defer(() => {
console.log('producer: tạo timer');
return interval(1_000).pipe(
take(3),
finalize(() => console.log('producer: dừng timer')),
);
});
const subscriptionA = ticks$.subscribe({
next: (value) => console.log('A:', value),
});
const subscriptionB = ticks$.subscribe({
next: (value) => console.log('B:', value),
});Dòng producer: tạo timer xuất hiện hai lần. A và B có hai timer, hai state đếm và hai teardown độc lập. Cùng nhận 0, 1, 2 không có nghĩa họ dùng chung producer.
Nếu A gọi unsubscribe() sớm, timer của A dừng nhưng B tiếp tục. Đây là hành vi đúng khi mỗi consumer cần isolation; nó chỉ là vấn đề nếu timer đại diện cho một side effect đắt đỏ mà bạn định dùng chung.
Chia sẻ một execution
Đặt share() sau source để những subscription đang chồng lấp dùng chung một source subscription:
import { defer, finalize, interval, share, take } from 'rxjs';
const sharedTicks$ = defer(() => {
console.log('producer: tạo timer');
return interval(1_000).pipe(
take(3),
finalize(() => console.log('producer: dừng timer')),
);
}).pipe(
share(),
);
const subscriptionA = sharedTicks$.subscribe({
next: (value) => console.log('A:', value),
});
const subscriptionB = sharedTicks$.subscribe({
next: (value) => console.log('B:', value),
});Lần này producer chỉ được tạo một lần, còn A và B vẫn nhận hai Subscription riêng. Nếu A rời đi, B vẫn giữ connection. Khi source complete hoặc subscriber cuối cùng rời đi, cấu hình mặc định của share() sẽ đóng connection và reset trạng thái.
Có hai giới hạn cần nhớ:
share()chỉ chia sẻ execution trong khoảng các subscription chồng lấp. Nếu source đã complete trước khi B đến, B có thể tạo execution mới.share()mặc định không replay. B đến sau emission0sẽ không nhận lại0; muốn replay phải chọn chiến lược khác và xác định cache lifecycle rõ ràng.
Ví dụ này nối toàn bộ mental model của nhóm: creation quyết định source, cold/hot mô tả ownership của producer, unicast/multicast mô tả cách phân phối, còn Subscription quyết định khi connection được đóng.
Lộ trình học đề xuất
Nên đọc theo thứ tự dưới đây vì mỗi bài trả lời một câu hỏi còn bỏ ngỏ từ bài trước:
| Bài | Câu hỏi chính | Kết quả cần đạt |
|---|---|---|
| Creation functions | Đưa value, collection, Promise, event và clock vào RxJS thế nào? | Chọn source theo timing, completion và cancellation thay vì chỉ theo type. |
| Cold và hot Observable | Producer được tạo riêng hay đã tồn tại và dùng chung? | Dự đoán side effect có chạy lại khi subscribe hay không. |
| Unicast và multicast | Một execution phát cho một hay nhiều consumer? | Biết khi nào giữ isolation, khi nào dùng share() hoặc Subject. |
| Subscription và teardown | Ai hủy execution và resource được dọn ra sao? | Thiết kế ownership, cancellation và cleanup rõ ràng. |
| Tự tạo Observable | Khi adapter có sẵn không đủ thì bọc producer thế nào? | Viết custom source giữ đúng contract và teardown đối xứng. |
1. Creation functions
Chọn cách tạo Observable theo loại source và lifecycle.
2. Cold và hot Observable
Xác định producer được tạo lúc nào và thuộc về ai.
3. Unicast và multicast
Hiểu một execution được phân phối cho các consumer ra sao.
4. Subscription và teardown
Dừng execution và giải phóng resource đúng chỗ.
5. Tự tạo Observable
Đóng gói callback hoặc resource tùy biến an toàn.
Checklist khi đọc một Observable
Trước khi đưa một stream vào production, hãy lần theo các câu hỏi này:
- Producer là gì? Timer, request, listener, socket, array hay một Subject?
- Producer được tạo lúc nào? Khi khai báo, khi Promise được tạo, hay trong mỗi lần subscribe?
- Subscribe lần hai có lặp side effect không? Nếu có, đó là chủ đích hay lỗi?
- Subscriber có dùng chung execution không? Nếu dùng chung, late subscriber nhận future value hay cả giá trị đã replay?
- Source kết thúc bằng cách nào? Complete, error, điều kiện operator hay phải unsubscribe?
- Teardown dọn resource gì? Chỉ ngừng notification, hay thật sự dừng timer/request/socket?
- Ai sở hữu lifecycle? Component, service, command, application hay chính source?
- Có cần
new Observable()không? Nếu creation function có sẵn đã mô tả đúng source, ưu tiên API đó.
Nếu chưa trả lời được câu 2, 3 hoặc 6, đừng vội thêm shareReplay() hay lưu Subscription vào một mảng chung. Những cách đó có thể che triệu chứng trong khi ownership vẫn chưa rõ.
Bài tập tự kiểm tra
Dùng ví dụ ticks$ phía trên và thử theo thứ tự:
- Chạy hai subscription vào bản cold, đếm số lần producer bắt đầu và teardown.
- Hủy A sau emission đầu tiên, xác nhận B vẫn chạy trên timer riêng.
- Thêm
share(), lặp lại và giải thích vì sao producer chỉ bắt đầu một lần. - Cho B subscribe sau khoảng 1,5 giây, ghi lại emission đầu tiên mà B nhận được.
- Hủy cả A và B trước
take(3), rồi subscribe C để kiểm tra source có bắt đầu lại từ0không. - Thay
interval()bằngfrom(Promise.resolve(42)), sau đó giải thích vì sao unsubscribe không thể hủy Promise nền.
Đừng chỉ ghi output. Với mỗi lần thử, hãy vẽ ba thứ: producer nào tồn tại, subscription nào đang active, và teardown nào sẽ chạy. Nếu ba thứ đó khớp với log, bạn đã có mental model đủ chắc để bước sang HTTP, WebSocket và UI lifecycle.
Học tiếp
Nếu bốn vai trò Producer, Observable, Observer và Subscription vẫn còn lẫn nhau, quay lại Observer pattern trước. Nếu đã rõ vai trò nhưng chưa quen callback next, error, complete, đọc Subscribe và Observer. Sau nhóm này, phần Operators sẽ giúp bạn biến đổi và kết hợp execution mà vẫn giữ lifecycle có thể suy luận được.
Observer pattern
Củng cố bốn vai trò cốt lõi trước khi đi sâu vào execution.
Subscribe và Observer
Thực hành cách nhận notification và đặt side effect ở boundary.
Vòng đời của stream
Phân biệt next, error, complete, unsubscribe và finalization.
Operators
Tiếp tục với cách biến đổi, lọc và kết hợp stream.