Học RxJS
Nền tảng

Khi nào dùng RxJS?

Cách nhận biết bài toán phù hợp và khi nên chọn giải pháp đơn giản hơn.

Bạn có một ô tìm kiếm: chờ người dùng ngừng gõ, bỏ từ khóa trùng, hủy request cũ khi có từ khóa mới, rồi dọn listener khi màn hình đóng. Từng việc đều không khó; phần khó là giữ đúng quan hệ giữa chúng theo thời gian. Đây là kiểu bài toán RxJS làm tốt.

Ngược lại, nếu bạn chỉ cần gọi một API rồi nhận một kết quả, async/await thường rõ hơn. RxJS đáng dùng khi stream và cách các stream phối hợp là phần cốt lõi của bài toán, chứ không phải chỉ vì code có bất đồng bộ.

Phiên bản trong bài

Các API và cách import dưới đây nhắm tới RxJS 7.8.x, thuộc nhánh RxJS 7.x ổn định. Bài không giả định API của RxJS 8.

Mục lục

Câu trả lời ngắn

Mình sẽ cân nhắc RxJS khi có đủ hai yếu tố:

  1. Dữ liệu là một chuỗi giá trị xuất hiện theo thời gian, có thể không biết trước lúc kết thúc.
  2. Code phải mô tả quan hệ giữa các giá trị đó: lọc, gộp, trì hoãn, chuyển sang tác vụ mới, retry, chia sẻ hoặc dừng theo một tín hiệu khác.

Một phép thử hữu ích là nhìn vào code imperative hiện tại. Nếu nó đang tích lũy nhiều cờ như isLoading, latestRequestId, nhiều timer, nhiều listener và nhiều nhánh cleanup để giữ đúng thứ tự, một pipeline RxJS có thể làm luật thời gian lộ rõ hơn. Nhưng nếu thay ba dòng await bằng một chuỗi operator dài, bạn đang thêm abstraction chứ chưa giảm độ phức tạp.

Mental model để ra quyết định

Nhìn bài toán như dữ liệu theo thời gian

Observable không chỉ là “Promise có nhiều giá trị”. Hãy hình dung một băng chuyền: producer đặt các giá trị lên theo thời gian, các operator biến đổi băng chuyền, còn Observer nhận next, error hoặc complete. Ví dụ này giúp hình dung luồng dữ liệu, nhưng đừng kéo phép so sánh quá xa: Observable còn có thể là lazy, tạo execution riêng cho từng subscriber hoặc chia sẻ producer tùy cách thiết kế.

Producer              Pipeline                         Observer
DOM event ──► debounce ──► distinct ──► switchMap ──► next(value)
HTTP       ──► catchError ──────────────────────────► error(err)
finite source ──────────────────────────────────────► complete()
                         ▲
                         │ điều khiển execution
                    Subscription
                         │
                         └── unsubscribe() ──► teardown

Ba notification có quy tắc rõ ràng:

  • next có thể xuất hiện không lần nào, một lần hoặc nhiều lần.
  • error kết thúc execution; sau đó không còn next hay complete.
  • complete cũng kết thúc execution; sau đó không còn next hay error.

Bạn có thể đọc kỹ hơn ở Reactive programming, Observer pattern và Vòng đời của stream.

Subscription là quyền sở hữu một execution

subscribe() trả về một Subscription. Đối tượng này không chỉ là “tay cầm để dừng”; nó cho biết đoạn code nào đang sở hữu execution và phải dọn tài nguyên lúc nào. Teardown có thể gỡ DOM listener, hủy timer, đóng socket hoặc abort request nếu producer thực sự hỗ trợ việc đó.

Có một khác biệt dễ bỏ sót: gọi unsubscribe() sẽ chạy teardown nhưng không gọi callback complete của Observer. complete là tín hiệu từ producer; unsubscribe là quyết định của consumer dừng nghe. Vì vậy, đừng đặt cleanup bắt buộc chỉ trong callback complete.

Cây quyết định

Sơ đồ dưới đây cố ý dùng ASCII vì project hiện chỉ đăng ký các MDX component mặc định, chưa cấu hình Mermaid.

Bắt đầu
   │
   ▼
