Học RxJS
Higher-order Observables

concatMap

Xử lý tuần tự và giữ nguyên thứ tự công việc.

Giả sử bạn gửi ba thay đổi liên tiếp cho cùng một đơn hàng: thêm sản phẩm, đổi số lượng, rồi xác nhận. Nếu cả ba request chạy song song, request xác nhận có thể đến server trước khi số lượng được cập nhật. Bạn cần giữ mọi thao tác, nhưng thao tác sau chỉ được bắt đầu khi thao tác trước đã xong. Đó là tình huống mình chọn concatMap.

Phạm vi bài viết

Các ví dụ dùng RxJS 7.x và TypeScript, với import từ rxjs. Bạn nên biết pipe, subscribe và khái niệm outer Observable, inner Observable trong Higher-order stream. Các mốc thời gian bên dưới là số liệu minh họa, không phải benchmark.

Mục lục

Khi nào chọn concatMap

concatMap ánh xạ từng giá trị của source thành một inner Observable, rồi đưa các giá trị inner phát ra xuống output. Khác biệt nằm ở cách nó quản lý subscription: chỉ một inner hoạt động tại một thời điểm; inner tiếp theo phải chờ inner hiện tại complete.

Mình chọn concatMap khi thứ tự thao tác mang ý nghĩa nghiệp vụ và mỗi thao tác cần được xử lý. Ví dụ: gửi các lệnh cập nhật cùng một tài nguyên, upload từng phần theo thứ tự, hoặc xử lý một danh sách công việc hữu hạn mà bước sau không được vượt bước trước.

Nếu bạn chỉ cần kết quả mới nhất, như tìm kiếm theo từ khóa đang gõ, xếp hàng mọi từ khóa cũ lại là điều không nên làm. Với dữ liệu độc lập cần throughput cao, xử lý tuần tự cũng gây chờ không cần thiết. Phần so sánh operator giúp bạn chọn chính sách khác.

Hãy hình dung một quầy chỉ phục vụ một người mỗi lượt. Người đến sau vẫn được ghi nhận, nhưng phải chờ người trước rời quầy. Giới hạn của ví von này là hàng đợi RxJS nằm trong bộ nhớ của một execution: nó không có lưu trữ bền vững, không tự đặt giới hạn và không sống tiếp sau khi subscription bị hủy.

Cách concatMap vận hành

Hàng đợi và điều kiện complete

Trong RxJS 7.x, concatMap(project) tương đương mergeMap(project, 1). Khi source phát một giá trị và chưa có inner đang chạy, operator gọi project rồi subscribe vào kết quả. Nếu đang bận, nó giữ giá trị outer trong buffer theo thứ tự FIFO, thay vì gọi ngay project cho giá trị ấy.

Outer nhận A, B, C

Đang xử lý: A ── complete ──► B ── complete ──► C
Hàng đợi:  [B, C]            [C]               []

Có ba điểm cần tách biệt:

Sự kiệnHành vi
Inner phát nextChuyển giá trị xuống output ngay; chưa lấy công việc kế tiếp
Inner completeLấy giá trị đầu hàng đợi, gọi project và subscribe vào inner tiếp theo
Outer completeNgừng nhận việc mới, nhưng vẫn xử lý inner đang chạy và hàng đợi
Outer hoặc inner error không được xử lýOutput lỗi và đóng subscription; không tiếp tục hàng đợi
Consumer unsubscribeHủy subscription outer và inner đang chạy; không xử lý phần còn chờ

Output chỉ complete bình thường khi outer đã complete, không còn inner hoạt động và hàng đợi đã rỗng. Vì vậy, một source hữu hạn có thể complete rất sớm trong khi output vẫn đang làm việc.

Một inner cũng có thể phát nhiều giá trị hoặc không phát giá trị nào. concatMap không ép mỗi công việc có đúng một kết quả: nó chuyển tiếp toàn bộ next của inner hiện tại. EMPTY không phát gì nhưng complete ngay, nên lượt kế tiếp được bắt đầu.

Đọc timeline

Sơ đồ dưới đây đặt mỗi inner vào đúng frame mà concatMap thực sự subscribe. Bạn có thể đọc theo cột dọc để thấy tại một thời điểm chỉ có một inner hoạt động.

Marble timeline
concatMap xếp các inner nối đuôi nhau
Mỗi vạch = 1 frame
concatMap xếp các inner nối đuôi nhau. outer: -a-b-c-|. inner của a: -^--x|. inner của b: -----^--x|. inner của c: ---------^--x|. output: ----x---x---x|

Ký hiệu subscribe trên các dòng inner đánh dấu lúc bắt đầu theo dõi inner, không phải một giá trị output. Mỗi inner mất bốn frame từ lúc subscribe đến complete và phát x sau ba frame. Outer complete ở frame 7, nhưng c chỉ bắt đầu ở frame 9; output complete ở frame 13.

Bạn không phải đợi tất cả công việc xong mới nhận output. Kết quả của a được phát ở thời điểm 4, rồi b ở 8 và c ở 12. Operator giữ thứ tự bằng cách tuần tự hóa subscription, không phải chạy song song rồi sắp xếp lại kết quả.

API và ví dụ tối thiểu

Dạng API nên dùng trong RxJS 7.x:

concatMap((value, index) => inner$)
Thành phầnÝ nghĩa
valueGiá trị source đang được đưa vào xử lý
indexChỉ số bắt đầu từ 0 của lần gọi project trong execution này
Kết quả của projectMột ObservableInput, thường là Observable; cũng có thể là Promise, iterable hoặc array
OutputObservable phát các giá trị do inner phát ra, không phải các inner Observable

Không có tham số chỉnh concurrency cho concatMap. Nếu cần nhiều inner đồng thời, hãy dùng mergeMap với giới hạn thích hợp. Overload resultSelector cũ đã deprecated trong RxJS 7; nếu cần kết hợp outer value với inner result, đặt map trong inner pipeline.

Ví dụ sau mô phỏng ba công việc có thời gian khác nhau:

import { concatMap, defer, map, of, timer } from 'rxjs';

const jobs = [
  { id: 'A', durationMs: 300 },
  { id: 'B', durationMs: 100 },
  { id: 'C', durationMs: 200 },
];

of(...jobs).pipe(
  concatMap((job) =>
    defer(() => {
      console.log('start', job.id);
      return timer(job.durationMs).pipe(map(() => job.id));
    }),
  ),
).subscribe({
  next: (id) => console.log('done', id),
  error: (error: unknown) => console.error(error),
  complete: () => console.log('all done'),
});
start A
done A
start B
done B
start C
done C
all done

of(...jobs) phát cả ba giá trị và complete đồng bộ. Nhưng vì inner của A chưa complete, B và C nằm trong hàng đợi. B mất ít thời gian hơn A cũng không giúp B vượt lên trước: timer của B chưa được subscribe cho tới khi A xong.

defer giúp phần tạo công việc nằm rõ ở thời điểm subscribe. Riêng việc tạo timer vốn đã lazy, nên ví dụ vẫn tuần tự nếu bỏ defer; nó được dùng ở đây để log start cùng lúc inner thực sự bắt đầu.

Gửi thay đổi tuần tự và chụp snapshot

Một bẫy dễ gặp là enqueue một object mutable rồi đọc nó sau vài giây. concatMap giữ reference của giá trị outer, không tự sao chép dữ liệu. Nếu object bị sửa khi đang chờ, request có thể gửi trạng thái mới thay vì trạng thái tại lúc người dùng bấm lưu.

Ví dụ dưới đây chụp snapshot trước khi enqueue. Nó chạy trong trình duyệt có RxJS 7.x; endpoint /api/drafts là API minh họa, cần thay bằng endpoint thật. Ví dụ có một subscriber để mỗi thao tác chỉ được gửi qua một hàng đợi.

