Học RxJS
Observable & Subscription

Cold và hot Observable

Phân biệt producer riêng và producer được chia sẻ.

Bạn subscribe hai lần vào một stream và vô tình gửi hai HTTP request. Ở chỗ khác, một subscriber đăng ký muộn lại bỏ lỡ dữ liệu vừa phát. Hai hiện tượng tưởng không liên quan này cùng dẫn về một câu hỏi: producer được tạo riêng cho từng subscription, hay đã tồn tại và được nhiều subscription dùng chung? Đó là ranh giới thực dụng giữa cold và hot Observable.

Phạm vi phiên bản

Bài này dùng public API của RxJS 7.x, được đối chiếu với tài liệu và mã nguồn RxJS 7.8.2. Cold và hot mô tả quan hệ giữa producer với thời điểm subscribe; chúng không nói stream đồng bộ hay bất đồng bộ.

Mục lục

Mental model: producer thuộc về ai

Producer là thứ thật sự sinh dữ liệu: vòng lặp qua một array, timer, HTTP request, WebSocket hoặc DOM. Observable là lớp nối consumer với producer đó. Muốn phân loại một stream, đừng nhìn tên biến có dấu $; hãy nhìn nơi producer được tạo.

  • Cold Observable tạo producer trong quá trình subscribe(). Mỗi subscription sở hữu một producer riêng.
  • Hot Observable dùng producer đã được tạo bên ngoài quá trình subscribe(). Nhiều subscription thường quan sát cùng producer.
COLD: producer riêng                    HOT: producer dùng chung

subscribe A ──► producer A ──► A                  ┌──► A
                                         producer ┤
subscribe B ──► producer B ──► B                  └──► B

A và B có execution độc lập.              A và B thấy cùng một execution.

Hình dung cold stream như gọi món theo bàn: mỗi bàn gọi thì bếp làm một phần riêng. Hot stream giống loa thông báo ở ga: loa vẫn thuộc về nhà ga, người đến sau chỉ nghe thông báo từ lúc họ có mặt. Phép so sánh dừng ở ownership và thời điểm tham gia; nó không mô tả error, completion hay cách cleanup của RxJS.

Ba câu hỏi sau thường đủ để xác định hành vi:

  1. Producer được tạo lúc khai báo stream hay bên trong lúc subscribe?
  2. Subscribe lần thứ hai có chạy lại side effect không?
  3. Subscriber đến muộn có nhận được emission cũ không, hay chỉ emission tương lai?

Câu thứ ba không tự quyết định cold/hot, nhưng nó phát hiện một yêu cầu liên quan: stream có cần replay hay không.

Cold Observable: mỗi subscription có producer riêng

Cold Observable tạo producer mới cho mỗi subscription. Vì mỗi consumer có execution độc lập, consumer A có thể bắt đầu, hoàn tất hoặc bị hủy mà không làm producer của B dừng theo.

Các source như of(), range() và interval() mặc định là cold. Với interval(1_000), hai subscriber vào ở hai thời điểm khác nhau sẽ có hai timer và mỗi người đều bắt đầu từ 0.

Hai subscription, hai execution

Ví dụ đồng bộ dưới đây làm ownership lộ rõ mà không phụ thuộc timing:

import { Observable } from 'rxjs';

let nextProducerId = 0;

const coldNumbers$ = new Observable<number>((subscriber) => {
  const producerId = ++nextProducerId;
  console.log(`producer ${producerId}: bắt đầu`);

  subscriber.next(10);
  subscriber.next(20);
  subscriber.complete();

  return () => console.log(`producer ${producerId}: teardown`);
});

coldNumbers$.subscribe({
  next: (value) => console.log('A:', value),
});

coldNumbers$.subscribe({
  next: (value) => console.log('B:', value),
});

Kết quả:

producer 1: bắt đầu
A: 10
A: 20
producer 1: teardown
producer 2: bắt đầu
B: 10
B: 20
producer 2: teardown

Hai consumer nhận cùng dãy số, nhưng không dùng chung execution. Dòng producer 2: bắt đầu mới là bằng chứng quan trọng. Nhìn giá trị giống nhau rồi kết luận “stream được share” là một lỗi suy luận phổ biến.

