Học RxJS
Observable & Subscription

Unicast và multicast

Hiểu cách emission được phân phối cho nhiều subscriber.

Hai component cùng cần một dữ liệu, nhưng vừa subscribe() xong bạn thấy hai timer, hai lần tính toán hoặc hai HTTP request. Đó không phải lỗi ngẫu nhiên: mỗi subscriber đang có một producer riêng. Unicast và multicast giúp bạn trả lời câu hỏi quan trọng hơn “có bao nhiêu subscriber?”: các subscriber đó đang dùng nhiều execution độc lập, hay cùng quan sát một execution được chia sẻ?

Phạm vi phiên bản

Bài này dùng public API của RxJS 7.x và đối chiếu hành vi với RxJS 7.8.2. Ví dụ ưu tiên share() và Subject; các multicasting API cũ như multicast(), publish() và refCount() độc lập đã bị deprecated trong RxJS 7.

Mục lục

Mental model: đếm producer, không chỉ đếm subscriber

Hãy hình dung một quán có hai bàn gọi cùng món. Với unicast, bếp nấu hai phần riêng: mỗi bàn có một lần thực hiện độc lập. Với multicast, bếp làm một mẻ rồi chia kết quả cho cả hai bàn đang chờ. Phép so sánh dừng ở việc chia sẻ lần thực hiện; trong RxJS, subscriber đến muộn có nhận được phần trước đó hay không còn phụ thuộc vào Subject hoặc chiến lược replay.

UNicast

subscriber A ──► producer A ──► 0, 1, 2
subscriber B ──► producer B ──► 0, 1, 2

MULTIcast

subscriber A ──┐
               ├──► shared producer ──► 0, 1, 2
subscriber B ──┘

Nói chính xác hơn:

  • Unicast nối một producer với một consumer trong mỗi execution.
  • Multicast cho nhiều consumer quan sát cùng một producer.
  • Mỗi consumer vẫn có Subscription riêng. A unsubscribe không đồng nghĩa B cũng bị hủy.
  • Chia sẻ producer không đồng nghĩa cache dữ liệu. Multicast thuần chỉ phân phối notification đang xảy ra.

Câu hỏi debug hữu ích nhất là: “dòng log bắt đầu producer chạy mấy lần?”. Nếu nó chạy hai lần cho hai lần subscribe, bạn đang có hai execution. Nếu nó chỉ chạy một lần trong lúc cả hai subscriber cùng active, execution đã được chia sẻ.

Unicast: mỗi subscriber có một execution

Plain Observable thường là unicast. Mỗi lần gọi subscribe(), RxJS chạy lại hàm producer cho subscriber đó. Vì vậy state nằm trong producer cũng được tạo riêng.

import { Observable, take } from 'rxjs';

const counter$ = new Observable<number>((subscriber) => {
  console.log('producer: start');

  let value = 0;
  const intervalId = setInterval(() => {
    subscriber.next(value++);
  }, 1_000);

  return () => {
    clearInterval(intervalId);
    console.log('producer: stop');
  };
}).pipe(take(3));

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

Kết quả thường có dạng:

producer: start
producer: start
A: 0
B: 0
A: 1
B: 1
A: 2
producer: stop
B: 2
producer: stop

Hai subscriber tình cờ nhận các số giống nhau, nhưng đó là hai biến value và hai interval khác nhau. Chúng có thể lệch thời điểm, bị hủy riêng và gây side effect riêng.

Unicast không phải điều cần “sửa” mặc định. Nó phù hợp khi mỗi consumer cần isolation: một form submit riêng, một workflow riêng, một phép tính có state riêng, hoặc một request mà mỗi caller chủ động khởi chạy và hủy. Vấn đề chỉ xuất hiện khi bạn tưởng mình đang dùng chung công việc đắt đỏ nhưng thực tế lại chạy nó nhiều lần.

Multicast: nhiều subscriber dùng chung producer

Multicast thêm một điểm phân phối giữa producer và các subscriber. Trong RxJS, điểm phân phối này thường là một Subject ở bên trong operator như share(), hoặc một Subject bạn quản lý trực tiếp.

                         ┌──► Subscription A ──► Observer A
source ──► connection ──► Subject
                         └──► Subscription B ──► Observer B

Ở đây chỉ có một connection từ Subject đến source. Mỗi notification source gửi vào Subject được chuyển đến các Observer đang đăng ký. A và B nhận cùng một emission từ cùng execution, thay vì tự tạo hai emission giống nhau.

Áp share() vào ví dụ trước:

import { Observable, share, take } from 'rxjs';

const sharedCounter$ = new Observable<number>((subscriber) => {
  console.log('producer: start');

  let value = 0;
  const intervalId = setInterval(() => {
    subscriber.next(value++);
  }, 1_000);

  return () => {
    clearInterval(intervalId);
    console.log('producer: stop');
  };
}).pipe(
  take(3),
  share(),
);

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

