Học RxJS
Bắt đầu

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

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ề Subscription

Bốn ý cần giữ trong đầu:

  1. Observable mô tả nguồn và pipeline.
  2. subscribe() bắt đầu một execution cụ thể và đăng ký consumer.
  3. Observer là tập callback phản ứng với next, error và complete.
  4. 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 = true

Observer<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
JS

Cá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 503

Khô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: true

Thờ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ụ:

  1. networkStatus$ mới chỉ mô tả hai event source và giá trị ban đầu.
  2. subscribe() khiến hai fromEvent() đăng ký listener lên window.
  3. Mỗi event trở thành một next, và Observer cập nhật DOM.
  4. destroyNetworkBanner() đóng subscription; teardown của fromEvent() gỡ cả hai listener.
  5. 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ốngLựa chọn mặc địnhVì sao
Cần xử lý giá trị, lỗi hoặc hoàn tấtsubscribe({ next, error, complete })Tên callback rõ và không phụ thuộc vị trí tham số.
Chỉ in giá trị trong demo ngắnsubscribe(value => ...)Gọn, vẫn hợp lệ trong RxJS 7.x.
Cần biến đổi dữ liệuOperator trong pipe()Giữ pipeline lazy và có thể compose; chưa cần side effect.
Source sống theo UI, socket hoặc timerLưu Subscription hoặc dùng primitive lifecycleCó nơi rõ ràng để hủy và chạy teardown.
Cần cleanup trên complete, error và unsubscribefinalize()Bao phủ cả ba đường kết thúc.
Hàm thư viện muốn cung cấp dữ liệu cho callerTrả 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ặc finalize().
  • 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:

  1. thêm console.log trước và sau subscribe() để xác định giá trị từ startWith() chạy đồng bộ ở đâu;
  2. gọi destroyNetworkBanner(), sau đó bật/tắt mạng và xác nhận banner không đổi nữa;
  3. 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ó callback complete.

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

Nguồn tham khảo

On this page