Dữ liệu có nhiều giá trị theo thời gian?
   ├── Không ──► Một tác vụ async duy nhất?
   │               ├── Có  ──► Promise + async/await
   │               └── Không ─► Function / Array / state thường
   │
   └── Có
       │
       ▼
Cần phối hợp time, hủy, gộp nguồn, retry hoặc chia sẻ?
       ├── Không ──► Event listener / async iterator đơn giản
       └── Có
           │
           ▼
Pipeline có rõ hơn state machine viết tay và team bảo trì được?
           ├── Không ──► Giữ primitive đơn giản
           └── Có  ──► Cân nhắc RxJS, rồi định nghĩa teardown

Đây không phải luật tuyệt đối. Một stream chỉ có một giá trị vẫn có thể nằm trong pipeline RxJS nếu nó cần phối hợp với các stream khác. Ngược lại, một WebSocket có vô số message vẫn chưa bắt buộc phải dùng RxJS nếu một callback nhỏ đã đủ và không có logic phối hợp.

Khi RxJS là lựa chọn tốt

Nhiều giá trị đến theo thời gian

DOM events, WebSocket messages, thay đổi route, trạng thái kết nối, sensor data và timer đều tự nhiên là stream. RxJS phù hợp khi bạn muốn xử lý chúng như một chuỗi thống nhất thay vì rải callback qua nhiều nơi.

Điểm quyết định không nằm ở số lượng event, mà ở việc bạn cần diễn đạt quy tắc trên timeline. Chẳng hạn, “chỉ lấy giá trị cuối sau 300 ms im lặng” là một quy tắc thời gian; debounceTime(300) nói đúng ý đó trực tiếp.

Quy tắc phối hợp mới là phần khó

RxJS tỏ ra hữu ích khi phải trả lời các câu như:

  • Query mới đến thì hủy query cũ hay để chạy song song?
  • Hai request phải chạy nối tiếp, song song hay request sau phụ thuộc request trước?
  • Chỉ cập nhật UI khi cả user, permission và route đã có giá trị mới nhất?
  • Khi mất kết nối thì retry bao nhiêu lần, chờ bao lâu và dừng theo tín hiệu nào?

Các operator làm lựa chọn này hiện ngay trong pipeline. Ví dụ, switchMap giữ inner Observable mới nhất; concatMap xếp hàng; mergeMap cho phép chạy đồng thời; exhaustMap bỏ qua nguồn mới khi tác vụ hiện tại chưa xong. Chọn sai operator vẫn gây race condition, nhưng ít nhất policy không bị giấu trong nhiều callback. Phần Higher-order Observable đi sâu vào nhóm quyết định này.

Hủy và teardown là yêu cầu thật

Autocomplete, route change, component unmount và reconnect đều cần dừng công việc cũ. RxJS cho bạn một mô hình teardown thống nhất qua unsubscription.

Tuy nhiên, unsubscription không phải phép màu. Nó chỉ dừng tài nguyên nền nếu Observable đã khai báo teardown tương ứng. Bọc một Promise không thể hủy bằng from(promise) sẽ ngừng chuyển kết quả đến subscriber sau khi unsubscribe, nhưng không tự làm request bên dưới biến mất. Với fetch, hãy nối teardown với AbortController, như ví dụ ở phần sau.

Pipeline được dùng ở nhiều nơi

Một Observable có thể đóng gói chính sách biến đổi và error handling để nhiều consumer dùng lại. Khi cần chia sẻ cùng một producer, multicasting giúp tránh tạo lại công việc cho từng subscriber.

Nhưng chia sẻ cũng đưa vào câu hỏi về cache, thời điểm kết nối, reset sau lỗi và vòng đời dữ liệu. Đừng thêm shareReplay chỉ để “cho nhanh”; hãy xác định rõ ai sở hữu cache và khi nào cache hết hạn. Xem Cold và hot Observable trước khi thiết kế phần này.

Khi nên chọn giải pháp đơn giản hơn