Lần này producer chỉ bắt đầu một lần:

producer: start
A: 0
B: 0
A: 1
B: 1
A: 2
B: 2
producer: stop

Vị trí đặt share() là một phần của thiết kế. Chỉ phần pipeline ở phía trước share() dùng chung một subscription; operator ở phía sau vẫn chạy riêng cho từng subscriber.

import { map, share, tap } from 'rxjs';

declare const source$: import('rxjs').Observable<number>;

const shared$ = source$.pipe(
  tap(() => console.log('chạy một lần cho mỗi emission của source')),
  map((value) => value * 2),
  share(),
  tap(() => console.log('chạy một lần cho mỗi subscriber')),
);

Nếu phép tính nặng cần được chia sẻ, đặt nó trước share(). Nếu logic phụ thuộc từng consumer, giữ nó sau share().

Unicast, multicast khác cold, hot thế nào?

Hai cặp thuật ngữ liên quan nhưng không phải hai cách gọi cho cùng một việc:

Cặp khái niệmCâu hỏi nó trả lời
Cold / hotProducer được tạo trong lúc subscribe hay đã tồn tại bên ngoài lần subscribe?
Unicast / multicastMột producer được một hay nhiều consumer quan sát?

Một cold Observable tạo producer mới cho mỗi subscription, nên nó luôn unicast. Khi dùng share(), bạn tạo một execution dùng chung trong khoảng thời gian có subscriber; kết quả được xem là hot và multicast. Ngược lại, nói “unicast” chỉ mô tả quan hệ một-một, không tự khẳng định producer cold.

Đừng dùng bốn nhãn như từ đồng nghĩa

“Cold = unicast” đúng theo chiều cold dẫn đến unicast, nhưng không đủ để đảo lại thành “mọi unicast đều cold”. Khi review code, hãy hỏi riêng hai câu: producer được tạo lúc nào, và bao nhiêu subscriber đang dùng chung producer đó?

Trang Cold và hot Observable tập trung vào thời điểm producer bắt đầu và ownership của producer. Trang này tập trung vào cách một execution được phân phối.

Chuyển một pipeline sang multicast bằng share

Trong code ứng dụng, mình chọn share() khi đã có một Observable và muốn các subscriber đồng thời dùng chung source. Operator này quản lý Subject, connection và reference count giúp bạn; ít state thủ công hơn việc tự nối source vào Subject.

Với cấu hình mặc định của RxJS 7.8.2, share():

  1. tạo Subject nội bộ khi cần;
  2. connect đến source khi subscriber đầu tiên đến;
  3. phân phối notification cho mọi subscriber đang active;
  4. unsubscribe source và reset khi số subscriber về 0 trước khi source kết thúc;
  5. reset state sau error hoặc complete, để subscriber đến sau có thể tạo execution mới.

Subscriber đến muộn không nhận lại dữ liệu cũ

share() mặc định dùng Subject, mà Subject không replay emission đã qua. Subscriber đến ở giữa stream chỉ nhận dữ liệu từ thời điểm nó đăng ký.

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

const shared$ = interval(1_000).pipe(
  take(4),
  share(),
);

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

setTimeout(() => {
  shared$.subscribe((value) => console.log('B:', value));
}, 2_500);

Kết quả thường là:

A: 0
A: 1
A: 2
B: 2
A: 3
B: 3

B bỏ lỡ 0 và 1; đây là hành vi đúng của live multicast. Nếu B cần giá trị gần nhất hoặc một buffer lịch sử, đó là yêu cầu replay, không chỉ multicast. Khi ấy hãy cân nhắc shareReplay() hoặc connector dùng ReplaySubject, đồng thời thiết kế rõ thời gian sống của cache.

Connection sống theo số subscriber

Mỗi lần subscribe vào kết quả của share(), reference count tăng. Khi một subscriber unsubscribe, chỉ subscription của nó đóng; các subscriber còn lại vẫn nhận dữ liệu. Nếu source chưa complete hay error và người cuối cùng rời đi, mặc định share() unsubscribe connection đến source rồi reset state.

refCount: 0 ── A subscribe ──► 1  (connect source)
refCount: 1 ── B subscribe ──► 2  (dùng connection cũ)
refCount: 2 ── A unsubscribe ► 1  (source vẫn chạy)
refCount: 1 ── B unsubscribe ► 0  (disconnect và reset)

Điều này đặc biệt quan trọng với timer, DOM listener và socket: chia sẻ không xóa nhu cầu quản lý lifecycle. Nó chỉ chuyển ownership của connection sang nhóm subscriber.

