Học RxJS
Lỗi & hoàn tất

repeat

Đăng ký lại sau khi source hoàn tất.

Bạn gọi API lấy trạng thái một job, request hoàn tất nhưng job vẫn đang chạy. Muốn hỏi lại sau một khoảng chờ, bạn cần một lượt request mới, không phải phát lại response cũ. repeat giúp nối các lượt như vậy: khi source complete, nó subscribe lại phần upstream.

Phạm vi bài viết

Ví dụ dùng TypeScript và API của RxJS 7.8.2, import từ rxjs. Bạn nên phân biệt complete, error và unsubscribe; xem vòng đời stream nếu cần ôn lại. Các ví dụ mô phỏng request bằng timer, không cần backend.

Mục lục

repeat đăng ký lại chứ không lưu dữ liệu

repeat phản ứng với complete bằng một subscription upstream mới. Nó chuyển tiếp các value nhận được, nhưng không lưu chúng để replay. Nếu producer tạo dữ liệu khác trong mỗi lần subscribe, output cũng khác.

Hình dung bạn gọi lại quán để hỏi tình trạng món ăn: mỗi cuộc gọi là một lần hỏi mới, không phải nghe lại bản ghi cuộc gọi trước. Ví von này chỉ đúng nếu source thực sự tạo công việc mới khi được subscribe; một Promise đã tạo sẵn thì không làm điều đó.

Lượt source 1:  A ── B ── complete
                              │ subscribe lại
Lượt source 2:                A ── B ── complete

Output:        A ── B ─────── A ── B ── complete

Sơ đồ trên minh họa repeat(2), không phải thời gian chính xác. Complete của lượt đầu bị repeat giữ lại để bắt đầu lượt tiếp; Observer phía sau chỉ nhận complete khi hết số lượt hoặc khi notifier yêu cầu kết thúc. Các value đã phát không bị xóa giữa các lượt.

import { of, repeat } from 'rxjs';

of('A', 'B').pipe(
  repeat(2),
).subscribe({
  next: (value) => console.log(value),
  complete: () => console.log('complete'),
});

Output:

A
B
A
B
complete

Với of, mọi thứ diễn ra đồng bộ trước khi subscribe() trả về. repeat không tự chuyển công việc sang background hay tự tạo khoảng nghỉ.

Cú pháp và cách đếm count

Chữ ký rút gọn, giữ nguyên contract của RxJS 7.8.2:

interface RepeatConfig {
  count?: number;
  delay?: number | ((count: number) => ObservableInput<any>);
}

function repeat<T>(
  countOrConfig?: number | RepeatConfig,
): MonoTypeOperatorFunction<T>;

MonoTypeOperatorFunction<T> nghĩa là kiểu value không đổi: đầu vào Observable<T>, đầu ra cũng là Observable<T>. Các tên kiểu trong signature thuộc RxJS; ví dụ sử dụng bên dưới không cần bạn tự khai báo lại interface này.

Cách dùngÝ nghĩa
repeat(3)Tối đa 3 lượt subscribe source, tính cả lượt đầu tiên.
repeat({ count: 3 })Cùng cách đếm với repeat(3).
repeat(1)Một lượt source, không đăng ký lại.
repeat(0)Complete ngay, không subscribe source. RxJS 7.8.2 cũng xử lý số âm như vậy.
repeat() hoặc repeat({})Không giới hạn số lượt; vẫn dừng khi gặp error hoặc bị hủy.
repeat({ count: 3, delay: 1_000 })Tối đa 3 lượt, chờ 1 giây giữa các lượt.

Hãy dùng số nguyên không âm cho count, hoặc bỏ nó khi chủ đích lặp không giới hạn. count là giới hạn số lượt, không phải số value: một lượt có thể phát không có value nào, một value hoặc nhiều value.

repeat và retry đếm khác nhau

repeat(3) là tối đa 3 lượt tổng cộng, còn retry(3) cho phép 3 lần thử lại sau lỗi, tức tối đa 4 attempt. Đừng dùng chung một biến cấu hình “số lần lặp” mà không ghi rõ nó đếm tổng lượt hay số lần thử lại.

Giới hạn count không ép lượt hiện tại complete. Nếu source của lượt đầu cứ sống mãi, repeat(3) không thể bắt đầu lượt thứ hai.

Điều gì kích hoạt một lượt mới

repeat chỉ xử lý complete của phần upstream ngay trước nó. Error và cancellation không phải complete.

Tình huống tại sourceHành vi của repeat
Complete và còn lượtĐóng lượt cũ, chờ delay nếu có, rồi subscribe lại.
Complete và hết lượtComplete output.
ErrorChuyển lỗi xuống Observer, không repeat.
Consumer unsubscribeHủy lượt đang chạy hoặc notifier đang chờ; không bắt đầu lượt mới.
Không complete, chẳng hạn interval() chưa bị giới hạnKhông repeat, dù vẫn chuyển tiếp value.
import { repeat, throwError } from 'rxjs';

throwError(() => new Error('request lỗi')).pipe(
  repeat(3),
).subscribe({
  error: (error: Error) => console.log('error:', error.message),
  complete: () => console.log('complete'),
});

// error: request lỗi

Chỉ có một lượt source và không có log complete. Nếu muốn thử lại sau lỗi, dùng retry. Nếu muốn tiếp tục hỏi sau một request thành công nhưng trạng thái vẫn chưa đạt yêu cầu, đó mới là công việc của repeat.

Chờ giữa các lượt bằng delay

Khoảng chờ cố định

delay dạng số tính bằng millisecond và bắt đầu sau khi lượt source complete. Nó không trì hoãn lượt đầu, không trì hoãn từng value và không tạo lịch cố định theo thời điểm bắt đầu request.

import { defer, of, repeat } from 'rxjs';

let round = 0;

const source$ = defer(() => of(++round));

source$.pipe(
  repeat({ count: 3, delay: 1_000 }),
).subscribe({
  next: (value) => console.log('lượt', value),
  complete: () => console.log('complete'),
});

Lượt 1 phát ngay; lượt 2 phát sau khoảng một giây; lượt 3 phát sau khoảng một giây nữa rồi output complete. Không có khoảng chờ thừa sau lượt cuối. Bộ đếm trong ví dụ đặt bên ngoài để nhìn rõ các lượt của một subscriber; phần polling bên dưới sẽ đặt state theo từng subscriber.

Giả sử một request mất 300 ms và delay là 1.000 ms, hai thời điểm bắt đầu liên tiếp cách nhau khoảng 1.300 ms. Đây là các con số minh họa, không phải cam kết timer chính xác. Vì phải chờ complete rồi mới subscribe tiếp, một vòng repeat không tạo các lượt source chồng lên nhau.

Notifier quyết định lúc đăng ký lại

Khi delay là callback, RxJS truyền vào số lượt source đã hoàn tất, bắt đầu từ 1. Callback chỉ chạy khi còn quyền repeat; với count: 3, nó nhận 1 và 2, không nhận 3.

import { defer, of, repeat, timer } from 'rxjs';

let round = 0;

const source$ = defer(() => of(++round));

source$.pipe(
  repeat({
    count: 3,
    delay: (completedRounds) => timer(completedRounds * 1_000),
  }),
).subscribe(console.log);

// 1 ngay; 2 sau khoảng 1 giây; 3 sau khoảng 2 giây tiếp theo.

Observable trả về là notifier, không phải nguồn dữ liệu của output:

Notifier làm gì?Kết quả
Phát next đầu tiênĐóng subscription notifier và bắt đầu lượt source tiếp theo; value của notifier không đi ra output.
Complete mà chưa phát value, chẳng hạn EMPTYComplete toàn output, không repeat nữa.
ErrorError toàn output.
Không phát và không kết thúc, chẳng hạn NEVERChờ mãi cho đến khi bị hủy.

Callback ném lỗi cũng làm output error. Đừng trả of(1_000) với kỳ vọng chờ một giây: nó phát đồng bộ và kích hoạt repeat ngay. Muốn chờ, trả timer(1_000); muốn dừng, trả EMPTY.

Tạo công việc mới với defer

Repeat một Observable không đảm bảo công việc bên dưới được chạy lại. Điều đó phụ thuộc vào producer được khai báo thế nào.

import { defer, from, repeat } from 'rxjs';

// Chạy trong môi trường có fetch.
// fetch bắt đầu ngay tại dòng này, chỉ tạo một Promise.
const promise = fetch('/api/job/status');
const reused$ = from(promise).pipe(repeat(3));

// Mỗi lượt subscribe gọi fetch mới.
const fresh$ = defer(() => from(fetch('/api/job/status'))).pipe(
  repeat({ count: 3, delay: 1_000 }),
);

