Reactive programming
Mô hình lập trình phản ứng với dữ liệu và sự kiện theo thời gian.
Khi một ô tìm kiếm, WebSocket và trạng thái đăng nhập cùng thay đổi theo thời gian, cách viết “đọc giá trị rồi xử lý” nhanh chóng biến thành nhiều callback và biến tạm khó đồng bộ. Reactive programming đổi câu hỏi từ “lúc nào mình phải gọi hàm tiếp theo?” thành “khi dữ liệu thay đổi, kết quả nào phải thay đổi theo?”.
Phiên bản trong bài
Các API và import bên dưới nhắm tới RxJS 7.8.x, thuộc dòng RxJS 7.x ổn định. Bài không giả định hành vi của RxJS 8.
Mục lục
- Reactive programming giải quyết điều gì?
- Mental model: giá trị chạy trên một dòng thời gian
- Từ mental model đến RxJS
- Ví dụ thực tế: tìm kiếm khi người dùng nhập
- Chọn reactive hay cách đơn giản hơn
- Những lỗi tư duy thường gặp
- Checklist thiết kế một stream
- Học tiếp
- Nguồn tham khảo
Reactive programming giải quyết điều gì?
Trong code tuần tự, bạn có đầu vào, gọi một hàm, rồi nhận kết quả. Mô hình này rất hợp với dữ liệu đã có sẵn. Nhưng UI event, timer, response từ server hay trạng thái kết nối không xuất hiện cùng một lúc; chúng tạo ra nhiều giá trị theo thời gian.
Cách imperative thường giữ biến trạng thái rồi cập nhật biến đó trong từng callback. Khi yêu cầu thêm debounce, bỏ kết quả request cũ, xử lý lỗi và dọn listener, luồng điều khiển bị rải ra nhiều nơi. Reactive programming gom các quy tắc ấy thành một pipeline: nguồn phát dữ liệu đi qua các phép biến đổi, sau đó consumer nhận kết quả.
Nguồn dữ liệu Quy tắc xử lý Consumer
──────────────► ───────────────────────────────► ─────────────────
input events trim → debounce → chỉ lấy mới nhất render kết quả
│
└── error / complete / teardownBạn có thể hình dung pipeline như một băng chuyền: mỗi trạm chỉ làm một việc rồi chuyển món hàng tiếp theo. Phép so sánh này giúp hiểu composition, nhưng nó dừng ở chỗ Observable còn có error, complete và cơ chế hủy; băng chuyền ngoài đời không mô tả đủ ba phần đó.
Điểm cần giữ trong đầu
Reactive programming không có nghĩa là “dùng RxJS ở mọi nơi”. Đây là cách mô hình hóa quan hệ phụ thuộc theo thời gian; RxJS là một bộ công cụ để biểu diễn và kết hợp các quan hệ đó bằng Observable.
Mental model: giá trị chạy trên một dòng thời gian
Thay vì coi một biến là một ô chứa giá trị hiện tại, hãy coi dữ liệu thay đổi là một chuỗi emission:
searchQuery$: ──"r"──"rx"──"rxjs"────────────────────► thời gian
│ │ │
filter >= 2 ký tự ─────"rx"──"rxjs"────────────────────►
debounce 300 ms ─────────────"rxjs"───────────────────►
request$: ──────────────[kết quả cho "rxjs"]────►Dấu $ chỉ là convention phổ biến để báo rằng biến chứa Observable, không phải cú pháp bắt buộc của RxJS. Giá trị không bị “đẩy vào biến $”; Observable mô tả cách các giá trị sẽ được tạo và chuyển tiếp khi có subscriber.
Bốn câu hỏi cho mọi pipeline
Khi đọc hoặc thiết kế một stream, hãy trả lời theo thứ tự này:
- Nguồn là gì? DOM event, mảng, timer, WebSocket hay response HTTP?
- Mỗi emission có kiểu gì?
InputEvent,string,Product[]hay một state object? - Quy tắc theo thời gian là gì? Lọc, biến đổi, trì hoãn, kết hợp hay thay thế công việc cũ?
- Stream kết thúc và dọn dẹp ra sao? Tự
complete, gặperror, hay cần chủ độngunsubscribe?
Bốn câu hỏi này quan trọng hơn việc nhớ tên operator. Khi lifecycle chưa rõ, một pipeline chạy đúng ở demo vẫn có thể giữ DOM listener hoặc request lâu hơn component sở hữu nó.
Push không đồng nghĩa với bất đồng bộ
Observable là một cơ chế push: producer quyết định lúc gửi next cho consumer. Tuy vậy, emission có thể đồng bộ hoặc bất đồng bộ tùy nguồn.
from([1, 2, 3])phát đồng bộ trong lúcsubscribe()đang chạy.fromEvent(button, 'click')chờ event tương lai.interval(1000)phát theo scheduler sau mỗi khoảng thời gian.
Vì vậy, đừng mặc định “đã dùng Observable thì code chạy ở background” hoặc chạy song song. RxJS không tự tạo thread mới; scheduler và API nguồn mới quyết định thời điểm thực thi.
Từ mental model đến RxJS
Trong RxJS, mental model trên thường ánh xạ thành ba phần:
- Observable mô tả nguồn và cách giá trị được phát.
- Operator nhận Observable đầu vào và trả về Observable mới, nhờ đó pipeline có thể composition bằng
pipe(...). - Observer nhận notification qua
next,errorvàcomplete; lời gọisubscribe(...)trả về một Subscription để quản lý execution và teardown.
Chi tiết quan hệ giữa các vai trò này nằm ở bài Observer pattern. Ở đây, điều quan trọng là tách phần mô tả pipeline khỏi phần kích hoạt side effect.
Một pipeline nhỏ nhưng đủ vòng đời
Ví dụ sau chuyển nhiệt độ sang cả hai đơn vị, chỉ giữ ngày nóng và ghi nhận lúc pipeline kết thúc:
import { filter, finalize, from, map } from 'rxjs';
type Temperature = {
celsius: number;
fahrenheit: number;
};
const hotDays$ = from([18, 21, 27, 31]).pipe(
map(
(celsius): Temperature => ({
celsius,
fahrenheit: (celsius * 9) / 5 + 32,
}),
),
filter((temperature) => temperature.celsius >= 25),
finalize(() => console.log('pipeline đã được dọn')),
);
const subscription = hotDays$.subscribe({
next: ({ celsius, fahrenheit }) =>
console.log(`${celsius}°C = ${fahrenheit}°F`),
error: (error: unknown) => console.error('Lỗi:', error),
complete: () => console.log('complete'),
});
console.log(`closed: ${subscription.closed}`);Kết quả:
27°C = 80.6°F
31°C = 87.8°F
complete
pipeline đã được dọn
closed: truefrom(array) phát đồng bộ và tự complete sau phần tử cuối. Bởi vậy, khi subscribe() trả về, Subscription đã đóng. Callback của finalize chạy trong teardown sau khi observer nhận complete; không cần gọi unsubscribe() lần nữa trong ví dụ hữu hạn này.
Điều gì thực sự xảy ra khi subscribe?
Pipeline phía trên chỉ là mô tả cho đến khi có subscriber. Có thể đọc execution theo chuỗi sau:
subscribe()
│
▼
from tạo emission ─► map ─► filter ─► observer.next
│
└── hết mảng ───────────────────► observer.complete
│
▼
teardown / finalizeVới Observable mặc định dạng cold, mỗi lần subscribe() thường tạo một execution riêng. Đây là lý do subscribe hai lần có thể chạy producer hai lần, chẳng hạn gửi hai HTTP request. Phân biệt cold/hot và chia sẻ execution là chủ đề riêng; đừng thêm share chỉ để “cho chắc” khi chưa xác định producer nào cần dùng chung.
Ví dụ thực tế: tìm kiếm khi người dùng nhập
Giả sử trang có một ô <input id="search" />. Yêu cầu là bỏ query quá ngắn, chờ người dùng ngừng gõ 300 ms, không gửi lại cùng một query và chỉ hiển thị kết quả của query mới nhất.
import {
Subject,
catchError,
debounceTime,
distinctUntilChanged,
finalize,
fromEvent,
map,
of,
switchMap,
takeUntil,
} from 'rxjs';
import { ajax } from 'rxjs/ajax';
type Product = {
id: string;
name: string;
};
const searchInput = document.querySelector<HTMLInputElement>('#search');
if (!searchInput) {
throw new Error('Không tìm thấy #search');
}
const destroy$ = new Subject<void>();
const results$ = fromEvent<InputEvent>(searchInput, 'input').pipe(
map((event) => (event.target as HTMLInputElement).value.trim()),
debounceTime(300),
distinctUntilChanged(),
switchMap((query) => {
if (query.length < 2) {
return of<Product[]>([]);
}
return ajax
.getJSON<Product[]>(`/api/products?q=${encodeURIComponent(query)}`)
.pipe(
catchError((error: unknown) => {
console.error('Request lỗi:', error);
return of<Product[]>([]);
}),
);
}),
takeUntil(destroy$),
finalize(() => console.log('Đã dọn pipeline tìm kiếm')),
);
results$.subscribe({
next: (products) => console.table(products),
error: (error: unknown) => console.error('Pipeline lỗi:', error),
complete: () => console.log('Pipeline tìm kiếm complete'),
});
// Gọi khi component hoặc màn hình sở hữu ô tìm kiếm bị hủy.
function destroySearch(): void {
destroy$.next();
destroy$.complete();
}Một chuỗi thao tác minh họa có thể là:
Người dùng gõ: r → rx → rxj → rxjs
Request: rxj (bị thay thế) → rxjs
UI nhận: [kết quả rxjs]
Hủy màn hình: complete → finalizeĐọc pipeline từ trên xuống
fromEvent đăng ký DOM listener khi có subscriber, còn map lấy giá trị hiện tại trong input. debounceTime và distinctUntilChanged xử lý hai vấn đề khác nhau: một operator giảm emission theo thời gian, operator còn lại bỏ giá trị liền trước bị lặp.
switchMap trả về mảng rỗng khi query ngắn, nhờ đó UI được xóa và request cũ cũng bị hủy. Với query hợp lệ, nó tạo một inner Observable HTTP. Khi query mới đến, switchMap unsubscribe inner Observable trước và chỉ chuyển tiếp dữ liệu từ inner Observable mới nhất. Với ajax của RxJS 7.8.x, teardown sẽ gọi abort() nếu XMLHttpRequest vẫn chưa hoàn tất, nên ví dụ này hủy cả việc giao kết quả cũ lẫn request đang chạy.
catchError được đặt bên trong switchMap để một request lỗi trở thành mảng rỗng mà stream input vẫn sống. Nếu đặt nó ngoài cùng và không tạo lại stream, lỗi của một request có thể kết thúc toàn bộ pipeline; lần gõ tiếp theo sẽ không còn subscriber để nhận.
switchMap không tự hủy mọi Promise
Nếu thay ajax.getJSON(...) bằng from(fetch(...)), unsubscribe sẽ ngăn kết quả cũ đi tiếp nhưng Promise và fetch không tự bị hủy. Muốn hủy I/O thật sự, nguồn phải có teardown hỗ trợ cancellation, chẳng hạn tích hợp AbortController hoặc dùng API Observable có cơ chế abort.
Lifecycle và teardown của ví dụ
fromEvent không tự complete chỉ vì người dùng ngừng gõ. Vì thế, owner của pipeline phải định nghĩa lúc kết thúc. Khi destroySearch() gọi destroy$.next():
takeUntilcomplete output Observable.- Observer nhận callback
complete. - RxJS unsubscribe chuỗi upstream; DOM listener được gỡ và request
ajaxđang chạy được abort. finalizechạy dù stream kết thúc do complete, error hay explicit unsubscribe.
Chỉ gọi destroy$.complete() mà không gọi next() sẽ không kích hoạt takeUntil; notifier phải phát một giá trị. Bạn cũng có thể giữ Subscription và gọi unsubscribe() trực tiếp. Khác biệt đáng nhớ là explicit unsubscribe dọn tài nguyên nhưng không gọi complete trên observer; finalize vẫn chạy.
Đọc thêm Vòng đời của stream và Subscription và teardown trước khi đưa stream dài hạn vào component.
Chọn reactive hay cách đơn giản hơn
Reactive phù hợp nhất khi bài toán có nhiều giá trị theo thời gian và cần kết hợp các quy tắc như debounce, retry, cancellation, ordering hoặc phối hợp nhiều nguồn. Một event handler đơn giản vẫn là lựa chọn tốt nếu nó chỉ làm một việc và có lifecycle rõ ràng.
| Bài toán | Lựa chọn mặc định | Khi cân nhắc Observable |
|---|---|---|
| Tính một giá trị đồng bộ | Hàm thuần | Khi giá trị đầu vào tự thay đổi theo thời gian |
| Một request rồi trả kết quả | Promise / async–await | Khi cần hủy theo stream, retry, polling hoặc phối hợp nhiều request |
| Một DOM event đơn lẻ | addEventListener | Khi cần debounce, combine hoặc teardown thống nhất với nguồn khác |
| Chuỗi event, WebSocket, timer | Observable | Khi consumer cần composition và lifecycle rõ ràng |
| State có nhiều nguồn cập nhật | Tùy độ phức tạp của ứng dụng | Khi state dẫn xuất từ nhiều stream và cần quy tắc nhất quán |
Mặc định, mình chọn giải pháp ít abstraction nhất vẫn diễn đạt đúng bài toán. Observable đáng giá khi thời gian và sự phối hợp là một phần của yêu cầu, chứ không phải chỉ để thay một Promise bằng cú pháp dài hơn. Bài Khi nào dùng RxJS? đào sâu quyết định này.
Những lỗi tư duy thường gặp
“Reactive” chỉ là viết callback bằng operator
Nếu pipeline vẫn chứa side effect ở mọi map, nhiều subscribe lồng nhau và state mutable nằm rải rác, bạn chỉ đổi cú pháp chứ chưa làm rõ data flow. map nên biến đổi giá trị; side effect có chủ đích nên đặt ở tap hoặc ở observer để người đọc nhận ra ranh giới.
Subscribe bên trong subscribe
Nested subscribe làm cancellation, error propagation và thứ tự hoàn thành bị chia thành nhiều nhánh thủ công. Khi emission này tạo ra một Observable khác, hãy chọn flattening operator theo semantics mong muốn: switchMap, concatMap, mergeMap hoặc exhaustMap. Với tìm kiếm chỉ cần kết quả mới nhất, switchMap là lựa chọn mặc định; xem bài switchMap.
Quên rằng error là terminal
Sau error, Observable không phát thêm next hay complete cho cùng subscription. Nếu một lỗi cục bộ không nên giết stream dài hạn, bắt lỗi tại đúng inner Observable như ví dụ tìm kiếm. Đừng trả về undefined từ catchError; callback phải trả về một ObservableInput. Xem thêm Error channel.
Dùng Subject làm biến toàn cục
Subject hữu ích khi bạn thực sự cần bridge imperative code vào một nguồn multicast. Nhưng expose Subject công khai cho mọi nơi gọi next() sẽ làm mất nguồn sự thật. Mặc định, hãy giữ Subject private và chỉ expose Observable chỉ-đọc; hoặc bắt đầu từ creation function như fromEvent, from và interval nếu producer đã tồn tại.
Tin rằng RxJS tự giải quyết memory leak
RxJS cung cấp teardown, nhưng bạn vẫn phải nối teardown đó với lifecycle của owner. Stream hữu hạn như HTTP thường tự kết thúc; stream vô hạn như DOM event, interval và WebSocket cần chiến lược hủy rõ ràng. Operator không thể biết lúc component của bạn không còn được dùng.
Checklist thiết kế một stream
Trước khi viết operator đầu tiên, bạn có thể dùng checklist ngắn này:
- Xác định owner: ai tạo stream và ai chịu trách nhiệm kết thúc nó?
- Ghi rõ kiểu của từng emission, thay vì chỉ nghĩ “đây là data”.
- Chọn operator theo semantics thời gian: song song, tuần tự, bỏ cũ hay bỏ mới.
- Đặt error boundary gần công việc có thể phục hồi.
- Giữ transformation thuần; cô lập side effect ở ranh giới.
- Với nguồn dài hạn, kiểm tra teardown có thật sự gỡ listener, timer, socket hoặc request không.
- Chỉ share/multicast khi nhiều subscriber phải dùng chung một producer.
Nếu chưa thể trả lời owner và điều kiện kết thúc, pipeline chưa sẵn sàng để đưa vào production.
Học tiếp
Observer pattern
Hiểu producer, Observable, Observer và Subscription.
Vòng đời của stream
Nắm quy tắc của next, error, complete và teardown.
Pipe và operator
Ghép các phép biến đổi thành pipeline dễ đọc.
DOM events
Áp dụng mental model reactive cho sự kiện giao diện.
Nguồn tham khảo
Các mô tả API trong bài đã được đối chiếu với tài liệu và mã nguồn chính thức của RxJS: