Học RxJS
Observable & Subscription

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

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 địnhObservable sẽ làm gì?Tự complete?
Một hoặc nhiều giá trị đã có sẵnof(value1, value2)Phát nguyên từng argumentCó
Array, array-like, iterable hoặc stringfrom(input)Phát lần lượt từng phần tửCó
Promisefrom(promise)Phát một giá trị hoặc báo lỗi theo PromiseCó, sau một giá trị thành công
DOM event, Node EventEmitter, target kiểu jQueryfromEvent(target, name)Phát mỗi khi handler được gọiThườ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êngtimer(due, period?)Phát 0 sau delay; có thể tiếp tục tăngCó nếu không truyền period
Cần tạo nguồn mới tại thời điểm subscribedefer(factory)Gọi factory riêng cho mỗi subscriberTheo 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 subscribe

of() 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: 30

of([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 xong

String 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
complete

Khá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;
  • NodeList hoặc HTMLCollection chứ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 B

Vì 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: complete

interval() 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:

SourceHành viTrường hợp dùng
EMPTYKhông phát giá trị, complete ngayBỏ qua một nhánh nhưng vẫn kết thúc sạch
NEVERKhông phát, không error, không completeTest hoặc giữ một nhánh im lặng có chủ đích
throwError(() => error)Gửi error ngay khi subscribeTrả 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 completeDữ 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ụ:

  1. 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.
  2. Mỗi tick khiến switchMap() subscribe vào một source mới.
  3. defer() gọi loadStock() tại thời điểm đó, nên request không chạy trước subscription.
  4. from() chuyển kết quả hoặc lỗi của Promise sang Observable.
  5. catchError() nằm trong inner stream, nên một request lỗi trở thành StockState và không giết clock polling bên ngoài.
  6. 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.

SourceThời điểm mặc địnhKết thúc mặc địnhCleanup đáng chú ý
of(...)Đồng bộComplete sau argument cuốiThường không có resource ngoài
from(array/iterable)Đồng bộComplete khi duyệt xongDừng duyệt nếu subscriber đã đóng
from(promise)Bất đồng bộ theo PromiseComplete sau resolve, error sau rejectUnsubscribe không hủy Promise
fromEvent(...)Khi target phát eventThường không completeUnsubscribe gỡ listener
interval(...)Sau mỗi periodKhông completeUnsubscribe hủy lịch tick
timer(due)Sau delayComplete sau một emissionUnsubscribe 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

  1. Dùng of(array) khi muốn từng phần tử. Downstream nhận một array duy nhất. Đổi sang from(array) hoặc giữ of(array) nếu array thật sự là một domain value.
  2. Dùng from(value) cho scalar. from(42) không hợp lệ; of(42) mới đúng.
  3. 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(...))).
  4. 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.
  5. Chờ complete từ 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.
  6. 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.
  7. Nhầm interval(1000) phát ngay. Tick 0 đến sau khoảng một giây. Dùng timer(0, 1000) nếu cần emission đầu tiên ngay.
  8. Dùng defer như cache. Nó làm điều ngược lại: tạo source mới theo từng subscriber.
  9. 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.
  10. 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ên scheduled() 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

  1. Chạy of([1, 2, 3]) và from([1, 2, 3]), rồi ghi lại số lần callback next được gọi.
  2. Đặt console.log trước và sau subscribe() cho of('A') và from(Promise.resolve('A')); giải thích thứ tự khác nhau.
  3. 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.
  4. Đổi interval(1_000) thành timer(0, 1_000) và quan sát thời điểm tick đầu tiên.
  5. 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

Nguồn tham khảo

On this page