Đây là ví dụ so sánh cách khai báo, chưa subscribe vào hai pipeline. Subscribe reused$ có thể phát cùng object Response ba lần, nhưng không gọi fetch ba lần. Đọc body của cùng Response nhiều lần còn có thể thất bại vì body đã được consume.

Mình dùng defer khi cần tạo lại Promise, đọc state mới hoặc chạy setup riêng cho mỗi lượt. Những producer vốn tạo execution mới theo subscription, như timer hoặc ajax, không bắt buộc phải bọc thêm chỉ để repeat.

defer không tự bổ sung cancellation cho Promise

defer(() => from(fetch(...))) tạo request mới theo subscription nhưng unsubscribe không tự abort fetch. Với HTTP thật, chọn adapter có teardown phù hợp, chẳng hạn fromFetch từ rxjs/fetch. Nếu cần bao phủ cả việc đọc body, dùng tùy chọn selector của fromFetch để body nằm trong vòng đời request Observable.

Polling đến khi job hoàn tất

Giả sử bạn cần hỏi trạng thái job sau mỗi request hoàn tất, dừng khi nhận done, và hủy khi màn hình đóng. Mình đặt repeat trong phạm vi một request đọc trạng thái; thao tác tạo job phải nằm ngoài vòng lặp để không tạo job mới mỗi lần hỏi.

Ví dụ sau mô phỏng request mất 200 ms, mỗi request phát một trạng thái rồi complete. Khoảng chờ 500 ms là lựa chọn minh họa, không phải khuyến nghị cho mọi API.

import {
  defer,
  finalize,
  map,
  repeat,
  Subject,
  takeUntil,
  takeWhile,
  timer,
} from 'rxjs';

type JobStatus = {
  state: 'running' | 'done';
  progress: number;
};

const stop$ = new Subject<void>();

const polling$ = defer(() => {
  // State mô phỏng được tạo riêng cho mỗi subscriber toàn phiên.
  let requestCount = 0;

  const request$ = defer(() => {
    const round = ++requestCount;

    return timer(200).pipe(
      map((): JobStatus => ({
        state: round >= 3 ? 'done' : 'running',
        progress: Math.min(round * 40, 100),
      })),
    );
  });

  return request$.pipe(
    repeat({ delay: 500 }),
    // inclusive=true: phát cả trạng thái done rồi dừng.
    takeWhile((status) => status.state !== 'done', true),
    // Đặt sau repeat để hủy toàn phiên, kể cả lúc đang chờ.
    takeUntil(stop$),
  );
}).pipe(
  finalize(() => console.log('cleanup phiên polling')),
);

const subscription = polling$.subscribe({
  next: (status) => console.log(status.state, status.progress),
  error: (error: unknown) => console.error('polling lỗi', error),
  complete: () => console.log('complete'),
});

// Khi màn hình đóng, chọn một đường hủy:
// stop$.next();                 // output complete, rồi cleanup
// subscription.unsubscribe();  // cleanup, không gọi complete

Nếu không hủy, output lần lượt là:

running 40
running 80
done 100
complete
cleanup phiên polling

Lượt đầu bắt đầu ngay, hai lượt tiếp bắt đầu sau khi lượt trước complete và đã chờ 500 ms. Khi done đi qua takeWhile, operator complete output và unsubscribe upstream, nên repeat không bắt đầu lượt thứ tư.

Nếu request thật error, phiên này dừng với error. Muốn chịu lỗi tạm thời, có thể đặt retry({ count: 2, delay: 1_000 }) trước repeat để từng request có tối đa ba attempt. Retry hết hạn vẫn làm toàn phiên error; đây không phải cơ chế biến lỗi thành trạng thái running.

Một job kẹt ở running có thể khiến polling sống mãi. Trong production, thêm giới hạn lượt hoặc deadline cho toàn phiên và thể hiện rõ “hết thời gian chờ” trên UI. Hết count chỉ tạo complete của Observable, không chứng minh job đã hoàn thành.

Vị trí operator quyết định phạm vi lặp

repeat đăng ký lại toàn đoạn upstream của nó, không chỉ creation function gần nhất. Side effect, bộ đếm và operator trong đoạn đó đều cần được xem lại theo từng lượt.

take và takeWhile trước hay sau repeat

Hai pipeline sau là ví dụ độc lập:

import { of, repeat, take } from 'rxjs';

of('A', 'B', 'C').pipe(
  take(2),
  repeat(2),
).subscribe(console.log);
// A B A B: take(2) được tạo lại cho mỗi lượt.

