Học RxJS
Nền tảng

RxJS là gì?

Vai trò của RxJS và vấn đề mà thư viện giải quyết.

Một ô tìm kiếm có thể phát ra hàng chục sự kiện khi người dùng gõ. Trong lúc đó, ứng dụng còn phải chờ timer, nhận WebSocket, gọi API và hủy những công việc đã lỗi thời. Nếu mỗi nguồn bất đồng bộ dùng một kiểu callback riêng, phần khó không còn là nhận dữ liệu mà là phối hợp thứ tự, thời gian, lỗi và cleanup.

RxJS đưa các nguồn ấy về cùng một mô hình: dòng giá trị theo thời gian. Bạn mô tả cách tạo, biến đổi và kết hợp dòng dữ liệu bằng Observable; RxJS lo việc chuyển từng notification đến consumer và cung cấp một cơ chế thống nhất để dừng luồng.

Phạm vi phiên bản

Các API và ví dụ trong bài nhắm đến RxJS 7.8.x (nhánh RxJS 7 ổn định), dùng public import từ rxjs. Bài không giả định API của RxJS 8.

Mục lục

RxJS giải quyết vấn đề gì

JavaScript đã có callback, Promise, async/await và event listener. RxJS không thay thế tất cả những công cụ đó. Nó hữu ích khi bài toán có nhiều giá trị xuất hiện qua thời gian và bạn cần áp dụng các quy tắc nhất quán lên chúng.

Hình dung ba yêu cầu của một ô tìm kiếm:

  1. Chỉ xử lý khi người dùng ngừng gõ 300 ms.
  2. Không gửi lại cùng một từ khóa liên tiếp.
  3. Khi màn hình bị đóng, dừng listener và mọi công việc còn chờ.

Với callback thuần, bạn phải tự giữ timer, so sánh giá trị cũ và nhớ cleanup. Với RxJS, mỗi yêu cầu trở thành một operator trong pipeline, còn việc dừng toàn bộ chuỗi đi qua một Subscription.

Nói ngắn gọn, RxJS là thư viện reactive programming cho JavaScript. Thư viện biểu diễn dữ liệu hoặc sự kiện dưới dạng Observable, rồi cho phép bạn tạo pipeline bằng các operator như map, filter, debounceTime hay switchMap.

Mental model của một pipeline

Đừng hình dung RxJS như một mảng biết chạy bất đồng bộ. Mô hình sát hơn là một dây chuyền có điểm bắt đầu, các trạm xử lý, consumer và công tắc ngắt:

Producer          Observable + operators                    Consumer
(timer, DOM,  ──► [source] ─► [map] ─► [filter] ────────►  Observer
 WebSocket...)          │                                      │
                        └──── tạo execution khi subscribe ─────┘
                                           │
                                           ▼
                                      Subscription
                                           │
                              unsubscribe / complete / error
                                           │
                                           ▼
                                        teardown

Mỗi mũi tên biểu diễn giá trị đi theo hướng producer đến consumer. Teardown đi ngược lại qua chuỗi để giải phóng tài nguyên mà từng tầng đã tạo.

Năm mảnh ghép cốt lõi

Mảnh ghépCâu hỏi nó trả lờiVai trò
ObservableDữ liệu đến từ đâu và được tạo thế nào?Mô tả nguồn cùng cách nối producer với consumer.
ObserverLàm gì khi có giá trị, lỗi hoặc hoàn tất?Consumer gồm các callback next, error, complete; có thể chỉ cung cấp callback cần dùng.
SubscriptionExecution này còn hoạt động không, dừng bằng cách nào?Đại diện cho một execution và có phương thức unsubscribe().
OperatorBiến đổi luồng ra sao?Nhận một Observable, trả về Observable mới để lọc, ánh xạ, kết hợp hoặc điều khiển thời gian.
NotificationProducer đang báo điều gì?next mang giá trị; error và complete là hai tín hiệu kết thúc loại trừ nhau.

Tên biến Observable thường có hậu tố $, chẳng hạn searchTerm$. Đây là convention phổ biến chứ không phải cú pháp bắt buộc; mục đích là giúp người đọc nhận ra biến đó cần được subscribe hoặc tiếp tục đưa qua pipe.

Observable là công thức chứ không phải hộp dữ liệu

Một Observable thường mô tả cách chạy, chưa phải kết quả đã nằm sẵn trong bộ nhớ. pipe(...) xây công thức mới; subscribe(...) mới nối consumer vào nguồn và bắt đầu một execution cho subscription đó.

