Học RxJS
Thời gian & Schedulers

Mô hình Scheduler

Vai trò của Scheduler trong việc lập lịch công việc.

Bạn gọi subscribe() rồi thấy một pipeline phát ngay trong cùng call stack, pipeline khác đợi sang lượt event loop, còn animationFrameScheduler lại chờ sát lần repaint tiếp theo. Nếu chỉ gắn nhãn “sync” và “async”, những khác biệt này rất dễ biến thành các mẹo phải học thuộc.

Scheduler cho mình một mental model rõ hơn: công việc nào đang chờ, đồng hồ nào quyết định thời điểm đến hạn, và cơ chế nào thực thi công việc đó. Hiểu ba mảnh này giúp bạn đọc đúng thứ tự notification, chọn chỗ đặt observeOn hoặc subscribeOn, và biết vì sao đổi scheduler không đồng nghĩa với chạy trên thread khác.

Phạm vi phiên bản

Bài này dùng public API của RxJS 7.x và đối chiếu hành vi với RxJS 7.8.2. Ví dụ import từ rxjs; bài không giả định API thử nghiệm của RxJS 8.

Mục lục

Scheduler giải quyết vấn đề gì

Một Observable có thể tạo công việc theo nhiều cách: đọc lần lượt một array, đặt timer, nghe DOM event hoặc chuyển notification sang một hàng đợi khác. Dù cơ chế bên dưới khác nhau, pipeline vẫn cần trả lời cùng một nhóm câu hỏi:

  • Công việc chạy ngay hay phải xếp hàng?
  • Nếu có nhiều công việc cùng chờ, chúng chạy theo thứ tự nào?
  • delay được đo bằng đồng hồ thật hay virtual time?
  • Khi consumer unsubscribe, công việc chưa chạy có bị hủy không?

Scheduler gom các quyết định đó vào một abstraction. Nó không thay producer và cũng không tự biến source thành concurrent. Nó chỉ cung cấp cách lập lịch một đơn vị công việc theo một quy tắc thời gian và thứ tự cụ thể.

Đây là lý do nhiều API thời gian nhận một scheduler tùy chọn. interval, timer, debounceTime hay delay cần một khái niệm “bây giờ” và một cơ chế đánh thức công việc khi đến hạn. Trong ứng dụng thông thường, RxJS chọn scheduler mặc định phù hợp; khi test hoặc cần đổi ranh giới thực thi, bạn mới truyền scheduler khác.

Mặc định đừng chọn scheduler bằng thói quen

Nếu pipeline đã đúng thứ tự và đúng thời gian, hãy giữ scheduler mặc định. Thêm scheduler tạo thêm hàng đợi, thay đổi timing và làm debug khó hơn; nó không phải nút “tăng tốc” cho RxJS.

Mental model của Scheduler

Repo này chưa cấu hình Mermaid, nên sơ đồ dùng ASCII để hiển thị đúng mà không cần component hoặc dependency bổ sung.

schedule(work, delay, state)
            │
            ▼
┌────────────────────────────────────────────────────────┐
│ Scheduler                                              │
│                                                        │
│  clock: now() ──► tính lúc đến hạn                     │
│                         │                              │
│  queue: lưu Action ────┼──► sắp thứ tự công việc      │
│                         │                              │
│  execution context ─────┴──► chạy Action đúng cơ chế   │
└────────────────────────────────────────────────────────┘
            │
            ▼
Subscription dùng để hủy Action

Đừng hiểu “execution context” ở đây là một thread mới. Trong các scheduler dựng sẵn của RxJS, nó thường là lựa chọn giữa thực thi đồng bộ, microtask-like queue, timer của event loop hoặc callback của requestAnimationFrame.

Ba trách nhiệm trong một abstraction

Scheduler là một cấu trúc dữ liệu. Nó giữ các Action đang chờ và quyết định thứ tự xử lý. Tiêu chí thường liên quan đến thời điểm đến hạn và thứ tự được enqueue; đây không phải một public priority queue để ứng dụng tự gán mức ưu tiên nghiệp vụ.

Scheduler là một execution context. Nó quyết định khi nào và qua cơ chế nào callback chạy. queueScheduler chạy trong lượt hiện tại, asapScheduler chờ phần code đồng bộ hiện tại kết thúc, còn asyncScheduler dùng timer trên event loop.

Scheduler có một clock. Phương thức now() trả về mốc thời gian theo đồng hồ của scheduler. Với scheduler thông thường, đồng hồ bám theo thời gian runtime; với virtual scheduler, “thời gian” có thể chỉ là một frame do test điều khiển. Vì vậy delay = 1000 có nghĩa là 1.000 đơn vị theo clock đó, không bắt buộc test phải đợi một giây ngoài đời.