of('A', 'B', 'C').pipe(
  repeat(2),
  take(4),
).subscribe(console.log);
// A B C A: take(4) giới hạn tổng value của output.

Tương tự, đặt takeWhile trước repeat có thể biến điều kiện “dừng” thành một complete mà repeat dùng để bắt đầu lại. Với polling đến trạng thái cuối, hãy đặt takeWhile sau repeat, như ví dụ ở trên.

takeUntil và việc hủy toàn vòng lặp

Nếu takeUntil(stop$) nằm trước repeat, stop$.next() complete một lượt source. Khi vẫn còn lượt, repeat có thể subscribe lại sau đó — ngược với ý định dừng toàn phiên. Với Subject thông thường, lượt đăng ký mới không nhận lại tín hiệu stop đã phát.

Đặt takeUntil(stop$) sau repeat để tín hiệu stop đóng cả output, unsubscribe lượt đang chạy hoặc notifier đang đợi. Chỉ gọi stop$.complete() không dừng takeUntil: notifier phải phát value. Xem thêm Subscription và teardown để phân biệt complete output với hủy subscription upstream.

catchError có thể biến lỗi thành lượt lặp mới

Khi fallback complete ở upstream, repeat thấy complete đó chứ không thấy lỗi ban đầu nữa:

import { catchError, EMPTY, repeat, throwError } from 'rxjs';

throwError(() => new Error('server lỗi')).pipe(
  catchError(() => EMPTY),
  repeat({ count: 3, delay: 500 }),
).subscribe({
  complete: () => console.log('complete nhưng không có dữ liệu'),
});

Source lỗi ba lượt, nhưng Observer chỉ thấy complete sau lượt cuối. repeat không bắt lỗi; chính catchError đã chuyển lỗi thành complete. Nếu bỏ giới hạn count và delay, pattern này còn có thể tạo vòng lặp đồng bộ không có value.

Mặc định mình để retry xử lý lỗi tạm thời và giữ lỗi cuối đi ra UI. Chỉ đặt fallback trước repeat khi nghiệp vụ thật sự coi fallback ấy là một lượt polling hợp lệ. Đặt catchError sau repeat thì nó xử lý lỗi của toàn chuỗi; fallback ở đó không tự khiến repeat phía trước khởi động lại.

finalize theo từng lượt hay toàn chuỗi

finalize trước repeat cleanup mỗi subscription upstream; finalize sau repeat cleanup toàn phiên, kể cả khi phiên bị hủy trong khoảng chờ.

import { defer, finalize, of, repeat } from 'rxjs';

let round = 0;

const source$ = defer(() => {
  const currentRound = ++round;
  return of(currentRound).pipe(
    finalize(() => console.log('cleanup lượt', currentRound)),
  );
});

source$.pipe(
  repeat({ count: 2, delay: 100 }),
  finalize(() => console.log('cleanup phiên')),
).subscribe(console.log);

Bạn sẽ có hai cleanup lượt và một cleanup phiên. Đặt loading chung của phiên ở phạm vi ngoài; đặt cleanup resource riêng của từng request trong request. Không dùng thứ tự giữa mọi finalizer để phối hợp nghiệp vụ; xem finalize để phân biệt scope và thứ tự teardown.

Những trường hợp không nên repeat trực tiếp

Source đồng bộ và lặp không giới hạn

of(1).pipe(repeat()) hoặc EMPTY.pipe(repeat()) liên tục complete rồi subscribe lại mà không nhường event loop. Trong RxJS 7.8.2, pattern này có thể gây tràn stack; nó cũng có thể khóa UI trước khi bạn có cơ hội gọi unsubscribe từ bên ngoài. Đặt takeUntil sau nó không khiến vòng lặp đồng bộ tự nhường thời gian cho sự kiện hủy.

Dùng giới hạn lượt hợp lý, hoặc delay dạng số/notifier bất đồng bộ khi thật sự cần lặp lâu dài. repeat({ delay: 0 }) dùng timer nên nhường lượt chạy, nhưng polling HTTP không nên chạy sát nhau như vậy: chọn khoảng chờ theo yêu cầu API và tải hệ thống.

Source hot đã complete hoặc dữ liệu đã cache

Subscribe lại một Subject đã complete chỉ nhận complete ngay, không khởi động producer lại. repeat() trên nguồn đó có thể trở thành vòng lặp đồng bộ. repeat không “mở lại” socket, DOM event source hay Subject đã đóng.