share() có các option như resetOnError, resetOnComplete và resetOnRefCountZero. Đừng tắt reset chỉ để “giữ stream sống” mà chưa xác định ai sẽ đóng connection; với source vô hạn, lựa chọn đó có thể để producer chạy dù không còn consumer.

Complete và error không tự retry subscriber hiện tại

Khi source complete hoặc error, share() chuyển terminal notification đó cho các subscriber hiện tại, và subscription của họ đóng. Việc reset mặc định chỉ cho phép subscriber mới tạo connection mới; nó không hồi sinh các subscription đã đóng.

Nếu muốn tự retry sau lỗi, đặt retry() có chủ đích trong pipeline. Nếu muốn lặp lại sau complete, dùng repeat(). Vị trí tương đối giữa các operator đó và share() quyết định cả nhóm dùng chung một lần retry hay mỗi nhánh tạo hành vi riêng, nên hãy viết test cho lifecycle thay vì suy đoán từ tên operator.

Subject: multicast bằng một boundary imperative

Subject vừa là Observer vừa là Observable

Subject<T> có hai mặt:

  • là Observer, nên code có thể gọi next, error và complete để đưa notification vào;
  • là Observable, nên nhiều consumer có thể subscribe() để nhận cùng notification.
import { Subject } from 'rxjs';

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

notifications$.subscribe((message) => console.log('toast:', message));
notifications$.subscribe((message) => console.log('audit:', message));

notifications$.next('Đã lưu hồ sơ');
notifications$.next('Đã đồng bộ dữ liệu');
notifications$.complete();

Kết quả:

toast: Đã lưu hồ sơ
audit: Đã lưu hồ sơ
toast: Đã đồng bộ dữ liệu
audit: Đã đồng bộ dữ liệu

Một Subject thường giống event bus cục bộ: producer đẩy event vào một đầu, nhiều listener nhận ở đầu kia. Nó phù hợp ở boundary imperative, chẳng hạn adapter nối callback API vào RxJS hoặc một service nhận command từ code không reactive.

Tuy nhiên, nếu mục tiêu chỉ là “chia sẻ pipeline này”, ưu tiên share() thay vì tự viết source$.subscribe(subject). Operator quản lý connection, error, complete và reset rõ hơn; Subject thủ công buộc bạn tự sở hữu tất cả các quyết định lifecycle đó.

Ẩn quyền ghi bằng asObservable

Đừng phát công khai cả Subject nếu consumer chỉ cần đọc. Nếu ai cũng gọi được next(), bạn khó truy ra producer nào đã tạo event và invariant rất dễ bị phá.

import { Observable, Subject } from 'rxjs';

class NotificationBus {
  private readonly messagesSubject = new Subject<string>();

  readonly messages$: Observable<string> =
    this.messagesSubject.asObservable();

  publish(message: string): void {
    const normalized = message.trim();

    if (normalized.length > 0) {
      this.messagesSubject.next(normalized);
    }
  }

  destroy(): void {
    this.messagesSubject.complete();
  }
}

asObservable() không tạo bản sao dữ liệu và không biến stream trở lại unicast. Nó chỉ che các method ghi khỏi public API, để class giữ quyền phát và kết thúc stream.

Quy tắc chọn nhanh

Có sẵn pipeline và chỉ muốn chia sẻ execution: dùng share(). Cần một cổng imperative để code bên ngoài đẩy event vào: cân nhắc Subject, giữ nó private và public một Observable chỉ đọc.

Ví dụ thực tế: chia sẻ luồng polling cho dashboard

Giả sử dashboard có bảng dữ liệu và badge trạng thái. Cả hai cùng cần snapshot mới mỗi 5 giây. Nếu mỗi widget subscribe vào một pipeline unicast, mỗi widget tạo timer và request riêng.

import { ajax } from 'rxjs/ajax';
import { map, share, switchMap, timer } from 'rxjs';

type SystemSnapshot = {
  services: Array<{ name: string; status: 'up' | 'down' }>;
  collectedAt: string;
};

const snapshot$ = timer(0, 5_000).pipe(
  switchMap(() =>
    ajax.getJSON<SystemSnapshot>('/api/system-snapshot'),
  ),
  share(),
);

const tableSubscription = snapshot$.subscribe({
  next: (snapshot) => {
    console.log('table:', snapshot.services);
  },
  error: (error: unknown) => {
    console.error('table error:', error);
  },
});

const badgeSubscription = snapshot$
  .pipe(
    map((snapshot) =>
      snapshot.services.every((service) => service.status === 'up'),
    ),
  )
  .subscribe({
    next: (allHealthy) => {
      console.log('badge:', allHealthy ? 'healthy' : 'degraded');
    },
    error: (error: unknown) => {
      console.error('badge error:', error);
    },
  });

// Khi bảng bị tháo, badge vẫn giữ shared connection hoạt động.
tableSubscription.unsubscribe();