import {
  Subject,
  concatMap,
  finalize,
  map,
  takeUntil,
  timeout,
} from 'rxjs';
import { ajax } from 'rxjs/ajax';

type Draft = { title: string; body: string };
type SaveCommand = { commandId: string; draft: Draft };
type SavedDraft = { id: string; revision: number };

const saveClicks$ = new Subject<Draft>();
const destroy$ = new Subject<void>();

function saveDraft$(command: SaveCommand) {
  return ajax.post<SavedDraft>(
    '/api/drafts',
    command.draft,
    {
      'Content-Type': 'application/json',
      'Idempotency-Key': command.commandId,
    },
  ).pipe(
    timeout({ first: 5_000 }),
    map((response) => response.response),
  );
}

const subscription = saveClicks$.pipe(
  map((draft): SaveCommand => ({
    commandId: crypto.randomUUID(),
    draft: { ...draft },
  })),
  concatMap((command) => saveDraft$(command)),
  takeUntil(destroy$),
  finalize(() => console.log('save pipeline đã đóng')),
).subscribe({
  next: (saved) => console.log('saved revision', saved.revision),
  error: (error: unknown) => console.error('Dừng hàng đợi lưu', error),
});

const draft: Draft = { title: 'Bản đầu', body: 'Nội dung' };
saveClicks$.next(draft);
draft.title = 'Bản sửa';
saveClicks$.next(draft);

// Gọi khi owner của pipeline bị destroy.
function destroy(): void {
  destroy$.next();
  destroy$.complete();
  saveClicks$.complete();
}

map trước concatMap chạy khi sự kiện đến, nên hai command có hai snapshot khác nhau. project của command thứ hai chỉ chạy khi request đầu complete. Với Draft chỉ có các field string, shallow copy là đủ; nếu có object hoặc array lồng nhau, bạn cần snapshot sâu phù hợp, chẳng hạn structuredClone với dữ liệu được hỗ trợ.

ajax.post tạo cold Observable: request bắt đầu khi subscribe. timeout giới hạn thời gian chờ response đầu tiên; với request Ajax thông thường, inner phát response rồi complete. Nếu timeout xảy ra, pipeline này dừng toàn bộ hàng đợi vì chưa có chính sách phục hồi lỗi. Mình chọn cách dừng cho các thao tác phụ thuộc nhau: tiếp tục sau khi bước trước thất bại có thể tạo trạng thái sai.

Idempotency là hợp đồng với backend

Header Idempotency-Key chỉ hữu ích khi server thực sự hỗ trợ nó. concatMap không biến request thành transaction, không tự retry và không bảo đảm exactly-once. Đừng dùng tên header này như bằng chứng rằng backend đã chống trùng lặp.

Xử lý lỗi mà không vô tình bỏ hàng đợi

catchError bên trong và bên ngoài

Không phải mọi hàng đợi đều nên dừng khi một job lỗi. Ví dụ, upload từng file độc lập có thể ghi nhận file thất bại rồi tiếp tục file kế tiếp. Khi đó, đặt catchError bên trong inner và trả về một Observable hoàn tất được.

Ví dụ tự chạy dưới đây mô phỏng A thành công, B lỗi và C thành công:

import { catchError, concatMap, map, of, throwError, timer } from 'rxjs';

type Job = { id: string; shouldFail: boolean };
type Result =
  | { status: 'ok'; jobId: string }
  | { status: 'failed'; jobId: string; error: unknown };

function runJob$(job: Job) {
  return timer(10).pipe(
    concatMap(() =>
      job.shouldFail
        ? throwError(() => new Error(`Job ${job.id} lỗi`))
        : of(job.id),
    ),
  );
}

