Creation functions
Tạo stream từ giá trị, collection, Promise, event và interval.
Dữ liệu trong ứng dụng hiếm khi bắt đầu dưới dạng Observable. Nó có thể là một giá trị mặc định, một mảng, Promise từ fetch, click trên DOM hoặc timer chạy định kỳ. Creation function là lớp nối các nguồn đó vào RxJS để phần còn lại của ứng dụng có thể xử lý chúng bằng cùng một pipeline.
Phạm vi phiên bản
Bài này dùng public API của RxJS 7.8.2. Các ví dụ import trực tiếp từ rxjs; phần tích hợp HTTP chuyên biệt như fromFetch chỉ được nhắc đến khi nói về cancellation.
Mục lục
- Mental model: creation function nối nguồn vào pipeline
- Chọn function theo loại nguồn
- Giá trị có sẵn với of
- Collection và Promise với from
- Event với fromEvent
- Thời gian với interval và timer
- Khởi tạo lazy theo từng subscriber với defer
- Các source đặc biệt
- Ví dụ thực tế: polling tồn kho
- Đồng bộ bất đồng bộ và lifecycle
- Những bẫy thường gặp
- Checklist chọn creation function
- Bài tập tự kiểm tra
- Học tiếp
- Nguồn tham khảo
Mental model: creation function nối nguồn vào pipeline
Creation function nhận một nguồn hoặc tham số mô tả nguồn, rồi trả về Observable<T>. Observable này vẫn tuân theo contract quen thuộc: phát next, sau đó có thể kết thúc bằng complete hoặc error; mỗi lần subscribe() tạo một subscription cần được quản lý vòng đời.
Giá trị / collection / Promise / event / clock
│
▼
creation function
of · from · fromEvent · timer · ...
│
▼
Observable<T>
│
pipe(map, filter, ...)
│
▼
subscribe()Hãy xem creation function như cửa nhập hàng của một kho: dữ liệu đến từ nhiều phương tiện khác nhau, nhưng qua cửa này đều được đưa vào cùng một quy trình xử lý. Ẩn dụ dừng ở chuyện “chuẩn hóa đầu vào”; nó không có nghĩa mọi nguồn sẽ giống nhau về thời điểm chạy, cancellation hay việc có tự complete hay không.
Creation function khác pipeable operator
Creation function đứng ở đầu pipeline và tạo Observable:
import { from, map } from 'rxjs';
const prices$ = from([100, 250, 400]).pipe(
map((price) => price * 1.1),
);Trong ví dụ này, from() là creation function còn map() là pipeable operator. from() không cần một source Observable đi trước; map() thì có. Một số function như combineLatest, concat hay merge cũng có dạng creation function, nhưng bài này tập trung vào việc đưa nguồn dữ liệu ban đầu vào RxJS. Phần Creation operators sẽ mở rộng sang nhóm tạo và phối hợp nhiều stream.
Chọn function theo loại nguồn
| Bạn đang có gì? | Chọn mặc định | Observable sẽ làm gì? | Tự complete? |
|---|---|---|---|
| Một hoặc nhiều giá trị đã có sẵn | of(value1, value2) | Phát nguyên từng argument | Có |
| Array, array-like, iterable hoặc string | from(input) | Phát lần lượt từng phần tử | Có |
| Promise | from(promise) | Phát một giá trị hoặc báo lỗi theo Promise | Có, sau một giá trị thành công |
DOM event, Node EventEmitter, target kiểu jQuery | fromEvent(target, name) | Phát mỗi khi handler được gọi | Thường không |
| Tick đều theo chu kỳ | interval(period) | Phát 0, 1, 2, ... sau mỗi chu kỳ | Không |
| Một lần sau delay, hoặc tick với delay ban đầu riêng | timer(due, period?) | Phát 0 sau delay; có thể tiếp tục tăng | Có nếu không truyền period |
| Cần tạo nguồn mới tại thời điểm subscribe | defer(factory) | Gọi factory riêng cho mỗi subscriber | Theo source mà factory trả về |
Mặc định, hãy chọn theo hình dạng và lifecycle của nguồn, không chọn theo tên nghe quen. Hai function dễ nhầm nhất là of() và from(): một bên giữ nguyên argument, bên kia duyệt input.
Giá trị có sẵn với of
of(...values) phát từng argument theo đúng thứ tự rồi complete ngay. Nếu không truyền scheduler, các emission xảy ra đồng bộ trong lời gọi subscribe().
import { of } from 'rxjs';
const status$ = of('queued', 'processing', 'done');
console.log('trước subscribe');
status$.subscribe({
next: (status) => console.log('next:', status),
complete: () => console.log('complete'),
});
console.log('sau subscribe');Kết quả:
trước subscribe
next: queued
next: processing
next: done
complete
sau subscribeof() hợp với giá trị mặc định, fixture trong test, fallback nhỏ hoặc nhánh cần trả một Observable từ dữ liệu đã có. Nếu dữ liệu vốn là collection và bạn muốn xử lý từng phần tử, from() diễn đạt ý định rõ hơn.
of không flatten collection
Đây là khác biệt cần nhớ:
import { from, of } from 'rxjs';
of([10, 20, 30]).subscribe({
next: (value) => console.log('of:', value),
});
from([10, 20, 30]).subscribe({
next: (value) => console.log('from:', value),
});Kết quả:
of: [10, 20, 30]
from: 10
from: 20
from: 30of([10, 20, 30]) có một emission là cả array. from([10, 20, 30]) có ba emission, mỗi emission là một phần tử. Không có lựa chọn nào tốt hơn tuyệt đối: chọn theo đơn vị dữ liệu mà downstream cần nhận.
Collection và Promise với from
from(input) chuyển một ObservableInput thành Observable. Trong RxJS 7.8.2, nhóm input này gồm Observable, Observable-like, array-like, Promise, iterable, async iterable và readable stream tương thích.
Array iterable và string
Với array hoặc iterable đồng bộ, from() duyệt từng phần tử theo thứ tự rồi complete:
import { from } from 'rxjs';
const roles = new Set(['reader', 'editor', 'admin']);
from(roles).subscribe({
next: (role) => console.log(role),
complete: () => console.log('đã duyệt xong'),
});Kết quả:
reader
editor
admin
đã duyệt xongString cũng là input hợp lệ và được phát thành từng ký tự:
import { from } from 'rxjs';
from('RxJS').subscribe({
next: (character) => console.log(character),
});Đừng truyền một scalar như from(42): number không phải ObservableInput, nên đây không phải cách tạo Observable từ một giá trị đơn. Dùng of(42).
Generator có thể bị dùng hết
Nếu bạn truyền cùng một generator object vào from() rồi subscribe nhiều lần, subscription đầu có thể đã tiêu thụ iterator. Khi mỗi subscriber cần một generator mới, dùng defer(() => from(createGenerator())) thay vì tạo generator một lần ở ngoài.
Promise phát một kết quả rồi complete
Với Promise resolve, from() phát giá trị đã resolve rồi complete. Nếu Promise reject, Observable gửi error và không gửi complete.
import { from } from 'rxjs';
const greeting$ = from(Promise.resolve('xin chào'));
console.log('trước subscribe');
greeting$.subscribe({
next: (value) => console.log('next:', value),
error: (error: unknown) => console.error('error:', error),
complete: () => console.log('complete'),
});
console.log('sau subscribe');Kết quả:
trước subscribe
sau subscribe
next: xin chào
completeKhác from(array), notification từ Promise xuất hiện bất đồng bộ qua Promise job queue. from() chuẩn hóa cách consumer nhận kết quả; nó không biến Promise thành công việc đồng bộ.
Promise vẫn eager và không tự bị hủy
Một hiểu nhầm phổ biến là from(fetch(...)) khiến request chỉ bắt đầu khi subscribe. Thực tế, fetch() đã được gọi để tạo Promise trước khi from() nhận nó:
import { from } from 'rxjs';
const responsePromise = fetch('/api/profile'); // Request bắt đầu ở đây.
const response$ = from(responsePromise);
// subscribe chỉ nối Observer với Promise đã tồn tại.
const subscription = response$.subscribe({
next: (response) => console.log(response.status),
error: (error: unknown) => console.error(error),
});Nếu unsubscribe trước khi Promise settle, RxJS không chuyển kết quả đến subscriber đã đóng. Tuy nhiên, Promise và công việc nền của nó vẫn tiếp tục vì Promise không có giao thức cancellation chuẩn.
Khi cần trì hoãn việc tạo Promise đến lúc subscribe, bọc nó bằng defer(). Khi cần hủy request mạng thật sự, dùng API hỗ trợ cancellation như AbortController hoặc fromFetch; chỉ đổi Promise thành Observable là chưa đủ.
Event với fromEvent
fromEvent(target, eventName) phù hợp khi API có cặp phương thức đăng ký và gỡ handler quen thuộc:
- DOM
addEventListener/removeEventListener; - Node-style
addListener/removeListener; - jQuery-style
on/off; NodeListhoặcHTMLCollectionchứa các DOM target.
Ví dụ biến nội dung ô tìm kiếm thành stream:
import { distinctUntilChanged, fromEvent, map } from 'rxjs';
const searchInput = document.querySelector<HTMLInputElement>('#search');
if (!searchInput) {
throw new Error('Không tìm thấy #search');
}
const query$ = fromEvent<InputEvent>(searchInput, 'input').pipe(
map(() => searchInput.value.trim()),
distinctUntilChanged(),
);
const subscription = query$.subscribe({
next: (query) => console.log('Tìm:', query),
});
export function destroySearch(): void {
subscription.unsubscribe();
}DOM event không tự complete chỉ vì người dùng ngừng tương tác. Lifecycle của màn hình hoặc component phải quyết định lúc hủy subscription.
Listener sống theo subscription
Mỗi lần subscribe vào Observable do fromEvent() trả về, RxJS đăng ký một handler. Khi subscription bị hủy, RxJS gỡ đúng handler đó.
subscribe A ──► add handler A ──► unsubscribe A ──► remove handler A
subscribe B ──► add handler B ──► unsubscribe B ──► remove handler BVì vậy, subscribe ba lần có thể tạo ba listener trên cùng target. Nếu nhiều consumer cần dùng chung một listener, đó là bài toán multicasting; hãy quyết định rõ bằng share() thay vì vô tình tạo subscription trùng lặp.
Event producer thường đã tồn tại bên ngoài RxJS
fromEvent() quản lý việc gắn và gỡ listener, nhưng nó không sở hữu browser, DOM hay EventEmitter. Nguồn event có thể tiếp tục hoạt động khi không có subscriber. Đây là lý do không nên đồng nhất “mỗi subscription có handler riêng” với “mỗi subscription tạo toàn bộ producer riêng”.
Khi event API không khớp fromEvent
Nếu thư viện dùng tên hàm hoặc chữ ký riêng, chẳng hạn register(handler) và dispose(handler), dùng fromEventPattern(addHandler, removeHandler). Nếu callback chỉ được gọi đúng một lần để trả kết quả, bindCallback() hoặc API Promise có thể hợp lý hơn.
Nguyên tắc là phải tìm được cả đường đăng ký lẫn gỡ đăng ký. Nếu chỉ bọc callback add mà không có teardown tương ứng, Observable có thể ngừng chuyển emission sau unsubscribe nhưng listener bên ngoài vẫn bị giữ lại.
Thời gian với interval và timer
Cả interval() và timer() đều tạo Observable số tăng dần dựa trên clock của scheduler. Chúng giống nhau ở output nhưng khác ở cách điều khiển emission đầu tiên.
interval chờ đủ một chu kỳ
interval(period) phát 0, 1, 2, ... vô hạn. Emission đầu tiên chỉ đến sau khi đã chờ đủ một period, không đến ngay lúc subscribe.
import { interval, take } from 'rxjs';
interval(1_000)
.pipe(take(3))
.subscribe({
next: (tick) => console.log('tick:', tick),
complete: () => console.log('complete'),
});Kết quả theo thời gian:
~1 giây: tick: 0
~2 giây: tick: 1
~3 giây: tick: 2
~3 giây: completeinterval() tự nó không complete. Ví dụ dùng take(3) để giới hạn vòng đời; trong ứng dụng, lifecycle UI hoặc điều kiện nghiệp vụ cũng có thể chịu trách nhiệm dừng stream.
timer điều khiển emission đầu tiên
timer(due) đợi due mili giây, phát 0 một lần rồi complete. timer(due, period) đợi theo due, sau đó tiếp tục phát số tăng dần theo period.
import { take, timer } from 'rxjs';
// Phát 0 ngay khi scheduler có thể chạy, rồi lặp mỗi giây.
timer(0, 1_000)
.pipe(take(3))
.subscribe({
next: (tick) => console.log('poll:', tick),
complete: () => console.log('dừng polling'),
});
// Phát 0 một lần sau khoảng 1,5 giây rồi complete.
timer(1_500).subscribe({
next: () => console.log('hết thời gian chờ'),
});Nếu cần polling chạy ngay rồi lặp, mình chọn timer(0, period). Nếu tick đầu tiên cũng phải chờ một chu kỳ, interval(period) nói rõ ý định hơn. Các mốc thời gian là tối thiểu mong đợi, không phải deadline chính xác tuyệt đối, vì event loop và scheduler có thể làm callback chạy muộn.
Khởi tạo lazy theo từng subscriber với defer
defer(factory) chưa gọi factory khi bạn khai báo Observable. Mỗi lần có subscriber, nó mới gọi factory, lấy ObservableInput được trả về rồi subscribe vào nguồn đó. Nhờ vậy, mỗi subscriber có thể đọc state mới nhất hoặc tạo Promise mới.
import { defer, from } from 'rxjs';
type Profile = {
id: string;
displayName: string;
};
function requestProfile(): Promise<Profile> {
return fetch('/api/profile').then(async (response) => {
if (!response.ok) {
throw new Error(`HTTP ${response.status}`);
}
return (await response.json()) as Profile;
});
}
const profile$ = defer(() => from(requestProfile()));
// Chưa có request nào ở dòng khai báo trên.
const firstSubscription = profile$.subscribe({
next: (profile) => console.log('first:', profile.displayName),
error: (error: unknown) => console.error('first error:', error),
});
const secondSubscription = profile$.subscribe({
next: (profile) => console.log('second:', profile.displayName),
error: (error: unknown) => console.error('second error:', error),
});Ví dụ này tạo hai request vì factory chạy hai lần. Đó là đúng khi mỗi subscriber cần execution riêng, nhưng sai nếu bạn định chia sẻ một request. Chia sẻ và cache cần quyết định riêng bằng multicasting operator phù hợp.
Nếu factory ném lỗi đồng bộ, RxJS chuyển lỗi đó vào error channel của subscription. Tuy nhiên, defer() chỉ quyết định khi nào tạo nguồn; nó không tự thêm cancellation cho Promise mà factory trả về.
defer không phải cơ chế cache
defer() tạo lại source cho mỗi subscription. Nếu mục tiêu là tránh request trùng lặp hoặc phát lại kết quả gần nhất, hãy thiết kế sharing/cache riêng và xác định rõ thời điểm reset.
Các source đặc biệt
Ngoài nhóm thường gặp phía trên, RxJS có vài source nhỏ nhưng hữu ích khi ghép nhánh điều kiện:
| Source | Hành vi | Trường hợp dùng |
|---|---|---|
EMPTY | Không phát giá trị, complete ngay | Bỏ qua một nhánh nhưng vẫn kết thúc sạch |
NEVER | Không phát, không error, không complete | Test hoặc giữ một nhánh im lặng có chủ đích |
throwError(() => error) | Gửi error ngay khi subscribe | Trả một Observable lỗi từ nhánh cần ObservableInput |
range(start, count) | Phát count số nguyên liên tiếp rồi complete | Dữ liệu số hữu hạn, demo hoặc test |
import { EMPTY, NEVER, range, throwError } from 'rxjs';
const skipped$ = EMPTY;
const waitingForever$ = NEVER;
const pageNumbers$ = range(1, 3); // 1, 2, 3 rồi complete
const invalid$ = throwError(() => new Error('Dữ liệu không hợp lệ'));NEVER không tự giải phóng subscription bằng complete, nên chỉ dùng khi bạn đã biết ai sẽ hủy. Với throwError, dùng error factory như trên để tạo error tại thời điểm subscribe; truyền thẳng error instance là overload đã deprecated trong RxJS 7.
Ví dụ thực tế: polling tồn kho
Giả sử dashboard cần tải tồn kho ngay khi mở, sau đó refresh mỗi 30 giây. Ta dùng timer(0, 30_000) làm clock, defer() để tạo request tại từng tick và from() để nối Promise vào pipeline.
import {
catchError,
defer,
from,
map,
of,
switchMap,
timer,
type Observable,
} from 'rxjs';
type StockItem = {
sku: string;
available: number;
};
type StockState =
| { kind: 'ready'; items: StockItem[] }
| { kind: 'error'; message: string };
function loadStock(): Promise<StockItem[]> {
return fetch('/api/stock').then(async (response) => {
if (!response.ok) {
throw new Error(`HTTP ${response.status}`);
}
return (await response.json()) as StockItem[];
});
}
const stock$: Observable<StockState> = timer(0, 30_000).pipe(
switchMap(() =>
defer(() => from(loadStock())).pipe(
map((items): StockState => ({ kind: 'ready', items })),
catchError((error: unknown) =>
of<StockState>({
kind: 'error',
message: error instanceof Error ? error.message : String(error),
}),
),
),
),
);
const subscription = stock$.subscribe({
next: (state) => {
if (state.kind === 'ready') {
console.table(state.items);
} else {
console.error('Không tải được tồn kho:', state.message);
}
},
});
export function stopStockPolling(): void {
subscription.unsubscribe();
}Dòng chảy của ví dụ:
timer(0, 30_000)phát tick đầu gần như ngay lập tức, rồi phát lại mỗi 30 giây.- Mỗi tick khiến
switchMap()subscribe vào một source mới. defer()gọiloadStock()tại thời điểm đó, nên request không chạy trước subscription.from()chuyển kết quả hoặc lỗi của Promise sang Observable.catchError()nằm trong inner stream, nên một request lỗi trở thànhStockStatevà không giết clock polling bên ngoài.stopStockPolling()hủy timer và ngăn subscriber nhận thêm kết quả.
Có một ranh giới quan trọng: nếu request cũ còn chạy khi tick mới đến, switchMap() sẽ unsubscribe inner Observable cũ nhưng Promise nền vẫn có thể tiếp tục. Kết quả cũ không được chuyển xuống downstream, song network request chỉ bị hủy thật sự nếu source tích hợp cancellation, chẳng hạn fromFetch hoặc custom Observable dùng AbortController.
Đồng bộ bất đồng bộ và lifecycle
Creation function không quyết định chung rằng “RxJS là async”. Chính loại nguồn và scheduler mới quyết định thời điểm emission.
| Source | Thời điểm mặc định | Kết thúc mặc định | Cleanup đáng chú ý |
|---|---|---|---|
of(...) | Đồng bộ | Complete sau argument cuối | Thường không có resource ngoài |
from(array/iterable) | Đồng bộ | Complete khi duyệt xong | Dừng duyệt nếu subscriber đã đóng |
from(promise) | Bất đồng bộ theo Promise | Complete sau resolve, error sau reject | Unsubscribe không hủy Promise |
fromEvent(...) | Khi target phát event | Thường không complete | Unsubscribe gỡ listener |
interval(...) | Sau mỗi period | Không complete | Unsubscribe hủy lịch tick |
timer(due) | Sau delay | Complete sau một emission | Unsubscribe hủy lịch chờ |
defer(factory) | Theo source factory trả về | Theo source đó | Tạo source mới cho mỗi subscription |
Vì vậy, đừng thêm setTimeout chỉ để “làm Observable async”, và cũng đừng giả định code sau subscribe() luôn chạy trước callback next. Với of() hoặc from(array), callback có thể chạy xong trước khi subscribe() trả về.
Những bẫy thường gặp
- Dùng
of(array)khi muốn từng phần tử. Downstream nhận một array duy nhất. Đổi sangfrom(array)hoặc giữof(array)nếu array thật sự là một domain value. - Dùng
from(value)cho scalar.from(42)không hợp lệ;of(42)mới đúng. - Tạo Promise trước
defer.const promise = fetch(...); defer(() => from(promise))vẫn dùng Promise đã chạy. Phải tạo Promise bên trong factory:defer(() => from(fetch(...))). - Tin rằng unsubscribe hủy Promise. RxJS ngừng gửi kết quả đến subscriber, nhưng công việc nền cần cơ chế cancellation riêng.
- Chờ
completetừ event hoặc interval. Hai nguồn này thường sống mãi. Dùng lifecycle,take(),takeUntil()hoặc điều kiện kết thúc phù hợp. - Subscribe nhiều lần vào
fromEvent. Mỗi subscription thêm một handler. Nếu cần chia sẻ, dùng multicasting có chủ đích. - Nhầm
interval(1000)phát ngay. Tick0đến sau khoảng một giây. Dùngtimer(0, 1000)nếu cần emission đầu tiên ngay. - Dùng
defernhư cache. Nó làm điều ngược lại: tạo source mới theo từng subscriber. - Bỏ qua đường error của Promise. Promise reject trở thành error terminal. Đặt
catchError()ở đúng cấp để quyết định fallback có kết thúc cả pipeline hay không. - Truyền scheduler vào overload cũ mà không kiểm tra deprecation. Trong RxJS 7, scheduler argument của
of()vàfrom()đã deprecated cho RxJS 8. Với code mới, ưu tiênscheduled()hoặc operator scheduling phù hợp.
Checklist chọn creation function
Trước khi viết code, trả lời năm câu hỏi:
- Downstream cần nhận cả collection hay từng phần tử?
- Công việc đã bắt đầu rồi, hay chỉ nên bắt đầu khi subscribe?
- Source tự complete hay cần consumer hủy?
- Unsubscribe có thực sự dừng resource bên ngoài không?
- Mỗi subscriber cần producer riêng hay phải chia sẻ producer?
Nếu chỉ nhớ một quy tắc: chọn creation function theo lifecycle, không chỉ theo kiểu TypeScript của dữ liệu. Hai nguồn cùng trả Observable<Response> vẫn có thể khác hoàn toàn: một nguồn tạo request mới mỗi lần subscribe, nguồn kia chỉ quan sát Promise đã chạy từ trước.
Bài tập tự kiểm tra
- Chạy
of([1, 2, 3])vàfrom([1, 2, 3]), rồi ghi lại số lần callbacknextđược gọi. - Đặt
console.logtrước và sausubscribe()choof('A')vàfrom(Promise.resolve('A')); giải thích thứ tự khác nhau. - Tạo
fromEvent(button, 'click'), subscribe hai lần, sau đó hủy từng subscription để kiểm tra handler nào còn hoạt động. - Đổi
interval(1_000)thànhtimer(0, 1_000)và quan sát thời điểm tick đầu tiên. - Bọc một hàm tăng counter bằng
defer()rồi subscribe hai lần; xác nhận factory chạy riêng cho từng subscriber.
Sau các bài này, hãy quay lại một pipeline thật trong dự án và đánh dấu ba điểm: source bắt đầu ở đâu, ai sở hữu subscription, và resource nào thực sự được dọn khi unsubscribe.
Học tiếp
Subscribe và Observer
Theo dõi notification và quản lý execution do creation function tạo ra.
Cold và hot Observable
Phân biệt producer riêng theo subscription với producer được chia sẻ.
Subscription và teardown
Dừng timer, gỡ listener và giải phóng resource đúng lúc.
Creation operators
Mở rộng sang các function tạo và phối hợp nhiều stream.
Nguồn tham khảo
- RxJS 7.8.2 — Observable guide
- RxJS 7.8.2 — mã nguồn và API
of - RxJS 7.8.2 — mã nguồn và API
from - RxJS 7.8.2 — chuyển đổi Array, Promise, iterable và async iterable
- RxJS 7.8.2 — mã nguồn và API
fromEvent - RxJS 7.8.2 — mã nguồn và API
interval - RxJS 7.8.2 — mã nguồn và API
timer - RxJS 7.8.2 — mã nguồn và API
defer - RxJS 7.8.2 — mã nguồn và API
throwError - RxJS 7.8.2 — cancellation trong
fromFetch