// Khi consumer cuối cùng rời đi, share() dừng connection đến timer.
badgeSubscription.unsubscribe();

Ở pipeline này, timer, switchMap và request nằm trước share(), nên bảng và badge dùng cùng chu kỳ polling và cùng response khi cả hai đang active. map() của badge nằm sau share(), vì đó là cách diễn giải dữ liệu riêng của badge.

Có một giới hạn cần nói rõ: nếu badge mount sau khi response vừa phát, share() không phát lại snapshot cũ; badge chờ lần polling tiếp theo. Nếu UI phải render ngay từ snapshot gần nhất, hãy chuyển yêu cầu thành cache có replay và xác định rõ cache reset lúc nào. Đừng thay share() bằng shareReplay(1) như một phản xạ, vì bạn vừa thay đổi cả memory và lifecycle của stream.

Chọn unicast hay multicast

Nhu cầuLựa chọn mặc địnhLý do
Mỗi consumer cần execution và state riêngGiữ Observable unicastIsolation rõ; hủy một execution không ảnh hưởng execution khác.
Nhiều consumer đồng thời cần cùng side effect hoặc computationshare()Dùng chung source subscription trong khi vẫn giữ Subscription riêng cho từng consumer.
Consumer đến muộn cần dữ liệu trước đóReplay strategy có giới hạnMulticast thuần không lưu lịch sử; cần chỉ rõ buffer và lifecycle cache.
Code imperative cần phát event cho nhiều listenerSubject private + asObservable()Boundary ghi rõ ràng mà consumer không có quyền gọi next().
Cần chủ động quyết định thời điểm connectconnectable() hoặc connect()Phù hợp khi first subscriber không phải tín hiệu bắt đầu đúng.
Mỗi lần gọi là một command độc lậpUnicastShare command có thể gộp nhầm những lần thực hiện lẽ ra tách biệt.

Mặc định, mình giữ unicast cho đến khi có lý do chia sẻ cụ thể: tránh side effect lặp, đồng bộ cùng một live source, hoặc giảm computation trùng. Multicast thêm state và lifecycle dùng chung; dùng nó chỉ vì “có nhiều subscriber” thường che mất yêu cầu thật.

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

  1. Đếm subscriber nhưng không đếm source subscription. Hai subscriber không tự động có nghĩa là hai producer, và cũng không tự động có nghĩa là một producer. Đặt log hoặc tap({ subscribe, finalize }) quanh source để quan sát connection thật.
  2. Dùng share() như cache. Subscriber đến muộn không nhận lại emission cũ. Multicast và replay là hai yêu cầu khác nhau.
  3. Share một request đã complete rồi mong dùng lại kết quả. Với reset mặc định, subscriber đến sau completion có thể tạo request mới. Nếu muốn cache response, hãy thiết kế replay, invalidation và error policy.
  4. Đặt share() sai vị trí. Operator đắt đỏ đứng sau share() vẫn chạy cho từng subscriber. Ngược lại, đưa logic phụ thuộc consumer lên trước share() có thể vô tình trộn state giữa các consumer.
  5. Public Subject cho cả ứng dụng. Mọi nơi đều có thể next, error hoặc complete, khiến data flow mất owner. Giữ Subject private và cung cấp method có ý nghĩa nghiệp vụ.
  6. Quên terminal notification ảnh hưởng cả nhóm. Một source error hoặc complete sẽ kết thúc tất cả subscriber đang dùng connection đó.
  7. Tắt resetOnRefCountZero cho source vô hạn mà không có owner. Source có thể tiếp tục chạy khi không còn ai nghe. Nếu thật sự cần connection sống lâu, phải có nơi chịu trách nhiệm đóng nó.
  8. Dùng multicasting API cũ cho code mới. Trong RxJS 7, ưu tiên share, connect và connectable; multicast, publish và các biến thể cũ đã deprecated.

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

Lấy ví dụ counter$ ở đầu bài và thử theo thứ tự:

  1. unsubscribe A sau emission 0, rồi kiểm tra B của bản unicast và bản share() có tiếp tục nhận hay không;
  2. cho B subscribe sau khoảng 2,5 giây để xác nhận B không nhận lại 0 và 1 trong bản multicast;
  3. để cả A và B unsubscribe trước take(3), sau đó subscribe C và quan sát producer bắt đầu lại từ 0 với cấu hình mặc định;
  4. di chuyển một map() có log từ trước xuống sau share() và đếm số lần log chạy.

Nếu bạn dự đoán đúng cả bốn kết quả, bạn đã nắm phần cốt lõi: multicast chia sẻ execution đang diễn ra, không tự động chia sẻ lịch sử, và mỗi subscriber vẫn sở hữu lifecycle downstream của riêng mình.

Nguồn tham khảo

Học tiếp

On this page