Ba trách nhiệm đi cùng nhau vì một queue thời gian không thể chỉ biết “đợi 500”: nó phải biết hiện tại là bao nhiêu, Action nào đến hạn trước và dùng cơ chế nào để thực thi.

Action là đơn vị công việc

Ở mức public contract, bạn thường làm việc với SchedulerLike thay vì tự tạo subclass từ lớp Scheduler:

interface SchedulerLike {
  now(): number;

  schedule<T>(
    work: (this: SchedulerAction<T>, state?: T) => void,
    delay?: number,
    state?: T,
  ): Subscription;
}

Lời gọi schedule(work, delay, state) tạo một Action:

  • work là callback cần chạy;
  • delay là khoảng chờ tương đối theo clock của scheduler, mặc định là 0;
  • state là dữ liệu được truyền vào callback;
  • giá trị trả về là một Subscription, nên có thể unsubscribe() để hủy công việc còn chờ.

Bên trong callback viết bằng function, this trỏ đến Action hiện tại. Bạn có thể gọi this.schedule(nextState, nextDelay) để lên lịch lại chính Action đó. Arrow function không có this riêng, nên không dùng được mẫu reschedule này.

Đừng kế thừa lớp Scheduler

Trong RxJS 7.8.2, lớp Scheduler được đánh dấu là internal implementation detail và deprecated cho mục đích mở rộng. Nếu thật sự cần một scheduler tùy biến, hãy triển khai contract SchedulerLike; với code ứng dụng, gần như luôn nên dùng scheduler dựng sẵn hoặc TestScheduler.

Observable không mặc định là bất đồng bộ

Observable mô tả cách phát notification; nó không cam kết callback sẽ chạy ở một thời điểm khác. of(1, 2, 3) phát đồng bộ khi được subscribe:

import { of } from 'rxjs';

console.log('trước');

of(1, 2, 3).subscribe((value) => console.log('next:', value));

console.log('sau');

Kết quả:

trước
next: 1
next: 2
next: 3
sau

Không có scheduler nào chen một ranh giới async vào ví dụ này. subscribe() chưa return cho đến khi of phát xong và complete.

Khi muốn chuyển việc duyệt input sang một scheduler, RxJS 7 cung cấp scheduled:

import { asapScheduler, scheduled } from 'rxjs';

console.log('trước');

scheduled([1, 2, 3], asapScheduler).subscribe({
  next: (value) => console.log('next:', value),
  complete: () => console.log('complete'),
});

console.log('sau');

Kết quả:

trước
sau
next: 1
next: 2
next: 3
complete

Ở đây source được tạo theo scheduler ngay từ đầu. Khác biệt này sẽ quan trọng khi so với observeOn: một bên điều khiển cách source lấy và phát dữ liệu, bên kia nhận notification đã phát rồi mới xếp lịch chuyển tiếp.

Bốn scheduler dựng sẵn

RxJS 7 có bốn scheduler bạn sẽ gặp thường xuyên:

Schedulerdelay = 0 chạy khi nàoHợp vớiChi tiết dễ nhầm
queueSchedulerĐồng bộ trong lượt hiện tại, Action lồng nhau vào queueTránh đệ quy call stack, tuần tự hóa công việc syncCó delay > 0 thì hành vi giống asyncScheduler
asapSchedulerSau phần code đồng bộ hiện tại, sớm nhất có thểDefer công việc mà không cần timer nghiệp vụKhông bảo đảm chen trước task đã được enqueue trước nó
asyncSchedulerQua timer trên event loopDelay, interval và operator dựa trên thời giandelay = 0 vẫn là bất đồng bộ
animationFrameSchedulerNgay trước lần browser repaint tiếp theoĐọc/ghi UI theo frame, animationCó delay > 0 thì dùng hành vi async; phụ thuộc môi trường browser

queueScheduler xếp hàng nhưng vẫn đồng bộ

Tên “queue” dễ khiến người đọc nghĩ nó bất đồng bộ. Thực tế, với delay = 0, Action đầu tiên chạy ngay trong lời gọi schedule(). Điểm khác là Action được schedule lồng bên trong không gọi đệ quy ngay; nó đợi Action hiện tại hoàn tất rồi mới chạy.

import { queueScheduler } from 'rxjs';

queueScheduler.schedule(() => {
  console.log('A: bắt đầu');

  queueScheduler.schedule(() => {
    console.log('B: công việc lồng');
  });

  console.log('A: kết thúc');
});