const jobs: Job[] = [
  { id: 'A', shouldFail: false },
  { id: 'B', shouldFail: true },
  { id: 'C', shouldFail: false },
];

of(...jobs).pipe(
  concatMap((job) =>
    runJob$(job).pipe(
      map((): Result => ({ status: 'ok', jobId: job.id })),
      catchError((error: unknown) =>
        of<Result>({ status: 'failed', jobId: job.id, error }),
      ),
    ),
  ),
).subscribe((result) => console.log(result.status, result.jobId));
ok A
failed B
ok C

Nếu không cần phát một kết quả thất bại, có thể trả về EMPTY sau khi ghi nhận lỗi. Nó complete ngay, nhường lượt cho job tiếp theo. Trả về NEVER thì làm ngược lại: inner không complete và hàng đợi bị kẹt.

Đặt catchError sau concatMap có ý nghĩa khác:

of(...jobs).pipe(
  concatMap((job) => runJob$(job)),
  catchError(() => of('pipeline thất bại')),
).subscribe(console.log);

Khi B lỗi, subscription của pipeline cũ đã bị đóng. catchError chuyển sang Observable thay thế và phát 'pipeline thất bại'; nó không khôi phục hàng đợi cũ để chạy C. Trong ví dụ này, bạn nhận A, rồi thông báo thất bại, rồi complete.

catchError bên trong cũng không bắt được lỗi outer hoặc lỗi xảy ra trước khi inner pipeline được tạo. Nếu hàm tạo công việc có thể throw đồng bộ, bọc nó bằng defer để biến exception thành error của inner rồi xử lý tại đó:

concatMap((job) =>
  defer(() => runJob$(job)).pipe(
    catchError((error: unknown) =>
      of({ status: 'failed', jobId: job.id, error }),
    ),
  ),
)

Đoạn này là mẫu operator, cần import defer, concatMap, catchError, of từ rxjs khi ghép vào pipeline.

Retry đúng phạm vi

Nếu retry từng công việc, đặt retry trong inner, trước catchError. Trong thời gian retry và delay, job hiện tại vẫn giữ lượt; các job khác chưa được bắt đầu.

// Áp dụng cho saveDraft$ và SaveCommand trong ví dụ lưu bản nháp.
// Cần import catchError, concatMap, defer, of, retry từ 'rxjs'.
concatMap((command: SaveCommand) =>
  defer(() => saveDraft$(command)).pipe(
    retry({ count: 2, delay: 1_000 }),
    catchError((error: unknown) =>
      of({ status: 'failed', commandId: command.commandId, error }),
    ),
  ),
)

count: 2 nghĩa là tối đa hai lần retry sau lần thử ban đầu. Mẫu trên chọn tiếp tục sau khi hết retry; nếu các command phụ thuộc nhau, bỏ catchError này để dừng pipeline và yêu cầu xử lý lỗi trước khi tiếp tục. Trong production, chỉ retry những lỗi có thể phục hồi, thay vì retry cả lỗi validation hay authorization.

Đặt retry bên ngoài concatMap sẽ resubscribe toàn bộ upstream. Với cold source như of(...jobs), những job đã thành công có thể chạy lại từ đầu. Với hot source như Subject, resubscribe không phát lại các sự kiện đã qua và không phục hồi buffer cũ.

Retry write request còn có rủi ro khác: client không nhận response không có nghĩa server chưa ghi dữ liệu. Nếu backend hỗ trợ idempotency, dùng cùng một command ID cho mọi lần thử của cùng thao tác; đừng tạo key mới trong mỗi lần defer chạy.

Complete và unsubscribe khác nhau

Giả sử bạn đang có A chạy, B và C chờ. Gọi complete() trên outer cho phép concatMap xử lý hết A, B, C rồi output mới complete. Gọi unsubscribe() trên subscription output lại dừng ngay phía RxJS: outer và inner đang chạy bị unsubscribe, B và C không được bắt đầu.