Tính độc lập này thường là điều bạn muốn khi:

  • mỗi consumer cần timer hoặc state riêng;
  • mỗi lần gọi phải đọc dữ liệu mới;
  • cancellation của consumer này không được ảnh hưởng consumer khác;
  • side effect cần chạy lại theo từng operation.

Đổi lại, subscribe ngoài ý muốn có thể nhân đôi công việc đắt đỏ. Một cold HTTP Observable được subscribe ở template và thêm một lần trong code có thể tạo hai request riêng.

Promise là trường hợp dễ hiểu sai

Promise chạy ngay khi được tạo, không đợi from() hay subscribe(). Vì vậy hai đoạn code trông gần giống nhau nhưng ownership khác hẳn.

import { defer, from } from 'rxjs';

type Profile = {
  id: string;
  displayName: string;
};

function loadProfile(): Promise<Profile> {
  console.log('HTTP: bắt đầu');

  return fetch('/api/profile').then(async (response) => {
    if (!response.ok) {
      throw new Error(`HTTP ${response.status}`);
    }

    return (await response.json()) as Profile;
  });
}

// Promise và HTTP request được tạo ngay tại đây.
const requestPromise = loadProfile();
const reusedRequest$ = from(requestPromise);

// Factory chỉ chạy khi có subscription; mỗi subscription gọi lại loadProfile().
const freshRequest$ = defer(() => from(loadProfile()));

Với reusedRequest$, request đã bắt đầu trước khi có subscriber và mọi subscriber quan sát kết quả của cùng một Promise. Với freshRequest$, mỗi subscription gọi loadProfile() một lần, nên đây là hành vi cold.

defer không nhân bản producer bên ngoài

defer() chỉ trì hoãn việc chạy factory. Nếu factory vẫn trả về cùng một Promise, socket hoặc Subject đã tạo sẵn, các subscription vẫn dùng producer đó. Muốn cold thật sự, factory phải tạo resource mới cho từng subscription.

Hot Observable: producer tồn tại ngoài subscribe

Hot Observable kết nối consumer tới producer đã tồn tại ngoài lần subscribe. Producer có thể phát trước khi có consumer, trong lúc có nhiều consumer, hoặc sau khi một consumer đã rời đi. Vì thế subscription không nhất thiết sở hữu toàn bộ vòng đời producer.

Subject là ví dụ rõ nhất trong RxJS: nó vừa nhận notification như một Observer, vừa multicast notification đó tới những subscriber đang đăng ký.

import { Subject } from 'rxjs';

const messages$ = new Subject<string>();

messages$.subscribe({
  next: (message) => console.log('A:', message),
});

messages$.next('một');

messages$.subscribe({
  next: (message) => console.log('B:', message),
});

messages$.next('hai');
messages$.complete();

Kết quả:

A: một
A: hai
B: hai

messages$ không tạo một producer mới khi B subscribe. B chỉ được thêm vào danh sách observer của Subject hiện có, nên A và B cùng nhận emission hai.

Ngoài Subject, producer bên ngoài như một WebSocket dùng chung, event emitter của SDK hoặc bus sự kiện toàn ứng dụng cũng thường dẫn tới hot stream. Điều cần kiểm tra là ai mở producer, ai đóng nó, và subscription riêng lẻ có quyền dừng resource chung hay không.

Subscriber đến muộn sẽ bỏ lỡ gì

B không nhận một vì Subject thường không lưu lịch sử. Hot không đồng nghĩa với cache hoặc replay. Nếu subscriber mới cần trạng thái hiện tại, bạn phải chọn semantics phù hợp, chẳng hạn BehaviorSubject, ReplaySubject hoặc shareReplay(); mỗi lựa chọn có quy tắc buffer và lifecycle riêng.

Đây cũng là lý do hot stream hợp với event hơn state theo mặc định. Click đã xảy ra thường không cần phát lại cho component vừa mount. Ngược lại, trạng thái đăng nhập hiện tại thường phải có ngay cho consumer mới; chỉ dùng Subject trần sẽ khiến consumer chờ đến thay đổi kế tiếp.