Bài toánLựa chọn mặc địnhKhi nào mới cân nhắc RxJS
Biến đổi một collection đồng bộArray.map, filter, reduceCác phần tử thực sự đến theo thời gian hoặc phải phối hợp với stream khác
Một request, một kết quảPromise và async/awaitRequest nằm trong pipeline cần hủy, retry, gộp hoặc phụ thuộc event
Một nút có một handler ngắnaddEventListenerEvent cần debounce, kết hợp nguồn khác hoặc cleanup theo lifecycle chung
State cục bộ đơn giảnState primitive của frameworkState là kết quả của nhiều stream và có policy chia sẻ rõ ràng
Vòng lặp dữ liệu async tuần tựfor await...ofCần graph nhiều nguồn, concurrency policy hoặc operator theo thời gian

Ví dụ này không cần RxJS:

interface User {
  id: string;
  name: string;
}

export async function loadUser(id: string): Promise<User> {
  const response = await fetch(`/api/users/${encodeURIComponent(id)}`);

  if (!response.ok) {
    throw new Error(`Không tải được user: HTTP ${response.status}`);
  }

  return response.json() as Promise<User>;
}

Có đúng một input, một request và một output. async/await giữ control flow tuyến tính, quen thuộc và dễ debug hơn. Chỉ chuyển nó thành Observable khi bối cảnh bên ngoài tạo ra nhu cầu stream thật, chẳng hạn id thay đổi liên tục và request cũ phải bị hủy.

Đừng chọn theo số dòng code

Pipeline ngắn chưa chắc dễ hiểu hơn, còn code imperative dài chưa chắc sai. Hãy so sánh số trạng thái phải giữ trong đầu, cách thể hiện policy thời gian và đường cleanup — đó mới là chi phí bảo trì thật.

Ví dụ thực tế typeahead có hủy request

Giả sử trang có một input #product-search, một danh sách #results, và endpoint GET /api/products?q=... trả về JSON dạng Product[]. Yêu cầu là chỉ tìm từ hai ký tự, chờ 300 ms sau lần gõ cuối, bỏ query liên tiếp giống nhau và hủy request cũ khi query mới bắt đầu.

Luồng xử lý

input event:  r ── rx ───── rxj ───────── rxjs ─────────────►
                         chờ 300 ms     chờ 300 ms
query sau debounce:         rxj             rxjs
                              │               │
request:                 start rxj ──X    start rxjs ──► kết quả
                                      query mới
                                      abort request cũ

Code TypeScript

import {
  Observable,
  catchError,
  debounceTime,
  distinctUntilChanged,
  fromEvent,
  map,
  of,
  switchMap,
} from 'rxjs';

type Product = {
  id: number;
  name: string;
};

function searchProducts(query: string): Observable<readonly Product[]> {
  return new Observable<readonly Product[]>((subscriber) => {
    const controller = new AbortController();
    console.log(`[request] start ${query}`);

    fetch(`/api/products?q=${encodeURIComponent(query)}`, {
      signal: controller.signal,
    })
      .then((response) => {
        if (!response.ok) {
          throw new Error(`HTTP ${response.status}`);
        }

        return response.json() as Promise<Product[]>;
      })
      .then((products) => {
        subscriber.next(products);
        subscriber.complete();
      })
      .catch((error: unknown) => {
        if (error instanceof DOMException && error.name === 'AbortError') {
          return;
        }

        subscriber.error(error);
      });

    return () => {
      console.log(`[request] teardown ${query}`);
      controller.abort();
    };
  });
}

const input = document.querySelector<HTMLInputElement>('#product-search');
const results = document.querySelector<HTMLUListElement>('#results');

if (!input || !results) {
  throw new Error('Thiếu #product-search hoặc #results');
}

const subscription = fromEvent<InputEvent>(input, 'input')
  .pipe(
    map(() => input.value.trim()),
    debounceTime(300),
    distinctUntilChanged(),
    switchMap((query) => {
      if (query.length < 2) {
        return of<readonly Product[]>([]);
      }

      return searchProducts(query).pipe(
        catchError((error: unknown) => {
          console.error('[request] failed', error);
          return of<readonly Product[]>([]);
        }),
      );
    }),
  )
  .subscribe({
    next: (products) => {
      const items = products.map((product) => {
        const item = document.createElement('li');
        item.textContent = product.name;
        return item;
      });

      results.replaceChildren(...items);
      console.log(`[ui] ${products.length} kết quả`);
    },
    error: (error: unknown) => console.error('[stream] failed', error),
    complete: () => console.log('[stream] complete'),
  });