Tương tự, request$.pipe(shareReplay(1), repeat(3)) có thể chỉ đọc response được cache ba lượt sau khi source đã complete, thay vì tạo request mới. Với polling dùng chung cho nhiều consumer, thường nên xây dựng vòng polling từ request cold trước, rồi mới chọn cơ chế chia sẻ cho toàn vòng polling. Cấu hình ref counting quyết định upstream có dừng khi consumer cuối rời đi hay không.

Tác vụ ghi có side effect

Không repeat thanh toán, tạo đơn hàng hay tạo job chỉ vì stream đã complete. Mỗi lượt subscribe có thể tạo thêm một side effect thật; RxJS không cung cấp rollback hay idempotency tự động. Với polling job, repeat endpoint đọc trạng thái, không repeat toàn flow “tạo job rồi đọc trạng thái”.

Cần nhịp cố định theo đồng hồ

repeat({ delay }) phù hợp với nhịp “xong lượt trước, nghỉ, rồi làm lượt sau”. Nếu cần tick theo lịch cố định, bắt đầu từ timer hoặc interval, rồi chọn flattening operator có chủ đích: exhaustMap bỏ tick khi request còn chạy, switchMap thay request cũ, còn concatMap có thể tích hàng đợi. Đừng đổi sang interval mà bỏ qua quyết định về concurrency.

Kiểm thử dữ liệu và số lần subscribe

Chỉ assert các value lặp lại chưa đủ để chứng minh có request mới. Test thêm subscription upstream và trường hợp hủy trong khoảng chờ. Ví dụ sau dùng TestScheduler có sẵn trong RxJS, không cần timer thật:

import { strict as assert } from 'node:assert';
import { repeat, takeUntil } from 'rxjs';
import { TestScheduler } from 'rxjs/testing';

function createScheduler(): TestScheduler {
  return new TestScheduler((actual, expected) => {
    assert.deepEqual(actual, expected);
  });
}

createScheduler().run(({ cold, expectObservable, expectSubscriptions }) => {
  const source = cold('-a|');
  const result = source.pipe(repeat({ count: 3, delay: 2 }));

  expectObservable(result).toBe('-a---a---a|');
  expectSubscriptions(source.subscriptions).toBe([
    '^-!',
    '----^-!',
    '--------^-!',
  ]);
});

createScheduler().run(({ cold, hot, expectObservable, expectSubscriptions }) => {
  const source = cold('-a|');
  const stop = hot('---s');
  const result = source.pipe(
    repeat({ delay: 5 }),
    takeUntil(stop),
  );

  expectObservable(result).toBe('-a-|');
  // Dừng trong khoảng chờ: không có subscription source thứ hai.
  expectSubscriptions(source.subscriptions).toBe('^-!');
});

createScheduler().run(({ cold, expectObservable, expectSubscriptions }) => {
  const source = cold('-a#');
  const result = source.pipe(repeat(3));

  expectObservable(result).toBe('-a#');
  expectSubscriptions(source.subscriptions).toBe('^-!');
});

Trong run, mỗi dấu - tương ứng một millisecond virtual time. Test đầu complete lượt đầu ở frame 2, chờ hai frame rồi subscribe tiếp ở frame 4; test thứ hai dừng ở frame 3 trước khi notifier phát; test cuối chứng minh error không repeat.

Với ứng dụng thật, thêm test repeat(0) không gọi producer, count: 1 không gọi callback delay và notifier EMPTY kết thúc output. Spy adapter gọi HTTP để xác nhận số request; nhiều subscription vào một Promise cũ không đồng nghĩa nhiều lần gọi server.

Checklist trước khi dùng

  • Source có complete sau mỗi lượt hay sẽ sống mãi?
  • count đang tính tổng lượt, không phải số lần thử lại chứ?
  • Mỗi subscription có tạo công việc mới, hay đang dùng Promise/cache cũ?
  • Vòng lặp chỉ chứa tác vụ được phép chạy lại, không chứa thao tác ghi ngoài ý muốn?
  • Điều kiện dừng và takeUntil đã nằm ngoài phạm vi repeat chưa?
  • Có khoảng chờ và giới hạn/deadline nếu source cứ complete mà điều kiện nghiệp vụ chưa đạt?
  • Error cần dừng, retry hay trở thành fallback có ý nghĩa?
  • Producer có teardown để hủy công việc thật, và test có kiểm tra số lượt subscribe không?

Học tiếp

Nguồn tham khảo

On this page