Hot không có nghĩa là luôn phát dữ liệu

Một WebSocket đang mở nhưng server chưa gửi message vẫn là producer dùng chung. Ngược lại, interval() phát liên tục nhưng mặc định vẫn cold vì mỗi subscription tạo timer riêng. “Nóng” không phải tốc độ hay tần suất phát; nó nói về producer được tạo ở đâu và được sở hữu thế nào.

Tương tự, cold/hot không quyết định sync/async:

  • of(1, 2, 3) thường cold và đồng bộ;
  • interval(1_000) cold và bất đồng bộ;
  • một Subject có thể hot và phát đồng bộ khi code gọi next();
  • một WebSocket hot và phát bất đồng bộ.

Cold, hot, unicast và multicast khác nhau thế nào

Hai cặp thuật ngữ trả lời hai câu hỏi gần nhau nhưng không giống nhau:

Khái niệmCâu hỏiHành vi
ColdProducer được tạo khi nào?Tạo một producer mới trong mỗi lần subscribe.
HotProducer được tạo khi nào?Producer tồn tại ngoài lần subscribe và thường được dùng chung.
UnicastMột producer gửi cho bao nhiêu consumer?Một producer phục vụ một consumer.
MulticastMột producer gửi cho bao nhiêu consumer?Một producer phục vụ nhiều consumer.

Theo glossary của RxJS, cold Observable luôn unicast vì mỗi subscription có producer riêng. Hot Observable gần như luôn multicast trong cách dùng thông thường. Tuy vậy, khi review code, vẫn nên tách hai câu hỏi “producer sinh ở đâu?” và “notification được phân phối cho ai?” để tránh biến các nhãn thành phỏng đoán.

Cold + unicast mặc định

producer A ──► subscriber A
producer B ──► subscriber B

Cold source qua share() trong lúc có subscriber

                  ┌──► subscriber A
một producer ─────┤
                  └──► subscriber B

share() không sửa source ban đầu. Nó trả về một Observable mới đặt cơ chế multicast ở giữa source và các downstream subscriber. Bài Unicast và multicast sẽ đi sâu vào cách notification được phân phối; trang này tập trung vào ownership và vòng đời producer.

Biến cold source thành stream được chia sẻ với share

Giả sử hai phần UI cần cùng một chuỗi tick. Nếu cả hai subscribe thẳng vào interval(), bạn có hai timer. Đặt share() sau source khiến những subscription đang chồng lấp dùng chung một source subscription.

import { defer, interval, share, take } from 'rxjs';

const sharedTicks$ = defer(() => {
  console.log('producer: tạo timer');
  return interval(1_000).pipe(take(3));
}).pipe(
  share(),
);

sharedTicks$.subscribe({
  next: (value) => console.log('A:', value),
  complete: () => console.log('A: complete'),
});

setTimeout(() => {
  sharedTicks$.subscribe({
    next: (value) => console.log('B:', value),
    complete: () => console.log('B: complete'),
  });
}, 1_500);

Kết quả xấp xỉ:

producer: tạo timer
A: 0
A: 1
B: 1
A: 2
B: 2
A: complete
B: complete

Chỉ có một dòng producer: tạo timer, vì B đến khi execution của A còn active. B không nhận 0: share() mặc định dùng Subject, nên nó chỉ chuyển tiếp notification từ thời điểm B đăng ký, không replay lịch sử.

Mặc định thực dụng

Dùng share() khi nhiều consumer đồng thời cần cùng side effect nhưng subscriber đến muộn chỉ cần dữ liệu tương lai. Nếu consumer mới phải nhận giá trị gần nhất, hãy xem shareReplay() thay vì tự giữ một biến cache bên cạnh stream.

Vòng đời reference count

share() trong RxJS 7.8.2 theo dõi số downstream subscriber đang active:

số subscriber     0 ──► 1 ──► 2 ──► 1 ──► 0
kết nối source    tắt   mở     giữ    giữ    đóng

