Time-based operators
Tạo và điều khiển stream phụ thuộc thời gian.
Bạn cần refresh dashboard mỗi 30 giây, trì hoãn một notification, hoặc báo lỗi khi server im lặng quá lâu. Nếu rải setTimeout và setInterval vào callback, phần khó không nằm ở việc tạo timer mà ở chỗ hủy timer, giữ đúng thứ tự và nối lỗi vào lifecycle của stream.
RxJS đưa những quy tắc đó vào pipeline. timer và interval tạo nhịp thời gian; delay và delayWhen dời emission; timeout đặt deadline; timestamp và timeInterval giúp quan sát thời điểm. Bài này tập trung vào cách chọn và ghép các API đó, còn nhóm operator giảm tần suất event được tách sang bài riêng.
Phạm vi phiên bản
Ví dụ trong bài nhắm tới RxJS 7.8.x và import public API từ rxjs. Các mốc thời gian dùng milliseconds trừ khi đoạn code ghi rõ khác.
Mục lục
- Mental model của thời gian trong RxJS
- Chọn API theo ý định
- Tạo nhịp với timer và interval
- Dời emission với delay
- Delay động với delayWhen
- Đặt deadline với timeout
- Đo thời gian với timestamp và timeInterval
- Scheduler liên quan như thế nào
- Ví dụ hoàn chỉnh polling có deadline
- Những bẫy thường gặp
- Checklist thiết kế pipeline theo thời gian
- Nguồn tham khảo
- Học tiếp
Mental model của thời gian trong RxJS
Time-based operator không sở hữu một chiếc đồng hồ chính xác riêng. Nó hỏi Scheduler hiện tại là mấy giờ và nhờ scheduler lên lịch công việc. Với scheduler mặc định, callback cuối cùng vẫn phải đi qua event loop của JavaScript, nên 1_000 ms có nghĩa là không chạy trước khoảng đó, chứ không đảm bảo chạy đúng tuyệt đối ở millisecond thứ 1.000.
Hãy tách ba ý thường bị trộn lẫn:
Tạo timeline timer, interval
source mới ───────────────────────────────► next
Dời notification delay, delayWhen
source next ──► giữ theo lịch ────────────► next muộn hơn
Giám sát timeline timeout, timestamp, timeInterval
source next ──► đo hoặc kiểm tra deadline ─► next / fallback / errordelay(500) không biến code nặng thành background work, và timeout(2_000) không làm request chạy nhanh hơn. Chúng chỉ thay đổi cách notification đi qua một Observable execution. Công việc nền có thực sự dừng khi unsubscribe hay không vẫn phụ thuộc teardown của producer.
Thời gian là yêu cầu, không phải độ chính xác tuyệt đối
Browser có thể throttle timer ở tab nền; Node.js và browser đều có thể chạy callback muộn khi event loop bận. Nếu nghiệp vụ cần mốc thời gian tuyệt đối, hãy lưu deadline và so sánh lại với clock thay vì đếm số tick rồi giả định không có drift.
Chọn API theo ý định
| Bạn muốn làm gì? | API mặc định | Điểm cần nhớ |
|---|---|---|
| Phát một lần sau một khoảng chờ | timer(due) | Phát 0, rồi complete |
| Phát ngay rồi lặp theo chu kỳ | timer(0, period) | Tick đầu được lên lịch gần như ngay |
| Chờ một chu kỳ rồi phát đều | interval(period) | Tick đầu chỉ đến sau một period |
| Dời mọi value cùng một khoảng | delay(due) | Dời notification, không dời lúc subscribe vào source |
| Mỗi value có thời gian chờ riêng | delayWhen(selector) | Duration Observable phải phát next |
| Giới hạn thời gian chờ value | timeout(config) | Timeout có thể error hoặc chuyển sang fallback |
| Gắn thời điểm tuyệt đối vào value | timestamp() | Mặc định là số milliseconds theo epoch |
| Đo khoảng cách giữa hai value | timeInterval() | Value đầu đo từ lúc subscribe |
| Giảm event dày đặc | debounceTime, throttleTime, auditTime, sampleTime | Chúng chọn hoặc bỏ value, không chỉ dời value |
Mình sẽ chọn API theo ý định nghiệp vụ trước, rồi mới nghĩ đến scheduler. Nếu câu yêu cầu là “không được im lặng quá 3 giây”, timeout({ each: 3_000 }) diễn đạt đúng hơn việc tự tạo timer rồi đua bằng cờ trạng thái.
Tạo nhịp với timer và interval
Một lần hoặc lặp lại với timer
timer(due) chờ due, phát số 0 một lần rồi complete. Truyền thêm period biến nó thành nguồn lặp: emission đầu theo due, các emission sau cách nhau theo period.
import { take, timer } from 'rxjs';
// 0 sau khoảng 1,5 giây, rồi complete.
timer(1_500).subscribe({
next: (value) => console.log('one-shot:', value),
complete: () => console.log('one-shot complete'),
});
// 0 được lên lịch ngay; 1 và 2 cách nhau khoảng một giây.
timer(0, 1_000)
.pipe(take(3))
.subscribe({
next: (value) => console.log('periodic:', value),
complete: () => console.log('periodic complete'),
});Timeline mong đợi:
time: 0 ms 1.000 ms 1.500 ms 2.000 ms
periodic$: (0)───────────(1)───────────────────────────(2)│
oneShot$: ─────────────────────────────(0)│due cũng có thể là một Date. Cách này hợp với deadline gần như “đến đầu phút tiếp theo”, nhưng không nên dùng một timer rất dài như lịch công việc bền vững. Process restart, sleep hoặc giới hạn của setTimeout có thể phá giả định đó; job quan trọng nên được lưu trong scheduler hoặc queue phù hợp.
Nhịp đều với interval
interval(period) phát 0, 1, 2, ... và không tự complete. Khác biệt dễ quên là tick 0 chỉ xuất hiện sau khi chờ đủ một period.
import { interval, take } from 'rxjs';
interval(1_000)
.pipe(take(3))
.subscribe(console.log);time: 0 ms 1.000 ms 2.000 ms 3.000 ms
interval$: ──────────────(0)────────────(1)────────────(2)│Nếu polling phải chạy ngay khi màn hình mở, mình dùng timer(0, period). Nếu yêu cầu nói rõ “chờ một chu kỳ rồi mới chạy”, interval(period) diễn đạt ý đó gọn hơn. Phần Creation functions so sánh chi tiết hai nguồn này với các cách tạo Observable khác.
Timer cũng có lifecycle
timer và interval mặc định là cold: mỗi subscriber tạo một lịch riêng và bắt đầu lại từ 0. Hai subscriber không tự dùng chung một timer chỉ vì chúng subscribe vào cùng một biến Observable.
Unsubscribe sẽ hủy công việc đã lên lịch cho execution đó. Vì interval không complete, bạn phải gắn nó với lifecycle bằng take, takeUntil, Subscription hoặc cơ chế cleanup của framework. Nếu nhiều consumer thật sự cần cùng một clock, hãy thiết kế việc share có chủ đích thay vì vô tình tạo nhiều timer.
Dời emission với delay
Delay cố định
delay(500) dời mỗi next khoảng 500 ms và giữ khoảng cách tương đối giữa các value. Source được subscribe ngay; chỉ notification tới downstream bị giữ lại.
import { delay, finalize, of, tap } from 'rxjs';
of('A', 'B')
.pipe(
tap((value) => console.log('source:', value)),
delay(500),
finalize(() => console.log('done')),
)
.subscribe((value) => console.log('output:', value));Output có dạng:
source: A
source: B
... khoảng 500 ms ...
output: A
output: B
doneof() phát đồng bộ, nên hai dòng source xuất hiện ngay trong call stack của subscribe. delay không trì hoãn side effect phía trước nó; nó chỉ dời các value khi chúng đi qua vị trí của operator trong pipeline. Completion đợi những value đang được giữ phát xong, nhưng error từ source đi qua ngay và các value đang chờ có thể bị hủy.
Muốn trì hoãn lúc bắt đầu công việc
Đừng đặt delay sau một Observable có side effect rồi kỳ vọng side effect bắt đầu muộn. Hãy dùng timer(due).pipe(switchMap(() => source$)), hoặc bọc việc tạo source bằng defer nếu nó phải được khởi tạo tại thời điểm đó.
import { defer, switchMap, timer } from 'rxjs';
const delayedRequest$ = timer(1_000).pipe(
switchMap(() => defer(() => fetch('/api/report'))),
);Ở đây fetch chỉ được gọi sau khi timer phát. Nếu consumer unsubscribe trong một giây đầu, switchMap chưa tạo request.
Mốc Date và thời điểm subscribe
delay(date) giữ các source value cho đến mốc Date đó. Nó vẫn subscribe vào source ngay, vì vậy hot source hoặc side effect upstream vẫn chạy trong lúc downstream chưa thấy value. Các value đến sau khi mốc đã qua được lên lịch với delay bằng 0.
Date tuyệt đối hữu ích khi nhiều value phải được mở tại cùng một mốc. Với độ trễ tương đối cho từng value, dùng một con số như delay(500) dễ đọc và ít phụ thuộc clock hệ thống hơn.
Delay động với delayWhen
delayWhen gọi selector cho từng source value. Value chỉ đi tiếp khi duration Observable tương ứng phát một next. Vì các duration chạy đồng thời qua cơ chế tương tự mergeMap, value có delay ngắn hơn có thể vượt value đến trước.
import { delayWhen, of, timer } from 'rxjs';
type Notice = {
id: string;
waitMs: number;
};
of<Notice>(
{ id: 'slow', waitMs: 600 },
{ id: 'fast', waitMs: 100 },
)
.pipe(delayWhen((notice) => timer(notice.waitMs)))
.subscribe((notice) => console.log(notice.id));Kết quả theo thời gian:
... khoảng 100 ms ... fast
... khoảng 500 ms nữa ... slowĐây là hành vi đúng nếu mỗi item có availability time riêng. Nếu thứ tự nguồn phải được giữ nguyên, delayWhen một mình chưa đủ; bạn cần mô hình tuần tự, chẳng hạn concatMap với timer(...).pipe(map(() => value)), và phải chấp nhận rằng thời gian chờ lúc đó cộng dồn.
Một bẫy riêng của RxJS 7 là duration complete mà không phát next sẽ làm mất source value. Vì vậy delayWhen(() => EMPTY) không có nghĩa “không chờ”; nó nuốt value. Dùng of(null) cho đường đi ngay, hoặc timer(ms) khi cần chờ.
import { delayWhen, of, timer } from 'rxjs';
const output$ = of(0, 250, 500).pipe(
delayWhen((waitMs) => (waitMs === 0 ? of(null) : timer(waitMs))),
);Đặt deadline với timeout
delay trả lời “hãy phát muộn hơn”; timeout trả lời “nếu chờ quá lâu thì coi là thất bại”. Khi deadline hết, timeout unsubscribe source hiện tại rồi error hoặc subscribe vào fallback do with trả về. Nó không phát source value trễ sau đó.
Phân biệt first và each
Cấu hình object giúp ý định rõ hơn overload nhận một con số:
firstgiới hạn thời gian từ lúc subscribe tới value đầu tiên. Nó có thể là milliseconds hoặc mộtDate.eachgiới hạn khoảng im lặng giữa các value. Nếu không cófirst,eachcũng áp dụng cho value đầu tiên.- Khi có cả hai,
firstchỉ quản lý value đầu;eachquản lý những value sau.
source$.pipe(
timeout({
first: 5_000,
each: 2_000,
}),
);Pipeline trên cho source tối đa 5 giây để khởi động, rồi không cho phép khoảng im lặng giữa hai value vượt 2 giây. Source complete trước deadline thì timer nội bộ được dọn và output complete bình thường.
Chuyển sang fallback
Giả sử status đầu đến ngay nhưng status tiếp theo chậm 2,5 giây. each: 2_000 sẽ hủy source trước khi value done xuất hiện và chuyển sang fallback.
import { concat, map, of, timeout, timer } from 'rxjs';
const status$ = concat(
of('accepted'),
timer(2_500).pipe(map(() => 'done')),
);
status$
.pipe(
timeout({
first: 1_000,
each: 2_000,
with: ({ seen, lastValue }) =>
of(`stale: đã nhận ${seen} value, cuối cùng là ${lastValue}`),
}),
)
.subscribe(console.log);Output:
accepted
... khoảng 2 giây ...
stale: đã nhận 1 value, cuối cùng là acceptedCallback with nhận seen, lastValue và meta nếu bạn truyền metadata trong config. Fallback là một Observable mới; khi fallback complete, output cũng complete. Source cũ đã bị unsubscribe nên done không còn cơ hội đi qua.
Xử lý TimeoutError
Nếu không có with, operator phát TimeoutError. Khi bắt lỗi bằng catchError, đừng đổi mọi lỗi thành timeout; source có thể đã lỗi vì lý do khác.
import {
TimeoutError,
catchError,
of,
throwError,
timeout,
} from 'rxjs';
const guarded$ = source$.pipe(
timeout({ first: 3_000 }),
catchError((error: unknown) => {
if (error instanceof TimeoutError) {
return of({ state: 'unavailable' as const });
}
return throwError(() => error);
}),
);Đặt catchError ở đâu quyết định phần nào của hệ thống sống tiếp. Trong polling, mình thường bắt lỗi bên trong operator xử lý từng tick để một lần timeout không giết cả clock. Nếu timeout phải dừng toàn bộ feature, bắt ở boundary ngoài hoặc để error đi thẳng tới subscriber.
Timeout không tự động abort mọi thứ
timeout luôn unsubscribe source, nhưng công việc nền chỉ dừng nếu source có teardown tương ứng. from(fetch(...)) không làm fetch có thể hủy. Nếu deadline phải đóng request thật, hãy tạo Observable nối unsubscribe với AbortController.abort().
Đo thời gian với timestamp và timeInterval
timestamp() gói mỗi value thành { value, timestamp }. Timestamp mặc định đến từ Date.now(), nên phù hợp để log thời điểm quan sát theo clock hệ thống.
timeInterval() trả { value, interval }, trong đó interval là thời gian từ emission trước đến emission hiện tại. Với value đầu tiên, nó đo từ lúc subscribe.
import { Subject, timeInterval, timestamp } from 'rxjs';
const events$ = new Subject<string>();
events$.pipe(timestamp()).subscribe(({ value, timestamp }) => {
console.log('at:', new Date(timestamp).toISOString(), value);
});
events$.pipe(timeInterval()).subscribe(({ value, interval }) => {
console.log('gap:', interval, 'ms', value);
});
events$.next('connected');
setTimeout(() => events$.next('heartbeat'), 750);Hai operator này quan sát timeline, không thay đổi thời điểm source phát. Con số thực tế có thể lệch vì timer và event loop. Với đo hiệu năng cần độ phân giải cao, telemetry phân tán hoặc clock đơn điệu, đừng mặc định timestamp() thay được công cụ chuyên dụng; hãy đưa timestamp provider phù hợp vào boundary của hệ thống.
Scheduler liên quan như thế nào
timer, interval, delay, timeout và timeInterval đều nhận scheduler hoặc dùng clock từ scheduler. Trong code ứng dụng, mặc định thường là asyncScheduler; timestamp mặc định dùng timestamp provider dựa trên Date.now().
Scheduler có hai vai trò liên quan nhưng khác nhau:
- cung cấp khái niệm “bây giờ” qua
now(); - quyết định khi nào work đã lên lịch được thực thi.
Đổi scheduler có thể đổi cả thứ tự execution lẫn cách đo thời gian, nên mình không truyền scheduler chỉ để “làm code async”. Trường hợp đáng dùng nhất là test: TestScheduler cho phép chạy nhiều giây logic bằng virtual time mà không đợi đồng hồ thật. Phần Mô hình Scheduler giải thích abstraction này; Virtual time đi vào cách kiểm thử.
Ví dụ hoàn chỉnh polling có deadline
Hình dung dashboard phải gọi /api/health ngay khi mở, sau đó mỗi 30 giây. Mỗi request chỉ được chờ 3 giây; request timeout phải bị abort thật; một lần lỗi không được làm polling dừng vĩnh viễn; khi trang rời đi, mọi timer và request đang chạy phải được dọn.
page open
│
▼
timer(0, 30s) ──► request ──► timeout 3s ──► state
│ │ │
│ └─ unsubscribe ─► AbortController.abort()
│
pagehide ─────────► takeUntil ───────► teardown toàn pipelineimport {
Observable,
catchError,
exhaustMap,
fromEvent,
map,
of,
takeUntil,
timeout,
timer,
} from 'rxjs';
type HealthResponse = {
status: 'ok';
};
type HealthState =
| { state: 'online'; checkedAt: number }
| { state: 'unavailable'; checkedAt: number };
function getHealth$(): Observable<HealthResponse> {
return new Observable<HealthResponse>((subscriber) => {
const controller = new AbortController();
fetch('/api/health', { signal: controller.signal })
.then((response) => {
if (!response.ok) {
throw new Error(`HTTP ${response.status}`);
}
return response.json() as Promise<HealthResponse>;
})
.then((body) => {
subscriber.next(body);
subscriber.complete();
})
.catch((error: unknown) => {
if (error instanceof DOMException && error.name === 'AbortError') {
return;
}
subscriber.error(error);
});
return () => controller.abort();
});
}
const pageHidden$ = fromEvent(window, 'pagehide');
const health$ = timer(0, 30_000).pipe(
exhaustMap(() =>
getHealth$().pipe(
timeout({ first: 3_000 }),
map(
(): HealthState => ({
state: 'online',
checkedAt: Date.now(),
}),
),
catchError((error: unknown) => {
console.error('Health check failed', error);
return of<HealthState>({
state: 'unavailable',
checkedAt: Date.now(),
});
}),
),
),
takeUntil(pageHidden$),
);
health$.subscribe((state) => console.log(state));timer(0, 30_000) tạo tick đầu ngay rồi lặp. exhaustMap đảm bảo không mở request mới nếu request trước vẫn còn chạy; với deadline 3 giây và period 30 giây, tình huống đó hiếm nhưng policy vẫn hiện rõ. timeout unsubscribe request sau 3 giây, và teardown của getHealth$ gọi abort().
catchError nằm trong exhaustMap, nên lỗi chỉ biến lần kiểm tra hiện tại thành trạng thái unavailable; outer timer vẫn chờ tick sau. Cuối cùng, takeUntil(pageHidden$) đóng outer subscription, hủy lịch polling và abort request nếu trang rời đi giữa chừng.
Trong ứng dụng thật, hãy phân biệt TimeoutError, lỗi HTTP và lỗi parse nếu UI hoặc retry policy khác nhau. Ví dụ gom chúng vào một trạng thái để tập trung vào lifecycle, không phải để khuyên bạn xóa mất thông tin lỗi.
Những bẫy thường gặp
- Tin rằng timer chạy đúng tuyệt đối. Scheduler đặt mốc sớm nhất; event loop, tab nền và tải hệ thống có thể làm callback muộn.
- Dùng
delayđể trì hoãn side effect upstream. Source vẫn được subscribe ngay. Dùngtimerkết hợpswitchMaphoặcconcatMap, vàdeferkhi cần tạo source muộn. - Chờ
intervalcomplete. Nó không tự complete. Luôn xác định owner và điều kiện teardown cho stream dài hạn. - Quên tick đầu của
interval.interval(1_000)phát0sau khoảng một giây;timer(0, 1_000)mới bắt đầu gần như ngay. - Dùng
EMPTYtrongdelayWhen. Trên RxJS 7, duration complete mà không cónextsẽ nuốt source value. - Cho rằng
timeouthủy được mọi request. Nó unsubscribe Observable; producer phải nối teardown tới API hủy thật nhưAbortController. - Bắt lỗi timeout ở ngoài polling.
catchErrorngoài outer pipeline có thể kết thúc clock sau lần lỗi đầu. Đặt error boundary theo lifecycle bạn muốn giữ. - Dùng
delaythay cho debounce hoặc throttle.delaygiữ mọi value; các operator rate limiting chọn hoặc bỏ value theo cửa sổ thời gian. - Test bằng timer thật. Test sẽ chậm và dễ flaky. Dùng
TestSchedulercho logic phụ thuộc scheduler, rồi chỉ để một ít integration test kiểm tra API timer thật. - Đếm tick để suy ra giờ hiện tại. Tick có thể drift. Khi cần deadline tuyệt đối, lưu timestamp đích và tính lại phần thời gian còn lại.
Checklist thiết kế pipeline theo thời gian
- Yêu cầu là tạo tick, dời notification, đặt deadline, đo thời gian hay giảm tần suất event?
- Tick đầu phải chạy ngay hay sau một chu kỳ?
- Khoảng thời gian là tương đối hay một mốc
Datetuyệt đối? - Source có thể phát mãi không, và ai sở hữu việc unsubscribe?
- Nếu timeout, pipeline phải error, dùng fallback hay tiếp tục ở tick sau?
- Unsubscribe có thật sự dọn timer, listener hoặc request bên dưới không?
- Nhiều subscriber cần timer riêng hay cùng chia sẻ một clock?
- Thứ tự value có được phép đổi khi mỗi value có delay khác nhau không?
- Logic đã có test bằng virtual time thay vì chờ timer thật chưa?
Nếu chưa trả lời được câu teardown và error boundary, pipeline chưa hoàn chỉnh dù happy path đã phát đúng value. Với time-based code, phần kết thúc execution quan trọng ngang với mốc bắt đầu.
Nguồn tham khảo
- RxJS 7.8.2 — mã nguồn
timervàinterval: lịch phát đầu tiên, chu kỳ và completion. - RxJS 7.8.2 — mã nguồn
delayvàdelayWhen: duration Observable, thứ tự và hành vi khi duration chỉ complete. - RxJS 7.8.2 — mã nguồn
timeout:first,each, fallback,TimeoutErrorvà việc unsubscribe source. - RxJS 7.8.2 — mã nguồn
timestampvàtimeInterval: cách gắn mốc và tính khoảng thời gian. - RxJS — Scheduler guide: vai trò của scheduler trong execution context và virtual time.