Từ “thường” ở đây quan trọng. Có nguồn tạo execution riêng cho từng subscriber (cold), nhưng cũng có nguồn đang phát độc lập và được nhiều subscriber dùng chung (hot). Đừng suy ra hành vi chia sẻ chỉ từ kiểu Observable; hãy xem nguồn được tạo và multicast thế nào. Phần này được giải thích riêng trong Cold và hot Observable.

Observable không đồng nghĩa với bất đồng bộ

RxJS có thể phát giá trị đồng bộ. of(1, 2, 3) gọi các callback ngay trong lúc subscribe() đang chạy, còn fromEvent chờ sự kiện tương lai. Thứ tự thực thi phụ thuộc vào source và scheduler, không phụ thuộc riêng vào kiểu Observable.

Ví dụ đầu tiên và vòng đời đầy đủ

Ví dụ sau cố ý dùng nguồn đồng bộ để bạn nhìn rõ thứ tự next, complete và teardown:

import { filter, finalize, map, of } from 'rxjs';

const prices$ = of(100, 200, 300).pipe(
  map((price) => price * 1.2),
  filter((price) => price >= 200),
  finalize(() => console.log('teardown')),
);

console.log('trước subscribe');

const subscription = prices$.subscribe({
  next: (price) => console.log('giá:', price),
  error: (error: unknown) => console.error('lỗi:', error),
  complete: () => console.log('complete'),
});

console.log('sau subscribe');
console.log('closed:', subscription.closed);

Kết quả:

trước subscribe
giá: 240
giá: 360
complete
teardown
sau subscribe
closed: true

Luồng chạy theo thứ tự này:

  1. of phát lần lượt 100, 200, 300 ngay khi có subscriber.
  2. map nhân từng giá trị với 1.2; filter loại 120.
  3. Sau giá trị cuối, of gửi complete. Không notification nào được gửi tiếp sau đó.
  4. finalize chạy khi subscription kết thúc và subscription.closed trở thành true.

prices$ không bị sửa tại chỗ. Mỗi operator trong pipe tạo một Observable mới bao quanh source trước đó. Cách ghép này giúp pipeline đọc từ trái sang phải và có thể tái sử dụng mà không lồng callback.

Ví dụ thực tế với ô tìm kiếm

Giả sử trang HTML có một input với id search. Pipeline dưới đây chỉ chuyển từ khóa hợp lệ cho consumer sau khi người dùng dừng gõ:

import {
  debounceTime,
  distinctUntilChanged,
  filter,
  fromEvent,
  map,
} from 'rxjs';

const input = document.querySelector<HTMLInputElement>('#search');

if (!input) {
  throw new Error('Không tìm thấy #search');
}

const searchTerm$ = fromEvent<InputEvent>(input, 'input').pipe(
  map((event) =>
    (event.currentTarget as HTMLInputElement).value.trim(),
  ),
  debounceTime(300),
  distinctUntilChanged(),
  filter((term) => term.length >= 2),
);

const subscription = searchTerm$.subscribe({
  next: (term) => console.log('Tìm:', term),
  error: (error: unknown) => console.error('Lỗi stream:', error),
});

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

Nếu người dùng gõ r, rx, rxjs với khoảng cách dưới 300 ms, rồi dừng lại, output là:

Tìm: rxjs

fromEvent đăng ký DOM listener khi subscribe. debounceTime(300) chỉ cho giá trị gần nhất đi tiếp sau 300 ms yên lặng; distinctUntilChanged bỏ từ khóa trùng liền nhau; filter loại chuỗi ngắn hơn hai ký tự. Khi pagehide xảy ra, unsubscribe() tháo listener và hủy phần việc đang chờ trong pipeline.

Ví dụ chưa gọi API có chủ ý: mục tiêu của trang overview là làm rõ luồng dữ liệu và lifecycle. Khi thêm request, bạn thường cần một flattening operator như switchMap để request mới thay request cũ; xem switchMap sau khi đã nắm Observable cơ bản.

RxJS khác Promise và event listener ở đâu

Không có công cụ thắng trong mọi tình huống. Chọn theo hình dạng của bài toán:

Công cụHình dạng dữ liệu phù hợpComposition và dừng công việc
Promise và async/awaitMột kết quả hoặc một lỗi; flow tuần tựDễ đọc cho tác vụ một lần. Promise settle đúng một lần; cancellation phải do operation bên dưới hỗ trợ, thường qua AbortController.
addEventListenerMột nguồn event đơn giảnÍt abstraction, nhưng bạn tự quản lý handler, state trung gian, timer và removeEventListener.
RxJS ObservableKhông, một hoặc nhiều giá trị theo thời gianOperator ghép được nhiều nguồn và quy tắc thời gian; Subscription gom đường teardown vào cùng một lifecycle.

Với một lần fetch rồi parse JSON, mình mặc định chọn async/await vì ít khái niệm hơn. Với autocomplete, gesture, polling, WebSocket hoặc nhiều nguồn event phải phối hợp, RxJS thường làm các quan hệ thời gian rõ hơn.

RxJS cũng không biến Promise thành “stream nhiều giá trị”. from(promise) chỉ phát kết quả duy nhất của Promise đó rồi complete; hơn nữa Promise đã bắt đầu thì unsubscribe khỏi Observable wrapper không tự động hủy công việc bên dưới. Muốn hủy network request thật sự, bạn vẫn cần API có cancellation và nối teardown với cơ chế đó.

Teardown là một phần của thiết kế

Một subscription có ba đường kết thúc đáng phân biệt:

  • Source gửi complete: kết thúc bình thường.
  • Source gửi error: kết thúc do lỗi chưa được xử lý.
  • Consumer gọi unsubscribe(): hủy execution vì không còn cần dữ liệu.

Cả ba đều dẫn đến teardown, nhưng unsubscribe không gửi complete cho Observer. Vì vậy, đừng đặt cleanup quan trọng chỉ trong callback complete; hãy dùng teardown của source hoặc operator finalize khi công việc phải chạy cả lúc complete, error và hủy chủ động.

Teardown có thể là clearInterval, tháo DOM listener, đóng kết nối hoặc hủy tác vụ mà API bên dưới cho phép. Với stream sống lâu, quyết định “ai sở hữu subscription và khi nào dừng nó” nên được đưa ra ngay lúc thiết kế, không phải vá thêm sau khi thấy memory leak.

Đọc sâu hơn tại Vòng đời của stream và Subscription và teardown.

Khi nên và không nên chọn RxJS

RxJS đáng giá khi ít nhất một trong các yếu tố sau xuất hiện cùng nhau:

  • Nguồn phát nhiều lần: DOM events, timer, WebSocket, message bus.
  • Logic phụ thuộc thời gian: debounce, throttle, timeout, retry, polling.
  • Nhiều tác vụ cạnh tranh hoặc cần hủy cái cũ khi cái mới đến.
  • Cần kết hợp nhiều nguồn và giữ quy tắc lỗi, hoàn tất, cleanup nhất quán.
  • Pipeline được tái sử dụng hoặc kiểm thử theo timeline.

Ngược lại, đừng dùng RxJS chỉ để bọc mọi giá trị. Một phép tính đồng bộ, một event handler ngắn hoặc chuỗi hai lệnh await thường rõ hơn nếu giữ nguyên JavaScript. Abstraction chỉ có lợi khi nó làm quan hệ dữ liệu và thời gian dễ thấy hơn chi phí phải học operator cùng lifecycle.

Nếu bạn đang cân nhắc cho một feature cụ thể, dùng checklist ở Khi nào dùng RxJS?.

Những hiểu lầm thường gặp

  • “Subscribe càng nhiều càng vô hại.” Với cold Observable, mỗi subscription có thể chạy lại producer, tạo thêm request, timer hoặc listener. Hãy xác định rõ bạn muốn execution riêng hay muốn chia sẻ.
  • “Unsubscribe sẽ gọi complete.” Không. Đây là hủy từ phía consumer; callback complete không chạy.
  • “Có Observable là tự động tránh memory leak.” RxJS cung cấp teardown, nhưng ứng dụng vẫn phải kích hoạt nó đúng lúc đối với stream không tự kết thúc.
  • “RxJS tự giải quyết backpressure.” Producer dạng push vẫn có thể phát nhanh hơn consumer xử lý. Bạn phải chọn chiến lược như lấy mẫu, gom lô, giới hạn tốc độ hoặc thiết kế queue phù hợp.
  • “Lồng subscribe là cách nối hai tác vụ.” Cách này tách lifecycle và error handling thành nhiều nhánh khó quản lý. Thường nên trả về Observable bên trong và dùng flattening operator; xem Nested subscribe.
  • “Subject là state store mặc định.” Subject hữu ích cho multicasting và bridge imperative code, nhưng lạm dụng nó làm data flow khó truy vết. Hãy bắt đầu từ Observable và operator, chỉ dùng Subject khi thật sự cần chủ động phát giá trị.

Lộ trình học tiếp

Đừng học thuộc danh sách operator ngay. Trước hết, hãy chắc rằng bạn nhìn một pipeline và chỉ ra được producer, nơi subscription bắt đầu, ba loại notification và đường teardown. Sau đó đi theo thứ tự sau:

Nguồn tham khảo

On this page