Vị trí của operator lifecycle vì thế thay đổi chính sách:

// Ngừng nhận job mới, nhưng xử lý hết inner và buffer hiện có.
source$.pipe(
  takeUntil(stopAccepting$),
  concatMap((job) => runJob$(job)),
);

// Hủy cả pipeline khi destroy$ phát giá trị.
source$.pipe(
  concatMap((job) => runJob$(job)),
  takeUntil(destroy$),
);

Đây là hai mẫu ghép pipeline; source$, notifier và runJob$ do ứng dụng cung cấp. takeUntil phản ứng khi notifier phát giá trị, không phải khi notifier chỉ complete. takeUntil ở sau concatMap hoàn tất output về phía consumer và teardown upstream; explicit unsubscribe() không gọi callback complete của consumer. Cả hai đều không xử lý tiếp hàng đợi.

Mình đặt lifecycle cancellation sau concatMap nếu công việc thuộc màn hình và phải dừng khi màn hình đóng. Nếu job cần sống lâu hơn màn hình, hãy chuyển ownership sang service hoặc worker có lifecycle thích hợp, không giữ subscription của component sống âm thầm.

Hủy subscription không hoàn tác thao tác trên server

Observable chỉ hủy được tài nguyên theo teardown của nó. Ajax có thể abort request đang chờ, nhưng server có thể đã ghi dữ liệu. Promise thường không bị hủy khi unsubscribe. Dùng finalize cho cleanup, không xem nó là thông báo rằng mọi job đã thành công.

Các bẫy khi dùng concatMap

Inner không complete

Một inner interval hoặc WebSocket sống dài hạn có thể phát rất nhiều giá trị nhưng không complete. Job kế tiếp vẫn phải chờ, vì điều kiện nhường lượt là complete chứ không phải next đầu tiên.

import { concatMap, interval, of, take } from 'rxjs';

of('A', 'B').pipe(
  concatMap((id) => interval(1_000).pipe(take(3))),
).subscribe(console.log);

Với take(3), mỗi lượt phát 0, 1, 2 rồi complete, nên B được chạy sau A. Bỏ take(3) thì chỉ inner của A hoạt động mãi. Biến id ở đây chỉ đại diện job, không được phát xuống output.

Chọn điều kiện kết thúc đúng nghiệp vụ: take(1) khi chỉ cần sự kiện đầu, takeUntil khi có tín hiệu kết thúc, hoặc timeout khi có SLA chờ. Nếu inner phát heartbeat liên tục nhưng không complete, timeout({ first: ... }) không giải quyết được việc nó chiếm lượt mãi: deadline đó chỉ kiểm soát thời gian chờ emission đầu tiên.

Hàng đợi tăng nhanh hơn tốc độ xử lý

concatMap không tạo backpressure cho source. Khi bận, nó tiếp tục nhận giá trị và buffer không có giới hạn cấu hình sẵn. Một nguồn sự kiện vô hạn nhanh hơn khả năng xử lý sẽ làm backlog, độ trễ và lượng bộ nhớ giữ lại tăng dần.

Giả sử mỗi giây có 10 job mới, nhưng mỗi job mất 200 ms. Một lượt xử lý chỉ làm được khoảng 5 job mỗi giây, nên backlog tăng khoảng 5 job mỗi giây trong mô hình minh họa này. Chỉ đổi từ mergeMap sang concatMap không làm nguồn phát chậm lại.

Chính sách cần đến từ ứng dụng:

Nhu cầuHướng xử lý
Chỉ quan tâm trạng thái mới nhấtCân nhắc switchMap, hoặc gộp sự kiện trước khi xử lý
Bỏ click mới khi đang bậnCân nhắc exhaustMap
Job độc lập, có thể chạy đồng thờiDùng mergeMap với concurrency phù hợp; vẫn cần kiểm soát backlog
Mọi job đều phải được xử lý, kể cả khi app restartDùng queue bền vững và cơ chế kiểm soát nhận việc ở cấp hệ thống
Các giá trị là thay đổi trạng thái có thể gộpCân nhắc debounceTime hoặc auditTime trước concatMap, chấp nhận mất các trạng thái trung gian