console.log('C: sau schedule');

Kết quả:

A: bắt đầu
A: kết thúc
B: công việc lồng
C: sau schedule

B không chen vào giữa A, nhưng toàn bộ queue vẫn được flush trước khi lời gọi schedule() ngoài cùng return. Cơ chế này thường được gọi là trampoline: thay vì lún sâu thêm vào call stack, công việc lặp lại được đặt cuối hàng đợi hiện tại.

asapScheduler và asyncScheduler khác nhau ở lượt chờ

Cả hai đều có thể dời công việc ra khỏi đoạn code đồng bộ hiện tại, nhưng chúng vào hai hàng đợi khác nhau. Với delay = 0, asapScheduler dùng cơ chế microtask-like và thường chạy trước công việc timer của asyncScheduler.

import { asapScheduler, asyncScheduler } from 'rxjs';

asyncScheduler.schedule(() => console.log('async'));
asapScheduler.schedule(() => console.log('asap'));

console.log('sync');

Kết quả theo mô hình scheduler của RxJS 7:

sync
asap
async

Đừng biến ví dụ này thành quy tắc “asap luôn thắng mọi thứ”. Thứ tự còn phụ thuộc công việc nào đã được enqueue trước trong cùng queue. Stance hữu ích hơn là: dùng asapScheduler khi mục tiêu chỉ là defer đến sau đoạn sync hiện tại; dùng asyncScheduler khi logic thật sự gắn với delay hoặc chu kỳ thời gian.

animationFrameScheduler bám theo nhịp vẽ

Trong browser, animationFrameScheduler lên lịch Action ngay trước lần repaint tiếp theo. Nó phù hợp khi downstream cập nhật vị trí, kích thước hoặc style và bạn muốn công việc đi cùng nhịp render.

import { animationFrameScheduler } from 'rxjs';

animationFrameScheduler.schedule(() => {
  console.log('chạy trước repaint tiếp theo');
});

Scheduler này không làm phép tính nặng trở nên rẻ hơn. Nếu callback chiếm 30 ms, main thread vẫn bị giữ 30 ms và frame vẫn có thể giật. Ngoài ra, khi truyền delay dương, animationFrameScheduler rơi về hành vi của asyncScheduler thay vì chờ đúng một frame tùy ý.

Browser API không phải môi trường nào cũng có

animationFrameScheduler dựa vào requestAnimationFrame. Đừng dùng nó như mặc định trong code chạy trên Node.js hoặc server-side rendering nếu môi trường chưa cung cấp API tương ứng.

Scheduler đi vào pipeline ở đâu

Cùng một scheduler nhưng đặt ở ba vị trí khác nhau sẽ điều khiển ba việc khác nhau:

subscribe() ─► subscribeOn ──[schedule]──► setup source
                                             │
                                             ▼
source ─► operators upstream ─► observeOn ──[schedule]──► observer
   ▲              │
   └──── source tự dùng scheduler để tạo và phát dữ liệu

Luồng subscription đi ngược lên source; notification đi từ source xuống Observer. Vì vậy subscribeOn và observeOn không phải hai tên cho cùng một thao tác.

Lập lịch ngay tại source

Khi source hoặc creation function nhận scheduler, scheduler có thể điều khiển chính quá trình lấy dữ liệu và phát notification. Ví dụ scheduled(array, scheduler) lên lịch việc duyệt array thay vì để array được duyệt hết đồng bộ rồi mới dời notification ở cuối pipeline.

Các operator dựa trên thời gian cũng có scheduler riêng, thường mặc định là asyncScheduler:

import { asyncScheduler, debounceTime, interval } from 'rxjs';

const ticks$ = interval(1_000, asyncScheduler);

const stableTicks$ = ticks$.pipe(
  debounceTime(300, asyncScheduler),
);

Trong production, truyền asyncScheduler tường minh như trên thường là dư thừa vì đó đã là mặc định. Tham số scheduler hữu ích hơn khi bạn muốn làm rõ contract của adapter hoặc thay clock thật bằng virtual time trong test.

observeOn dời notification downstream

observeOn(scheduler) đặt một ranh giới tại vị trí của nó trong pipeline. Source và các operator phía trước vẫn chạy theo cơ chế cũ; chỉ next, error và complete đi xuống phía sau observeOn mới được schedule lại.

import { asyncScheduler, Observable, observeOn, tap } from 'rxjs';

const source$ = new Observable<number>((subscriber) => {
  console.log('producer: bắt đầu');
  subscriber.next(1);
  subscriber.next(2);
  subscriber.complete();
});

console.log('trước subscribe');