window.addEventListener(
  'pagehide',
  () => subscription.unsubscribe(),
  { once: true },
);

Kết quả và vòng đời

Giả sử request rxj đã bắt đầu, rồi người dùng nhập rxjs trước khi request đó trả về; API sau cùng trả về hai sản phẩm. Console minh họa sẽ có dạng:

[request] start rxj
[request] teardown rxj
[request] start rxjs
[ui] 2 kết quả
[request] teardown rxjs

debounceTime chỉ phát query sau 300 ms yên lặng, còn distinctUntilChanged bỏ hai query liên tiếp giống nhau. Khi rxjs đến, switchMap unsubscribe inner Observable của rxj; teardown gọi AbortController.abort(), nên request cũ thực sự bị hủy thay vì chỉ bị bỏ kết quả.

Request rxjs phát một mảng qua next, sau đó inner Observable complete. RxJS chạy teardown cả khi inner hoàn tất bình thường, nên bạn vẫn thấy dòng teardown rxjs; gọi abort() sau khi fetch đã xong không làm thay đổi kết quả.

Outer Observable tạo bởi fromEvent là long-lived và không tự complete trong ví dụ này. Khi trang phát pagehide, subscription.unsubscribe() gỡ listener do fromEvent đăng ký và hủy inner request nếu nó còn chạy. Dòng [stream] complete không xuất hiện chỉ vì unsubscribe.

catchError được đặt bên trong switchMap. Vì vậy, một request lỗi được đổi thành mảng rỗng nhưng stream input vẫn sống để nhận query tiếp theo. Nếu đặt error handling sai chỗ, lỗi của một request có thể kết thúc toàn bộ typeahead. Xem Error channel để hiểu quy tắc này.

Các bẫy thường gặp

  • Dùng RxJS cho mọi thứ. Một giá trị đồng bộ không cần trở thành of(value), và một request độc lập không cần pipeline chỉ để trông “reactive”.
  • Nested subscribe. subscribe bên trong subscribe giấu dependency, error path và teardown. Hãy chọn flattening operator theo concurrency policy; xem Nested subscribe.
  • Quên dọn stream không tự kết thúc. DOM event, interval và WebSocket thường sống lâu. Gắn subscription vào lifecycle của màn hình hoặc service; xem Subscription và teardown.
  • Cho rằng unsubscribe luôn hủy công việc nền. Nó chỉ làm vậy khi producer trả về teardown có khả năng hủy, như clearInterval, removeEventListener hoặc AbortController.abort().
  • Biến Subject thành global event bus. Khi ai cũng có thể gọi next, luồng dữ liệu và quyền sở hữu trở nên khó truy vết. Ưu tiên Observable chỉ đọc ở boundary và giới hạn nơi được phép phát giá trị.
  • Chọn flattening operator theo thói quen. switchMap hợp với “mới nhất thắng”, nhưng sai cho thao tác ghi không được phép mất. Đọc Race condition trước khi orchestration request quan trọng.

Checklist trước khi đưa RxJS vào tính năng

  • Bài toán có một sequence theo thời gian, không chỉ một giá trị đơn lẻ.
  • Có ít nhất một nhu cầu phối hợp rõ ràng: time, cancellation, concurrency, combination, retry hoặc sharing.
  • Pipeline làm policy dễ đọc hơn các cờ và callback hiện tại.
  • Team hiểu các operator được chọn và có thể test timeline quan trọng.
  • Mỗi long-lived subscription có owner và thời điểm teardown.
  • Producer bên dưới thật sự hỗ trợ cleanup cần thiết.
  • Error boundary được đặt có chủ đích để lỗi không giết nhầm stream dài hạn.
  • Đã thử viết phương án bằng Promise, event listener hoặc state primitive và biết vì sao nó kém phù hợp hơn.

Nếu chưa đánh dấu được ba mục đầu, mình sẽ giữ giải pháp đơn giản. RxJS nên mua lại độ phức tạp mà bài toán vốn đã có; nó không nên tạo thêm một lớp phức tạp chỉ để code trông đồng nhất.

Học tiếp

Nguồn tham khảo

On this page