Đừng debounce các lệnh nghiệp vụ bắt buộc phải giữ mọi thao tác. Và giới hạn concurrency không đồng nghĩa giới hạn buffer: mergeMap cũng có thể giữ nhiều giá trị đang chờ.

Promise chạy sớm và hot Observable

concatMap tuần tự hóa subscription, không quay ngược thời gian để trì hoãn công việc đã bắt đầu:

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

// Các fetch đã bắt đầu ngay khi tạo array, trước concatMap.
const requests = [fetch('/api/a'), fetch('/api/b')];
from(requests).pipe(
  concatMap((request) => from(request)),
).subscribe();

// Tạo fetch khi đến lượt: chưa có request B khi A đang chạy.
of('/api/a', '/api/b').pipe(
  concatMap((url) => defer(() => fetch(url))),
).subscribe();

Đây là các snippet trình duyệt minh họa thời điểm khởi tạo; ứng dụng thật cần xử lý HTTP status, đọc body và xử lý lỗi. fetch resolve khi đã có response headers, không nhất thiết khi body đã đọc hết. Nếu “xong công việc” bao gồm đọc JSON, giữ cả bước đọc body trong inner:

concatMap((url: string) =>
  defer(async () => {
    const response = await fetch(url);
    if (!response.ok) {
      throw new Error(`HTTP ${response.status}`);
    }
    return response.json();
  }),
)

Mẫu này vẫn không tự abort fetch khi unsubscribe; muốn hủy request cần thiết kế teardown với AbortController hoặc dùng Observable HTTP có hỗ trợ cancellation.

Với hot Observable, subscription bắt đầu muộn có thể bỏ lỡ những emission đã xảy ra trước đó, trừ khi source có replay phù hợp. Bảo đảm chỉ một inner subscription hoạt động không có nghĩa chỉ một producer tồn tại. Vì vậy, nếu mục tiêu là tuần tự hóa side effect, hãy tạo cold công việc khi đến lượt thay vì enqueue các công việc đã chạy.

Tuần tự không đồng nghĩa với transaction

Một subscription giữ thứ tự các giá trị nó nhận được. Nó không bảo đảm thứ tự giữa hai tab, hai client hoặc hai subscription khác nhau. Subscribe hai lần vào pipeline cold thường tạo hai execution và hai hàng đợi riêng.

Ngay cả trong một execution, “request đã complete” cũng chỉ có ý nghĩa theo hợp đồng API. Nếu server trả 202 Accepted rồi xử lý nền, inner complete không chứng minh thao tác nghiệp vụ đã hoàn tất. Bạn cần chờ trạng thái hoàn tất thật sự nếu job kế tiếp phụ thuộc kết quả ấy.

Cuối cùng, concatMap không tự tạo thread hay chuyển công việc CPU nặng sang background. Nếu inner chạy đồng bộ, công việc tuần tự có thể vẫn chặn event loop. Dùng worker hoặc chiến lược scheduling phù hợp cho bài toán đó, thay vì coi operator là cơ chế xử lý song song.

So sánh các flattening operator

Cùng nhận một giá trị mới khi inner cũ còn chạy, bốn operator có bốn chính sách khác nhau:

OperatorLàm gì với giá trị mới?Trường hợp thường hợp lý
concatMapXếp hàng, chờ inner cũ completeMọi thao tác đều cần giữ, thứ tự quan trọng
mergeMapBắt đầu inner mới nếu còn slot; nếu hết slot thì bufferJob độc lập, muốn concurrency
switchMapUnsubscribe inner cũ rồi bắt đầu inner mớiChỉ kết quả mới nhất còn ý nghĩa
exhaustMapBỏ giá trị mới trong lúc inner cũ còn hoạt độngKhông muốn enqueue click trùng khi đang submit