source$
  .pipe(
    tap((value) => console.log('upstream:', value)),
    observeOn(asyncScheduler),
    tap((value) => console.log('downstream:', value)),
  )
  .subscribe({
    complete: () => console.log('complete'),
  });

console.log('sau subscribe');

Kết quả:

trước subscribe
producer: bắt đầu
upstream: 1
upstream: 2
sau subscribe
downstream: 1
downstream: 2
complete

Producer vẫn chạy đồng bộ và đã phát cả hai giá trị trước dòng sau subscribe. observeOn chỉ xếp lịch việc giao các notification đó cho downstream.

observeOn không chia nhỏ producer đồng bộ

Đặt observeOn(asyncScheduler) sau một source phát hàng nghìn giá trị đồng bộ không ngăn source tạo hàng nghìn notification ngay lập tức; nó còn có thể enqueue hàng nghìn Action downstream. Nếu cần chia quá trình tạo dữ liệu thành các lượt, hãy lập lịch tại source hoặc thiết kế producer theo chunk.

subscribeOn dời thời điểm subscribe source

subscribeOn(scheduler) schedule việc gọi subscribe vào source. Vì setup của cold Observable thường xảy ra lúc subscribe, operator này có thể dời cả side effect khởi tạo producer.

import { asapScheduler, merge, of, subscribeOn } from 'rxjs';

const deferred$ = of('A1', 'A2').pipe(
  subscribeOn(asapScheduler),
);

const direct$ = of('B1', 'B2');

merge(deferred$, direct$).subscribe({
  next: (value) => console.log(value),
  complete: () => console.log('complete'),
});

Kết quả:

B1
B2
A1
A2
complete

merge gặp deferred$ trước, nhưng subscription của nhánh đó được xếp lịch. Trong lúc chờ, direct$ được subscribe và phát đồng bộ, nên nhóm B xuất hiện trước.

Cách nhớ ngắn gọn:

Công cụNó schedule điều gì?Phần nào bị ảnh hưởng?
Scheduler tại sourceQuá trình tạo/lấy và phát dữ liệuChính producer hoặc creation function đó
subscribeOnLời gọi subscribe vào sourceSetup source và các side effect xảy ra khi subscribe
observeOnViệc giao next, error, completeChỉ downstream kể từ vị trí operator

Hủy một Action đã lên lịch

schedule() trả về một Subscription. Nếu Action chưa chạy, unsubscribe() loại hoặc vô hiệu hóa công việc đang chờ:

import { asyncScheduler } from 'rxjs';

const pending = asyncScheduler.schedule(() => {
  console.log('dòng này không chạy');
}, 1_000);

pending.unsubscribe();

Với Action lặp lại, callback có thể reschedule chính nó qua this.schedule(...):

import { asyncScheduler } from 'rxjs';

const countdown = asyncScheduler.schedule(
  function (remaining) {
    console.log(remaining);

    if (remaining > 1) {
      this.schedule(remaining - 1, 1_000);
    }
  },
  0,
  3,
);

setTimeout(() => countdown.unsubscribe(), 1_500);

Output xấp xỉ:

3
2

Action cho 1 đã được lên lịch ở mốc khoảng hai giây, nhưng subscription bị hủy trước đó. Thời điểm thực có thể lệch vì event loop bận; scheduler giữ quy tắc thứ tự và thời điểm sớm nhất, chứ không biến timer JavaScript thành đồng hồ real-time cứng.

Có một giới hạn hiển nhiên: cancellation không quay ngược công việc đã chạy. Với queueScheduler, Action không delay có thể hoàn tất ngay trong lời gọi schedule(), nên unsubscribe() ở dòng sau chỉ đóng subscription đã xong.

Scheduler không tạo parallelism

Đây là hiểu lầm đáng dẹp sớm nhất: đổi scheduler không chuyển JavaScript sang CPU core khác. asapScheduler, asyncScheduler và animationFrameScheduler thay đổi thời điểm callback được main event loop gọi; chúng không tạo Web Worker và không làm hai callback chạy đồng thời.

Giả sử một phép tính chiếm CPU 200 ms. Đưa nó vào asyncScheduler có thể cho đoạn code hiện tại kết thúc trước khi phép tính bắt đầu, nhưng khi Action chạy, main thread vẫn bị chặn 200 ms. Nếu UI cần giữ responsive, bạn phải giảm lượng việc mỗi chunk, dùng thuật toán khác hoặc chuyển việc nặng sang Web Worker/worker thread rồi bọc kết quả thành Observable.