Theo cấu hình mặc định:

  1. subscriber đầu tiên làm share() subscribe vào source;
  2. subscriber tiếp theo dùng chung kết nối đang có;
  3. một subscriber rời đi chỉ đóng subscription của chính nó;
  4. khi số subscriber về 0 trước lúc source kết thúc, share() unsubscribe source và reset trạng thái;
  5. khi source complete hoặc error, trạng thái cũng reset; subscription đến sau sẽ tạo execution mới.

Cơ chế này thường được gọi là reference counting. Nó thuận tiện, nhưng tạo một trade-off: nếu mọi consumer tạm rời đi, công việc upstream bị hủy. Với HTTP request có thể abort hoặc stream kết nối đắt đỏ, hãy quyết định hành vi đó có đúng với nghiệp vụ không thay vì thêm share() theo thói quen.

share chỉ chia sẻ các subscription chồng lấp

Nếu source hoàn tất đồng bộ, subscription thứ nhất có thể kết thúc trước khi dòng subscribe thứ hai chạy:

import { defer, of, share } from 'rxjs';

const sharedOnce$ = defer(() => {
  console.log('producer: chạy');
  return of(42);
}).pipe(
  share(),
);

sharedOnce$.subscribe((value) => console.log('A:', value));
sharedOnce$.subscribe((value) => console.log('B:', value));

Kết quả:

producer: chạy
A: 42
producer: chạy
B: 42

Đây không phải lỗi của share(). Execution đầu đã complete và reset trước khi B đăng ký, nên không còn cửa sổ nào để chia sẻ. Nếu mục tiêu là giữ kết quả cho subscriber đến sau, bạn cần replay/cache semantics chứ không chỉ multicast.

Khi subscriber đến muộn cần giá trị gần nhất

shareReplay() vừa share source vừa chuyển notification qua một ReplaySubject. Ví dụ sau dùng buffer một phần tử để nhiều consumer dùng chung request và consumer đến sau nhận profile gần nhất:

import { defer, from, shareReplay } from 'rxjs';

type Profile = {
  id: string;
  displayName: string;
};

const profile$ = defer(() =>
  from(
    fetch('/api/profile').then(async (response) => {
      if (!response.ok) {
        throw new Error(`HTTP ${response.status}`);
      }

      return (await response.json()) as Profile;
    }),
  ),
).pipe(
  shareReplay({ bufferSize: 1, refCount: true }),
);

Trong RxJS 7.8.2, cấu hình này có các ý nghĩa chính:

  • những subscriber chồng lấp dùng chung một source subscription;
  • subscriber đến muộn có thể nhận tối đa một giá trị đã buffer;
  • nếu source chưa complete và reference count về 0, refCount: true ngắt source;
  • sau khi source complete thành công, giá trị và completion được giữ để phát lại cho subscriber sau trên cùng instance profile$;
  • nếu source error, trạng thái được reset để lần subscribe sau có thể thử một execution mới.

shareReplay không phải chiến lược cache hoàn chỉnh

bufferSize: 1 giới hạn số emission, không định nghĩa TTL, refresh hay invalidation nghiệp vụ. refCount: true giúp dừng source còn active khi không ai nghe, nhưng không tự làm mới một request đã complete. Nếu profile có thể stale, hãy thiết kế trigger refresh hoặc cache boundary rõ ràng.

Mặc định mình dùng share() cho event sống và chỉ chọn shareReplay() khi có yêu cầu rõ rằng late subscriber cần dữ liệu trước đó. Replay “cho chắc” dễ giữ state cũ lâu hơn dự kiến và che mất câu hỏi ai chịu trách nhiệm refresh.

Chọn cold hay hot trong code thực tế

Không có loại nào tốt hơn tuyệt đối. Cold ưu tiên isolation; hot ưu tiên sharing. Chọn dựa trên semantics của operation, không dựa trên việc code nào ngắn hơn.

