Observer pattern
Mối quan hệ giữa producer, Observable, Observer và Subscription.
Một ô tìm kiếm có thể phát ra hàng chục sự kiện khi người dùng gõ, còn component chỉ muốn phản ứng với những giá trị mình quan tâm và dừng nghe khi bị huỷ. Nếu nối trực tiếp mọi callback với mọi nguồn sự kiện, phần khởi động, xử lý lỗi và dọn tài nguyên rất nhanh bị rải khắp code. Observer pattern tách bên phát dữ liệu khỏi bên nhận; RxJS bổ sung Observable và Subscription để mối liên hệ đó có vòng đời rõ ràng và ghép nối được.
Phạm vi phiên bản
Các ví dụ trong bài dùng public API ổn định của RxJS 7.x và TypeScript. Bài tập trung vào mental model; các biến thể dự kiến cho RxJS 8 không được dùng ở đây.
Mục lục
- Mental model: đăng ký một phiên làm việc
- Bốn vai trò cốt lõi
- Điều gì xảy ra khi subscribe
- Teardown: dọn tài nguyên đúng chỗ
- Ví dụ thực tế: lắng nghe ô tìm kiếm
- Chọn API nào
- Những lỗi tư duy thường gặp
- Bước tiếp theo
- Tài liệu tham khảo
Mental model: đăng ký một phiên làm việc
Hình dung một kênh thông báo ở quán ăn: bếp là producer, quy ước chuyển phiếu là Observable, nhân viên phục vụ nhận phiếu là Observer, còn tờ đăng ký ca trực là Subscription. Khi ca kết thúc, tờ đăng ký bị huỷ và tài nguyên liên quan phải được dọn.
Ẩn dụ này chỉ giúp nhớ vai trò. Trong RxJS, một Observable thường có thể tạo producer riêng cho từng lần subscribe, trong khi bếp ngoài đời thường được nhiều người dùng chung. Việc producer riêng hay dùng chung dẫn tới khái niệm cold, hot, unicast và multicast ở phần sau của lộ trình.
next(value)
┌──────────┐ được bao bởi ┌────────────┐ ─────────────► ┌──────────┐
│ Producer │ ───────────────► │ Observable │ error(err) │ Observer │
└──────────┘ └─────┬──────┘ ─────────────► └──────────┘
▲ │ complete() │
│ │ │
└──── teardown ◄──────────────┴──── Subscription ◄─────────┘
unsubscribe()Điểm quan trọng nhất: Observable không phải dữ liệu đã chạy sẵn. Nó là một template mô tả cách nối consumer với producer. Với Observable thông thường, subscribe() mới thiết lập kết nối và bắt đầu execution.
Bốn vai trò cốt lõi
| Vai trò | Câu hỏi nó trả lời | Ví dụ |
|---|---|---|
| Producer | Giá trị thật sự đến từ đâu và xuất hiện khi nào? | Timer, DOM event, WebSocket, mảng dữ liệu hoặc callback của SDK |
| Observable | Kết nối producer với consumer theo quy tắc nào? | fromEvent(input, 'input'), interval(1000), hoặc new Observable(...) |
| Observer | Consumer phản ứng ra sao với từng loại notification? | Object có các handler next, error, complete |
| Subscription | Execution đang chạy được quản lý và huỷ bằng cách nào? | Giá trị trả về từ observable.subscribe(observer) |
Producer và Observable không nhất thiết là một. Với fromEvent, DOM là producer; Observable là lớp bao cung cấp contract và khả năng kết hợp operator. Với new Observable, code tạo timer nằm trong hàm khởi tạo execution nên ranh giới giữa hai khái niệm trông gần nhau hơn, nhưng vẫn nên hỏi riêng: “Ai sinh dữ liệu?” và “Ai quản lý cách đăng ký?”.
Observer là consumer được biểu diễn bằng object. Một observer đầy đủ có ba handler:
import type { Observer } from 'rxjs';
const observer: Observer<number> = {
next: (value) => console.log('giá trị:', value),
error: (error: unknown) => console.error('lỗi:', error),
complete: () => console.log('hoàn tất'),
};RxJS cũng chấp nhận partial observer, chẳng hạn chỉ có next. Tuy vậy, ở boundary của ứng dụng như gọi API hoặc kết nối bên ngoài, mình thường khai báo error rõ ràng để lỗi không trở thành unhandled error khó truy vết.
Observer không điều khiển nhịp phát
Observable dùng mô hình push: producer quyết định lúc gửi giá trị, observer chỉ phản ứng khi nhận được. observer.next cũng không phải lệnh để consumer “xin” giá trị kế tiếp như iterator.next().
Điều gì xảy ra khi subscribe
Mỗi lời gọi subscribe tạo một subscription. Theo thứ tự khái niệm, RxJS sẽ:
- nhận observer từ consumer;
- thiết lập Observable execution và kết nối tới producer;
- chuyển các notification hợp lệ tới observer;
- trả về một
Subscriptionđể consumer có thể huỷ execution; - chạy teardown khi execution kết thúc hoặc bị huỷ.
Ví dụ đồng bộ tối thiểu
Ví dụ sau cố ý dùng new Observable để lộ toàn bộ vòng đời. Trong code ứng dụng, creation function có sẵn thường là lựa chọn gọn và an toàn hơn.
import { Observable, type Observer } from 'rxjs';
const numbers$ = new Observable<number>((subscriber) => {
console.log('producer: bắt đầu');
subscriber.next(10);
subscriber.next(20);
subscriber.complete();
return () => console.log('producer: teardown');
});
const observer: Observer<number> = {
next: (value) => console.log('observer: next', value),
error: (error: unknown) => console.error('observer: error', error),
complete: () => console.log('observer: complete'),
};
console.log('trước subscribe');
const subscription = numbers$.subscribe(observer);
console.log('sau subscribe, closed =', subscription.closed);Kết quả:
trước subscribe
producer: bắt đầu
observer: next 10
observer: next 20
observer: complete
producer: teardown
sau subscribe, closed = truenext xuất hiện trước dòng “sau subscribe”, vì Observable có thể chạy đồng bộ. RxJS không tự biến mọi việc thành async. Sau complete, subscription đóng và teardown chạy; vì thế subscription.closed là true.
Đừng gọi biến được truyền vào constructor là observer của ứng dụng. Đó là Subscriber do RxJS cung cấp cho phía producer: nó chuyển tiếp notification tới observer và bảo vệ Observable contract. Consumer vẫn chỉ truyền observer vào subscribe.
Ba notification và một hành động huỷ
Một execution tuân theo grammar ngắn gọn:
next* (error | complete)?Nói cách khác, nó có thể phát không, một hoặc nhiều next; sau đó có nhiều nhất một terminal notification là error hoặc complete. Khi một terminal notification đã xảy ra, observer không nhận thêm giá trị.
next(value)mang dữ liệu và có thể xảy ra nhiều lần.error(error)báo execution thất bại rồi đóng subscription.complete()báo producer kết thúc bình thường rồi đóng subscription.unsubscribe()là hành động của consumer để ngừng quan sát; nó không phải notification.
Điểm cuối thường gây nhầm: gọi unsubscribe() không gọi handler complete. complete là tín hiệu từ producer, còn unsubscribe là consumer nói “tôi không cần nữa”. Cả hai đều dẫn tới finalization, nhưng ý nghĩa nghiệp vụ khác nhau.
Để đi sâu vào quy tắc terminal notification, đọc Vòng đời của stream.
Teardown: dọn tài nguyên đúng chỗ
Producer có thể giữ timer, event listener, socket hoặc request đang chạy. Khi tạo custom Observable, hàm subscribe nên trả về teardown để giải phóng tài nguyên đó. RxJS chạy teardown khi subscription complete, error hoặc bị unsubscribe.
import { Observable } from 'rxjs';
const heartbeat$ = new Observable<number>((subscriber) => {
let beat = 0;
const timerId = setInterval(() => {
beat += 1;
subscriber.next(beat);
}, 1_000);
return () => {
clearInterval(timerId);
console.log('teardown: đã đóng timer');
};
});
const subscription = heartbeat$.subscribe({
next: (beat) => console.log('heartbeat', beat),
complete: () => console.log('hoàn tất'),
});
setTimeout(() => subscription.unsubscribe(), 2_500);Kết quả xấp xỉ:
heartbeat 1
heartbeat 2
teardown: đã đóng timerDòng hoàn tất không xuất hiện vì consumer đã unsubscribe chứ producer không gọi complete(). Nếu bỏ clearInterval, observer đã đóng sẽ không nhận thêm giá trị, nhưng timer vẫn đánh thức event loop và lãng phí tài nguyên. Đây là lý do “không còn thấy log” chưa chứng minh việc cleanup đã đúng.
Quy tắc ownership
Code gọi subscribe() sở hữu Subscription và phải biết khi nào không còn cần nó. Riêng code tạo producer sở hữu teardown và phải biết cách giải phóng từng tài nguyên đã mở.
Creation function và operator của RxJS thường đã cài teardown đúng. Chỉ dùng new Observable khi bạn thực sự cần bọc một API callback hoặc resource chưa có adapter phù hợp. Xem cách tổ chức cleanup phức tạp hơn tại Subscription và teardown và Tự tạo Observable.
Ví dụ thực tế: lắng nghe ô tìm kiếm
Giả sử trang có một ô nhập với id="search". Bạn muốn chỉ xử lý khi người dùng ngừng gõ 300 ms và bỏ qua query trùng liên tiếp:
import {
debounceTime,
distinctUntilChanged,
fromEvent,
map,
} from 'rxjs';
const input = document.querySelector<HTMLInputElement>('#search');
if (!input) {
throw new Error('Không tìm thấy #search');
}
const query$ = fromEvent<InputEvent>(input, 'input').pipe(
map((event) =>
(event.currentTarget as HTMLInputElement).value.trim(),
),
debounceTime(300),
distinctUntilChanged(),
);
const subscription = query$.subscribe({
next: (query) => console.log('tìm kiếm:', query),
error: (error: unknown) => console.error('ô tìm kiếm lỗi:', error),
});
export function destroySearchBox(): void {
subscription.unsubscribe();
}Nếu người dùng gõ nhanh r → rx → rxjs, rồi dừng hơn 300 ms, output là:
tìm kiếm: rxjsỞ đây:
- DOM
inputevent là producer; query$là Observable mô tả pipeline;- object truyền vào
subscribelà observer; subscriptionquản lý lần lắng nghe cụ thể;fromEventđăng ký event listener khi subscribe và gỡ listener trong teardown khidestroySearchBox()chạy.
DOM event không tự complete khi element biến mất. Vì vậy lifecycle của UI phải gọi destroySearchBox; chỉ xoá element khỏi DOM không thay thế cho unsubscribe. Khi nối query với HTTP, nên dùng flattening operator thay vì subscribe lồng nhau; xem Nested subscribe.
Chọn API nào
| Tình huống | Lựa chọn mặc định | Lý do |
|---|---|---|
| Có array, Promise, iterable hoặc Observable-compatible input | from(...) | Adapter có sẵn, semantics và teardown đã được thư viện xử lý |
| Có một nhóm giá trị tĩnh | of(...) | Thể hiện rõ emission hữu hạn rồi complete |
| Có DOM/EventTarget | fromEvent(...) | Tự thêm và gỡ listener theo subscription |
| Có callback API hoặc resource riêng chưa được hỗ trợ | new Observable(...) | Cho phép định nghĩa setup, notification và teardown chính xác |
Cần đủ next, error, complete | subscribe({ ... }) | Tên handler rõ, tránh nhầm vị trí callback |
| Chỉ cần xử lý giá trị đơn giản | subscribe(value => ...) | Ngắn gọn và vẫn hợp lệ trong RxJS 7.x |
Trong RxJS 7.x, các overload subscribe(next, error, complete) nhiều đối số đã bị deprecated. Đừng viết callback rời theo vị trí; observer object dễ đọc hơn và tránh phải truyền null làm chỗ trống.
Một câu hỏi khác là khi nào cần lưu Subscription. Với source hữu hạn, đồng bộ như of(1, 2, 3), execution đã đóng ngay sau subscribe. Với timer, DOM event, WebSocket hoặc source sống lâu, bạn cần chiến lược huỷ rõ ràng: giữ Subscription, gắn nó vào lifecycle của framework, hoặc dùng operator kết thúc thích hợp. Không nên gọi unsubscribe() ngay sau subscribe() chỉ như một nghi thức; làm vậy có thể huỷ công việc async trước khi nó phát giá trị.
Những lỗi tư duy thường gặp
Nghĩ rằng khai báo Observable là đã chạy
Observable thường lazy. Tạo const data$ = new Observable(...) chỉ tạo template; side effect trong hàm subscribe chưa chạy cho tới khi có consumer đăng ký.
Nghĩ rằng mọi Observable đều async
Ví dụ of(1, 2, 3) phát và complete đồng bộ theo mặc định. Nếu thứ tự thực thi quan trọng, đừng đoán dựa vào kiểu Observable; hãy xem producer và scheduler thực tế.
Subscribe nhiều lần nhưng mong một execution duy nhất
Với cold Observable, mỗi subscribe thường tạo producer riêng. Hai subscription vào một HTTP Observable có thể gây hai request. Nếu cần chia sẻ một producer, hãy học rõ trade-off của multicast thay vì thêm cờ thủ công; bắt đầu tại Cold và hot Observable và Unicast và multicast.
Quên teardown trong custom Observable
Chặn notification sau unsubscribe chưa đủ; timer, listener và socket vẫn phải được đóng. Mỗi lệnh setup giữ resource nên có hành động cleanup đối xứng trong teardown.
Dùng complete như cleanup duy nhất
Source sống lâu có thể không complete, còn unsubscribe không gọi callback complete. Đặt cleanup resource ở teardown hoặc dùng operator finalize cho side effect cần chạy khi subscription kết thúc theo bất kỳ đường nào.
Subscribe bên trong subscribe
Nested subscribe tách một flow thành nhiều subscription có ownership và error handling rời nhau. Trừ khi bạn chủ ý quản lý từng execution độc lập, hãy ưu tiên operator kết hợp hoặc flattening để teardown lan truyền qua một subscription chính.
Bước tiếp theo
Bạn có thể tự kiểm tra mental model bằng cách lấy ví dụ ô tìm kiếm và trả lời bốn câu: producer là gì, Observable được tạo ở đâu, observer xử lý notification nào, và ai chịu trách nhiệm unsubscribe. Nếu một câu chưa rõ, lifecycle trong code cũng thường chưa rõ.
Subscribe và Observer
Thực hành các dạng observer và cách đăng ký.
Vòng đời của stream
Hiểu sâu next, error, complete và trạng thái đóng.
Subscription và teardown
Quản lý cleanup và ownership trong code thực tế.