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ì
- Mental model của một pipeline
- Ví dụ đầu tiên và vòng đời đầy đủ
- Ví dụ thực tế với ô tìm kiếm
- RxJS khác Promise và event listener ở đâu
- Teardown là một phần của thiết kế
- Khi nên và không nên chọn RxJS
- Những hiểu lầm thường gặp
- Lộ trình học tiếp
- Nguồn tham khảo
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:
- Chỉ xử lý khi người dùng ngừng gõ 300 ms.
- Không gửi lại cùng một từ khóa liên tiếp.
- 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
│
▼
teardownMỗ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ép | Câu hỏi nó trả lời | Vai trò |
|---|---|---|
| Observable | Dữ 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. |
| Observer | Là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. |
| Subscription | Execution 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(). |
| Operator | Biế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. |
| Notification | Producer đ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: trueLuồng chạy theo thứ tự này:
ofphát lần lượt100,200,300ngay khi có subscriber.mapnhân từng giá trị với1.2;filterloại120.- Sau giá trị cuối,
ofgửicomplete. Không notification nào được gửi tiếp sau đó. finalizechạy khi subscription kết thúc vàsubscription.closedtrở thànhtrue.
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: rxjsfromEvent đă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ợp | Composition và dừng công việc |
|---|---|---|
Promise và async/await | Mộ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. |
addEventListener | Mộ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 Observable | Không, một hoặc nhiều giá trị theo thời gian | Operator 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
completekhô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
subscribelà 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:
Reactive programming
Xây mô hình tư duy dữ liệu thay đổi theo thời gian.
Observer pattern
Hiểu quan hệ giữa Observable, Observer và Subscription.
Cài đặt RxJS
Thiết lập môi trường TypeScript để chạy ví dụ.
Pipe và operator
Ghép các phép biến đổi thành pipeline dễ đọc.
Nguồn tham khảo
- RxJS — Introduction — tổng quan chính thức về Observable, Observer, Subscription và operators.
- RxJS 7.8.2 source —
Observablevà source —Subscription— hành visubscribe, trạng tháiclosedvà finalizer. - RxJS 7.8.2 source —
of— phát các đối số đồng bộ rồi complete. - RxJS 7.8.2 source —
fromEvent— hành vi đăng ký và tháo event handler. - RxJS 7.8.2 source —
finalize— callback khi complete, error hoặc unsubscribe. - Gói
rxjstrên npm — phiên bản phát hành và cách cài package.