Scheduler cũng không phải cơ chế backpressure. Nếu producer phát nhanh hơn consumer xử lý, observeOn chỉ có thể dồn thêm Action vào queue. Bạn vẫn phải chọn chiến lược nghiệp vụ như sampling, buffering có giới hạn, dropping hoặc điều tiết producer.

Cách chọn scheduler

Mình dùng thứ tự quyết định này:

  1. Không truyền scheduler nếu default của source/operator đã đúng. Đây là lựa chọn phổ biến nhất.
  2. Chọn queueScheduler khi công việc phải giữ đồng bộ nhưng cần xếp hàng Action lồng nhau để tránh đệ quy sâu.
  3. Chọn asapScheduler khi chỉ cần defer đến sau đoạn sync hiện tại, không có delay nghiệp vụ.
  4. Chọn asyncScheduler cho timer, delay và công việc lặp theo thời gian trên event loop.
  5. Chọn animationFrameScheduler cho cập nhật UI cần bám lần repaint tiếp theo trong browser.
  6. Chọn virtual scheduler hoặc TestScheduler khi test logic thời gian mà không muốn chờ đồng hồ thật.

Một bảng hỏi nhanh khi review code:

Câu hỏiNếu câu trả lời là có
Source/operator hiện tại đã có đúng timing chưa?Giữ mặc định
Cần đổi lúc producer bắt đầu setup?Xem subscribeOn
Producer giữ nguyên nhưng downstream cần đổi ranh giới giao notification?Xem observeOn
Cần dời chính quá trình duyệt hoặc tạo dữ liệu?Lập lịch tại source, chẳng hạn scheduled
Công việc gắn với repaint của browser?animationFrameScheduler
Mục tiêu là chạy CPU song song?Scheduler không giải quyết; dùng worker
Mục tiêu là test thời gian xác định?Dùng virtual time/TestScheduler

Những bẫy thường gặp

  1. Đồng nhất Observable với async. of, range và custom Observable có thể phát trong cùng call stack. Hãy đọc source và scheduler thay vì đoán từ kiểu Observable.
  2. Xem queueScheduler như một event-loop queue. Với delay bằng 0, nó flush đồng bộ. Queue ở đây chủ yếu quản lý Action lồng nhau.
  3. Dùng observeOn để “làm nhẹ” producer. Source vẫn có thể phát hết đồng bộ trước ranh giới đó, còn downstream nhận thêm một hàng đợi lớn.
  4. Nhầm subscribeOn với observeOn. Một cái schedule việc nối vào source; cái kia schedule việc giao notification xuống dưới.
  5. Cho rằng asyncScheduler tạo thread nền. Callback vẫn chạy trên JavaScript event loop của môi trường hiện tại.
  6. Quên hủy Action sống lâu. Action lặp lại hoặc delay dài vẫn là resource của subscription. Gắn nó vào lifecycle sở hữu và unsubscribe khi không còn cần.
  7. Dùng animationFrameScheduler cho mọi timer UI. Nó hợp với repaint, không phải deadline nghiệp vụ hay polling trên server.
  8. Tin rằng delay là thời điểm tuyệt đối. Timer chỉ đảm bảo không chạy sớm hơn mốc đến hạn; event loop bận có thể làm Action chạy muộn.
  9. Tự viết custom scheduler quá sớm. Scheduler liên quan đến ordering, clock, cancellation và reentrancy. Dùng scheduler dựng sẵn hoặc TestScheduler trước khi nhận thêm độ phức tạp này.

Bài tập tự kiểm tra

Đọc đoạn code sau mà chưa chạy:

import {
  asapScheduler,
  asyncScheduler,
  observeOn,
  of,
  subscribeOn,
  tap,
} from 'rxjs';

console.log('A');

of(1)
  .pipe(
    tap(() => console.log('B')),
    subscribeOn(asapScheduler),
    observeOn(asyncScheduler),
  )
  .subscribe({
    next: () => console.log('C'),
    complete: () => console.log('D'),
  });

console.log('E');

Thứ tự là:

A
E
B
C
D

Lý do:

  1. A chạy đồng bộ.
  2. subscribeOn(asapScheduler) chưa cho source được subscribe ngay, nên dòng cuối chạy và in E.
  3. Ở lượt asap, of(1) được subscribe và tap phía upstream in B.
  4. observeOn(asyncScheduler) schedule lại next(1) và complete, nên C rồi D đến ở lượt async sau đó.

Nếu bạn đoán B xuất hiện cùng C, hãy quay lại ranh giới trong sơ đồ: observeOn chỉ ảnh hưởng downstream, còn tap nằm trước nó.

Nguồn tham khảo

Học tiếp

On this page