switchMap hủy subscription cũ, không bảo đảm rollback một write đã đến server. Vì vậy, đừng dùng nó cho lệnh ghi chỉ vì muốn giao diện “luôn mới nhất” mà chưa xét semantics của backend.

concatMap(project) và mergeMap(project, 1) có cùng hành vi tuần tự trong RxJS 7.x. Mình dùng tên concatMap khi muốn người đọc nhận ngay ý định giữ thứ tự; dùng mergeMap khi concurrency là một lựa chọn thiết kế cần biểu diễn rõ.

Kiểm chứng bằng TestScheduler

Test theo virtual time giúp kiểm tra thứ tự mà không phải chờ timer thật. Ví dụ dưới đây dùng RxJS 7.x trong môi trường test Node; dùng assertion của test runner tương đương nếu bạn không dùng node:assert.

import { deepStrictEqual } from 'node:assert';
import { concatMap } from 'rxjs';
import { TestScheduler } from 'rxjs/testing';

const scheduler = new TestScheduler((actual, expected) => {
  deepStrictEqual(actual, expected);
});

scheduler.run(({ cold, expectObservable, expectSubscriptions }) => {
  const source$ = cold('-a-b-c-|');
  const inner$ = cold('---x|');
  const result$ = source$.pipe(concatMap(() => inner$));

  expectObservable(result$).toBe('----x---x---x|');
  expectSubscriptions(inner$.subscriptions).toBe([
    '-^---!',
    '-----^---!',
    '---------^---!',
  ]);
});

Trong run, một frame của các marble trên tương ứng 1 ms virtual time. ! đánh dấu unsubscribe. Inner thứ hai chỉ subscribe ở frame 5, đúng lúc inner đầu complete; inner thứ ba bắt đầu ở frame 9. Assertion subscription quan trọng vì output đúng thứ tự chưa chắc đã chứng minh các side effect không chạy đồng thời.

Bạn nên thêm các test theo chính sách thực tế:

  • Inner thứ hai lỗi: pipeline dừng hay phát kết quả lỗi rồi xử lý inner thứ ba?
  • Outer complete lúc còn backlog: output có đợi drain hết không?
  • Owner bị destroy: active inner có teardown và job đang chờ có bị bỏ không?
  • Inner không complete: job tiếp theo có đúng là chưa được subscribe?

Virtual time kiểm soát timer Observable như timer và delay trong run, nhưng không tự điều khiển native Promise hoặc request mạng thật. Để test thứ tự HTTP, dùng cold Observable giả lập; kiểm tra cancellation và hành vi backend bằng test tích hợp riêng.

Checklist trước khi dùng

  • Mọi giá trị thực sự cần giữ lại, thay vì chỉ quan tâm giá trị mới nhất?
  • Inner là công việc tạo khi đến lượt, không phải Promise hoặc producer đã chạy trước đó?
  • Inner complete ở đúng thời điểm nghiệp vụ coi là hoàn tất?
  • Dữ liệu mutable đã được snapshot trước khi enqueue?
  • Lỗi phải dừng hàng đợi hay ghi nhận và tiếp tục? catchError đã đặt đúng phạm vi?
  • Retry có thể lặp write không, và backend có hợp đồng idempotency không?
  • Lifecycle cần drain hay cancel? Vị trí takeUntil đã thể hiện đúng lựa chọn?
  • Tốc độ nhận job có thể vượt tốc độ xử lý lâu dài không, và chính sách backlog nằm ở đâu?

Nếu còn phân vân, hãy viết một marble test có ba job, một job chậm và một job lỗi. Chọn chính sách trước, rồi mới chọn operator: đó là cách tránh biến hàng đợi thành hành vi tình cờ.

Học tiếp

Nguồn tham khảo

On this page