Tình huốngLựa chọn mặc địnhVì sao
Mỗi thao tác phải chạy độc lậpCold source, thường tạo bằng defer() hoặc creation function phù hợpMỗi consumer có state, error và cancellation riêng.
Nhiều consumer đồng thời cần cùng side effectshare()Một source subscription trong cửa sổ các consumer cùng active.
Consumer đến muộn cần giá trị gần nhấtshareReplay({ bufferSize: 1, refCount: true }) sau khi xác định lifecycleCó sharing và replay hữu hạn, nhưng vẫn phải thiết kế refresh.
Event toàn cục đã tồn tạiHot Observable hoặc Subject được encapsulateObservable không giả vờ sở hữu producer bên ngoài.
State hiện tại phải có ngay khi subscribeStore abstraction, BehaviorSubject hoặc replay có chủ đíchSubject trần chỉ phát event tương lai.
Mỗi subscription phải gửi request mớidefer(() => from(requestFactory()))Factory tạo Promise/request mới đúng lúc subscribe.

Một API thư viện thường nên trả Observable<T> và giữ quyết định share ở boundary biết rõ số consumer và lifecycle. Share quá sớm ở tầng thấp khiến caller khó yêu cầu execution độc lập; share quá muộn có thể làm side effect đắt đỏ chạy nhiều lần.

Trước khi thêm share hay shareReplay, hãy viết ra bốn điều:

  1. side effect nào đang bị chia sẻ;
  2. cửa sổ nào các consumer phải dùng chung execution;
  3. subscriber đến muộn cần future value hay cả giá trị cũ;
  4. điều gì xảy ra khi subscriber cuối rời đi, source complete hoặc source error.

Nếu chưa trả lời được, operator sharing chỉ đang che một lifecycle chưa rõ.

Những bẫy thường gặp

  1. Gắn cold với async và hot với sync. Hai trục này độc lập. of() cold nhưng sync; Subject hot vẫn có thể phát sync.
  2. Cho rằng subscribe hai lần luôn chạy producer hai lần. Điều này đúng với cold source, không đúng với một hot hoặc shared Observable.
  3. Cho rằng cùng giá trị nghĩa là cùng execution. Hai cold producer có thể cùng phát 42; hãy log lúc producer được tạo hoặc đo side effect thực tế.
  4. Dùng from(existingPromise) rồi mong request bắt đầu khi subscribe. Promise đã chạy. Dùng defer(() => from(createPromise())) nếu mỗi subscription cần operation mới.
  5. Mong share() phát lại giá trị cũ. share() mặc định chỉ chuyển emission tương lai cho late subscriber.
  6. Mong share() cache qua các subscription tuần tự. Khi source complete hoặc reference count về 0, cấu hình mặc định reset. Chỉ các subscription chồng lấp mới dùng chung execution.
  7. Thêm shareReplay(1) vào mọi HTTP stream. Bạn đang chọn cả replay và cache-after-completion. Hãy định nghĩa invalidation, error và refresh trước.
  8. Expose Subject để mọi nơi gọi next(). Quyền phát dữ liệu bị phân tán và producer không còn owner rõ ràng. Giữ Subject private, chỉ expose asObservable() nếu consumer chỉ cần đọc.
  9. Để subscription con đóng resource chung. Với hot producer do service sở hữu, unsubscribe một consumer chỉ nên gỡ consumer đó. Lifecycle của WebSocket hoặc event bus cần owner riêng.

Bài tập tự kiểm tra

Lấy ví dụ sharedTicks$ và thử lần lượt:

  1. bỏ share() rồi đếm số dòng producer: tạo timer;
  2. đổi thời điểm B subscribe từ 1_500 ms thành 3_500 ms, sau khi A đã complete;
  3. thay share() bằng shareReplay({ bufferSize: 1, refCount: true }) và quan sát giá trị B nhận ngay khi vào muộn;
  4. giữ source là interval() không take(), unsubscribe cả A lẫn B rồi subscribe C để kiểm tra timer có bắt đầu lại từ 0 không.

Mỗi lần chạy, đừng chỉ ghi output. Hãy vẽ producer nào thuộc subscription nào, lúc nào reference count đổi và notification nào được replay. Nếu ba thứ đó khớp với log, bạn đã có mental model đủ chắc để review các stream HTTP, WebSocket và UI state.

Học tiếp

Nguồn tham khảo

On this page