Học RxJS
Higher-order Observables

exhaustMap

Bỏ qua yêu cầu mới khi tác vụ hiện tại chưa xong.

Bạn bấm nút gửi form, request chưa trả về, rồi bấm thêm hai lần. Nếu mỗi click đều tạo request, backend có thể nhận ba lần gửi cho cùng một thao tác. exhaustMap phù hợp khi bạn muốn giữ tác vụ đã bắt đầu và bỏ qua trigger mới cho đến khi tác vụ đó complete.

Phạm vi bài viết

Ví dụ dùng API của RxJS 7.8.2, import operator từ rxjs. Bạn nên biết pipe, subscribe và vòng đời next / error / complete. Outer Observable là stream trigger; inner Observable là công việc được tạo từ một trigger được chấp nhận.

Mục lục

Mental model của exhaustMap

exhaustMap(project) vừa ánh xạ một outer value thành inner Observable, vừa flatten những next của inner xuống output. Điểm khác biệt nằm ở quyết định có nhận outer value hay không: trong mỗi subscription, operator chỉ cho phép tối đa một inner đang active.

Hình dung một quầy chỉ nhận lượt mới khi đã xử lý xong lượt hiện tại, nhưng không phát số chờ. Người đến lúc quầy bận phải quay lại sau. Trong pipeline, emission bị bỏ qua cũng vậy: operator không lưu nó để phát lại. Điểm dừng của phép ví von là quầy thực tế có thể từ chối bằng lời; exhaustMap không tự phát thông báo “đã bỏ qua” cho UI.

Bận thì bỏ qua chứ không xếp hàng

Khi đang rảnh, operator gọi project(value, index) rồi subscribe vào inner được trả về. Trong lúc inner active, outer vẫn được theo dõi nhưng những next mới bị bỏ qua; projection không được gọi cho các value đó.

Rảnh ── outer next ──► gọi project và subscribe inner ──► Bận
  ▲                                                       │
  └──────────────── inner complete ───────────────────────┘

Bận + outer next → bỏ qua, không gọi project, không buffer
Bận + inner next → chuyển value xuống output, vẫn bận

Không có hàng đợi, cũng không có “giữ lại click cuối”. Khi inner complete, operator chỉ chấp nhận emission đến sau đó. Nếu không có emission mới, không có tác vụ tiếp theo.

Chỉ complete mới mở cửa cho trigger tiếp theo

Một inner có thể phát một value, nhiều value hoặc không có value nào. Số lượng next không quyết định lúc mở cửa; complete mới quyết định.

Ví dụ, một stream upload phát progress 10%, 50%, 100% nhưng vẫn chưa complete thì exhaustMap vẫn bận. Một inner EMPTY complete ngay mà không phát value nào thì lại mở cửa ngay.

Nếu inner error, mặc định cả output error và đóng, không phải quay về trạng thái rảnh để nhận click tiếp. Muốn tiếp tục lắng nghe sau một request thất bại, bạn cần xử lý lỗi trong inner; phần submit form bên dưới minh họa vị trí đó.

Cú pháp và projection

Chữ ký thường dùng, viết gọn từ API RxJS 7.8.2:

import type { ObservableInput, ObservedValueOf, OperatorFunction } from 'rxjs';

declare function exhaustMap<T, O extends ObservableInput<any>>(
  project: (value: T, index: number) => O,
): OperatorFunction<T, ObservedValueOf<O>>;
Thành phầnÝ nghĩa
valueOuter value được nhận khi không có inner active
indexBắt đầu từ 0, tăng cho mỗi lần gọi projection; value bị bỏ qua không làm tăng index
Giá trị trả về của projectMột ObservableInput, thường là Observable; cũng có thể là Promise, array hoặc iterable
OutputNhững value do inner được chấp nhận phát ra, không phải bản thân inner Observable

Không có tham số concurrency: giới hạn một inner active là chính sách cố định của operator. Nhưng mergeMap(project, 1) không tương đương, vì nó giữ các value đến trong lúc bận để xử lý sau, còn exhaustMap bỏ chúng.

Nếu cần giữ cả outer value và kết quả inner, đặt map bên trong:

import { exhaustMap, map, of, Subject } from 'rxjs';

const jobIds$ = new Subject<string>();

jobIds$.pipe(
  exhaustMap((jobId) =>
    of({ ok: true }).pipe(
      map((result) => ({ jobId, result })),
    ),
  ),
).subscribe(console.log);

jobIds$.next('job-1');
// { jobId: 'job-1', result: { ok: true } }

Overload resultSelector của RxJS 7 đã deprecated; mình chọn inner map vì dễ đọc và không phụ thuộc overload cũ. Inner of trong ví dụ này complete đồng bộ; nó chỉ minh họa cách kết hợp dữ liệu, chưa tạo khoảng thời gian bận.

Timeline của một lượt xử lý

Giả sử outer phát a, b, c, d. Mỗi inner được nhận phát x sau hai frame và complete sau năm frame. Đây là thời gian minh họa, không phải thời lượng request thực tế.

Marble timeline
exhaustMap bỏ trigger đến trong lúc đang bận
Mỗi vạch = 1 frame
exhaustMap bỏ trigger đến trong lúc đang bận. outer: a-b-c-----d---|. inner a: ^-x--|. inner d: ----------^-x--|. output: --x---------x--|

Dấu subscribe trong sơ đồ là thời điểm operator bắt đầu theo dõi inner; dấu complete kết thúc timeline của stream tương ứng. Diễn biến:

  1. Frame 0: nhận a, gọi projection và subscribe inner a.
  2. Frame 2: inner a phát x; outer b bị bỏ qua vì inner chưa complete.
  3. Frame 4: outer c cũng bị bỏ qua; không có inner b hay inner c được tạo.
  4. Frame 5: inner a complete, operator rảnh trở lại.
  5. Frame 10: nhận d và subscribe inner d.
  6. Frame 14: outer complete nhưng output chưa complete, vì inner d còn active.
  7. Frame 15: inner d complete; lúc này output mới complete.

Ở ví dụ này không có outer emission trùng thời điểm inner complete. Nếu hai notification trùng frame, thứ tự xử lý thực tế quyết định emission có được nhận hay không; đừng dựa vào một mốc thời gian bằng nhau để thiết kế nghiệp vụ.

Ví dụ tối giản với timer

Ví dụ này chạy trong môi trường đã có RxJS, không cần DOM:

import { exhaustMap, map, Subject, timer } from 'rxjs';

const triggers$ = new Subject<string>();

triggers$.pipe(
  exhaustMap((label, index) => {
    console.log('accepted:', label, index);
    return timer(100).pipe(map(() => `done: ${label}`));
  }),
).subscribe({
  next: (value) => console.log(value),
  complete: () => console.log('complete'),
});

triggers$.next('A');
triggers$.next('B');
triggers$.next('C');
triggers$.complete();

// Ngay lập tức: accepted: A 0
// Sau khoảng 100 ms: done: A
// Sau đó: complete

A, B, C đến liên tiếp trong cùng lượt chạy đồng bộ, trong khi inner của A đang chờ timer. Vì vậy chỉ A được nhận, không phụ thuộc vào độ chính xác của clock. Timer phát một lần rồi complete; outer đã complete trước đó nên output cũng complete sau timer.

Bạn sẽ không thấy accepted: B hay accepted: C. Đây là cách quan sát trực tiếp khác biệt giữa “không subscribe kết quả” và “không gọi projection”.

Submit form mà không gửi chồng request

Mình chọn exhaustMap cho submit khi yêu cầu là “lượt đầu đang xử lý thì không nhận thêm”, chứ không phải “mọi lần bấm đều là một lệnh phải được lưu”. Ví dụ dưới đây giả sử trang có các phần tử:

<form id="profile-form">
  <input name="displayName" required />
  <button id="save-button" type="submit">Lưu</button>
</form>
<p id="save-status" role="status"></p>

Đoạn TypeScript chạy trong browser sau khi DOM đã sẵn sàng. Endpoint POST /api/profile là API minh họa: nhận { displayName: string } và trả JSON { displayName: string }; bạn cần thay endpoint và validation theo ứng dụng.

import {
  catchError,
  defer,
  EMPTY,
  exhaustMap,
  finalize,
  fromEvent,
  Subject,
  takeUntil,
  tap,
  timeout,
} from 'rxjs';
import { ajax } from 'rxjs/ajax';

type SavedProfile = { displayName: string };

function mountProfileForm(): () => void {
  const form = document.querySelector<HTMLFormElement>('#profile-form');
  const button = document.querySelector<HTMLButtonElement>('#save-button');
  const status = document.querySelector<HTMLElement>('#save-status');

  if (!form || !button || !status) {
    throw new Error('Thiếu phần tử của profile form');
  }

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

  const subscription = fromEvent<SubmitEvent>(form, 'submit').pipe(
    // Luôn ngăn submit native, kể cả event sẽ bị exhaustMap bỏ qua.
    tap((event) => event.preventDefault()),
    exhaustMap(() =>
      defer(() => {
        const fields = new FormData(form);
        const displayName = String(fields.get('displayName') ?? '').trim();
        if (!displayName) {
          throw new Error('Tên hiển thị không được để trống');
        }

        button.disabled = true;
        status.textContent = 'Đang lưu…';

        return ajax.post<SavedProfile>(
          '/api/profile',
          { displayName },
          { 'Content-Type': 'application/json' },
        );
      }).pipe(
        timeout({ first: 10_000 }),
        tap(({ response }) => {
          status.textContent = `Đã lưu: ${response.displayName}`;
        }),
        catchError((error: unknown) => {
          status.textContent = error instanceof Error
            ? `Không lưu được: ${error.message}`
            : 'Không lưu được. Bạn có thể thử lại.';
          return EMPTY;
        }),
        finalize(() => {
          button.disabled = false;
        }),
      ),
    ),
    // Khi màn hình bị tháo, hủy cả listener lẫn request đang active.
    takeUntil(destroy$),
  ).subscribe();

  return () => {
    destroy$.next();
    destroy$.complete();
    subscription.unsubscribe();
  };
}

const dispose = mountProfileForm();
// Gọi dispose() trong lifecycle unmount/destroy của màn hình.

Tạo request trong projection

Request và việc đọc form đều nằm trong projection, vì bạn chỉ muốn làm chúng cho event được nhận. defer chạy factory khi inner được subscribe, đồng thời chuyển lỗi đồng bộ trong factory vào error channel của inner để catchError bên trong xử lý được.

Ngược lại, preventDefault() phải nằm trước exhaustMap. Nếu đặt nó trong projection, những submit event bị bỏ qua sẽ không bị ngăn hành vi submit native. Không phải mọi side effect đều cần đặt cùng một chỗ: xử lý event thuộc outer, tạo tác vụ thuộc inner.

Nút disabled giúp người dùng hiểu trạng thái, còn exhaustMap bảo vệ pipeline nếu event vẫn đến qua một đường khác. Không nên chỉ bỏ qua click một cách im lặng rồi để người dùng đoán ứng dụng có đang làm gì không.

Đặt catchError và finalize bên trong

catchError bên trong thay một inner lỗi bằng EMPTY. EMPTY complete ngay, nên lượt xử lý kết thúc và outer tiếp tục nhận các submit về sau. Lỗi không được coi là thành công: UI hiển thị thông báo lỗi, chỉ không để lỗi request giết listener.

Nếu đặt catchError(() => EMPTY) sau exhaustMap, lỗi request làm toàn bộ chuỗi trước nó bị unsubscribe; Observable thay thế complete và listener submit không còn hoạt động. UI có thể trông như đã phục hồi nhưng click tiếp theo không tạo request.

finalize bên trong chạy cho từng lượt được nhận, kể cả complete, error hoặc unsubscribe. Đặt nó sau exhaustMap chỉ theo dõi việc kết thúc toàn bộ subscription, không phải từng request. Và đừng đặt tap(() => button.disabled = true) trước exhaustMap: khi đang bận, outer vẫn chạy tap dù event sẽ bị bỏ qua ở bước sau.

Retry có thể gửi lại một thao tác đã thành công

exhaustMap không bảo đảm exactly-once hay chống trùng trên backend. Hai tab, hai subscription, một lần retry hoặc request timeout sau khi server đã ghi dữ liệu đều có thể tạo thao tác trùng. Với giao dịch quan trọng, cần idempotency key hoặc cơ chế chống trùng phía server; đừng thêm retry cho POST chỉ vì pipeline đã có exhaustMap.

Đặt teardown sau exhaustMap

takeUntil(destroy$) nằm sau exhaustMap để khi notifier phát next, downstream kết thúc và unsubscribe cả outer lẫn inner active. Nếu đặt nó trước, notifier chỉ khiến outer complete; exhaustMap vẫn chờ inner hiện tại complete.

Notifier phải phát một value. Chỉ gọi destroy$.complete() không kích hoạt takeUntil. Trong ví dụ, hàm dispose phát next trước rồi mới complete notifier; unsubscribe() là cleanup trực tiếp bổ sung và có thể gọi lặp lại an toàn.

Unsubscribe khỏi ajax đang chạy sẽ abort XHR phía client. Điều đó không phải rollback: server có thể đã nhận và xử lý request trước khi client hủy.

Vòng đời và các trường hợp biên

Outer complete không hủy inner đang chạy

complete của outer có nghĩa là “không có trigger mới nữa”, không phải “hủy công việc hiện tại”. Output chờ inner active kết thúc. Khi không có inner active, outer complete được chuyển xuống output ngay.

NotificationHành vi mặc định
Outer next, đang rảnhGọi projection và subscribe inner
Outer next, đang bậnBỏ value, không gọi projection
Inner nextPhát xuống output, vẫn giữ inner active
Inner completeMở cửa cho trigger mới; complete output nếu outer đã complete
Outer completeChờ inner active; nếu rảnh thì complete output ngay
Outer hoặc inner errorError output, teardown các subscription đang active
Projection throwError output nếu không được xử lý phù hợp
Downstream unsubscribeTeardown outer và inner, không phát complete cho observer

Muốn chặn lỗi đồng bộ khi tạo request bằng catchError trong inner, hãy dùng cấu trúc defer(() => createRequest()).pipe(catchError(...)). Nếu createRequest() throw ngay trong projection trước khi trả về inner, catchError nằm trên inner chưa được tạo sẽ không thể bắt lỗi đó.

Inner không complete và giới hạn chờ

interval, một stream WebSocket dài hạn hoặc NEVER có thể giữ cửa đóng vô hạn. Kể cả inner phát rất nhiều value, trigger mới vẫn không được nhận. Nếu outer complete trong lúc đó, output cũng không complete.

Khi một trigger chỉ cần một kết quả, giới hạn inner theo yêu cầu đó:

import { exhaustMap, interval, Subject, take } from 'rxjs';

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

refresh$.pipe(
  exhaustMap(() => interval(1000).pipe(take(1))),
).subscribe(console.log);

refresh$.next();
// Khoảng 1 giây sau: 0; inner complete và cửa mở lại.

take(1) kết thúc sau value đầu tiên, nhưng không cứu được một inner không bao giờ phát. Với request một response như ajax ở ví dụ submit, timeout({ first: 10_000 }) giới hạn thời gian chờ response đầu tiên; con số 10 giây chỉ là cấu hình minh họa.

Với inner phát progress liên tục, timeout({ first: ... }) không giới hạn tổng thời gian tác vụ. timeout({ each: ... }) kiểm tra khoảng chờ giữa các emission, cũng không phải deadline tổng. Nếu nghiệp vụ cần deadline cố định, thiết kế riêng cơ chế kết thúc theo timer và xác định rõ timeout sẽ báo lỗi hay chỉ complete; nếu complete im lặng, UI không được nhầm nó thành thành công.

Inner đồng bộ có thể nhận mọi trigger

exhaustMap không đồng nghĩa với “chỉ lấy value đầu tiên”:

import { exhaustMap, of } from 'rxjs';

of('A', 'B', 'C').pipe(
  exhaustMap((value, index) => of(`${index}: ${value}`)),
).subscribe(console.log);

// 0: A
// 1: B
// 2: C

Mỗi inner of complete trước khi outer phát value kế tiếp. Operator đã rảnh trở lại nên nhận cả ba value. Khoảng bận được quyết định bởi lifecycle của inner, không phải khoảng cố định như throttleTime.

Nếu muốn chờ người dùng ngừng gõ thì dùng debounceTime; nếu muốn giới hạn tần suất theo clock thì xét throttle hoặc audit. Đừng dùng exhaustMap như một tên khác của debounce.

Promise và side effect tạo quá sớm

Promise được hỗ trợ, nhưng việc tạo Promise có thể bắt đầu công việc ngay. Đặt side effect ở upstream sẽ làm công việc chạy trước khi exhaustMap quyết định bỏ qua:

import { exhaustMap, from, map, Subject } from 'rxjs';

const urls$ = new Subject<string>();

// Không dùng cấu trúc này nếu mục tiêu là ngăn fetch chồng nhau.
urls$.pipe(
  map((url) => fetch(url)), // Mỗi URL đã bắt đầu một fetch.
  exhaustMap((promise) => from(promise)),
).subscribe();

Phiên bản đúng vị trí tạo công việc:

import { defer, exhaustMap, Subject } from 'rxjs';

const urls$ = new Subject<string>();

urls$.pipe(
  exhaustMap((url) => defer(() => fetch(url))),
).subscribe();

Bây giờ URL bị bỏ qua không gọi fetch. Tuy nhiên, unsubscribe khỏi Observable bọc Promise không tự cancel Promise hoặc abort fetch. Muốn hủy request thực sự, bạn cần tích hợp AbortController hoặc dùng API Observable có teardown phù hợp.

Một chi tiết khác: Promise của fetch resolve khi response headers đã sẵn sàng, không phải khi bạn xử lý xong body. Nếu muốn giữ cửa đóng đến khi parse JSON xong, hãy đưa cả bước đó vào inner, chẳng hạn defer(() => fetch(url).then(response => response.json())), và xử lý HTTP status theo yêu cầu ứng dụng.

Mỗi subscription có trạng thái bận riêng

Hai subscriber vào cùng một pipeline chứa exhaustMap có hai bộ trạng thái inner riêng. Với outer hot như Subject, một trigger có thể làm cả hai subscription gọi projection và tạo hai request.

Nếu nhiều consumer cần cùng một execution, có thể đặt share() sau exhaustMap, tùy lifecycle và nhu cầu replay. Đặt share() chỉ trước exhaustMap chia sẻ outer, không chia sẻ quyết định nhận trigger hay execution của inner.

Dù chia sẻ execution, chính sách này vẫn thuộc một pipeline trong một ứng dụng; nó không khóa tài nguyên chung giữa tab, máy hoặc người dùng khác nhau.

Chọn giữa bốn flattening operator

Giả sử a đã bắt đầu một inner và b đến trước khi inner đó complete:

OperatorXử lý bKhi nào mình chọn?
exhaustMapBỏ qua b, giữ inner aSubmit hoặc refresh mà trigger mới trong lúc bận là dư thừa
switchMapUnsubscribe inner a, chạy inner bSearch theo input; kết quả mới nhất thay thế kết quả cũ
concatMapGiữ b trong hàng đợi, chạy sau aCác lệnh phải được xử lý đủ và theo thứ tự
mergeMapChạy b song song nếu còn slot concurrency; nếu không thì bufferCông việc độc lập, cho phép xử lý đồng thời

Mình chọn exhaustMap khi bỏ một trigger là hành vi mong muốn, không phải khi muốn che request chậm. Nếu người dùng sửa form trong lúc lưu và lần bấm tiếp theo phải lưu phiên bản mới, bỏ lần bấm ấy có thể làm mất ý định của họ. Khi đó hãy khóa việc sửa form, thông báo rõ cần gửi lại sau, hoặc chọn chiến lược khác.

Với autosave, “luôn giữ bản mới nhất” thường gần với switchMap, nhưng hủy subscription của client không hoàn tác một write đã tới server; backend vẫn cần xử lý ordering hoặc version. Nếu yêu cầu là “đang chạy thì nhớ đúng một bản mới nhất để chạy tiếp”, cả exhaustMap lẫn concatMap mặc định đều không mô tả đúng chính sách đó: một bên bỏ hết, bên kia giữ tất cả.

Polling là trường hợp hợp lý khác: interval(...).pipe(exhaustMap(() => request$)) bỏ tick khi request cũ còn chạy. Lịch vẫn theo clock của outer, không phải “đợi một khoảng cố định sau mỗi response”; với yêu cầu thứ hai, hãy thiết kế polling theo completion.

Kiểm chứng bằng marble test

Test sau dùng TestScheduler của RxJS và assertion của Node.js. Nó kiểm tra output, số lần gọi projection và thời gian subscription của inner, không chỉ kiểm tra value cuối cùng.

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

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

scheduler.run(({ cold, expectObservable, expectSubscriptions, flush }) => {
  const outer = cold('a-b-c-----d---|');
  const inner = cold('--x--|');
  const accepted: Array<{ value: string; index: number }> = [];

  const result = outer.pipe(
    exhaustMap((value, index) => {
      accepted.push({ value, index });
      return inner;
    }),
  );

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

  flush();
  deepStrictEqual(accepted, [
    { value: 'a', index: 0 },
    { value: 'd', index: 1 },
  ]);
});

Trong run, mỗi dấu - biểu diễn một frame thời gian ảo; ^ là subscribe, ! là unsubscribe. Inner chỉ có hai subscription: tại frame 0 và 10. Nếu operator gọi projection cho b hoặc c, assertion accepted sẽ fail ngay cả khi output tình cờ vẫn trông đúng.

Khi áp dụng vào ứng dụng, mình bổ sung test cho ba ranh giới:

  • Request đầu lỗi, lỗi được xử lý bên trong, rồi trigger sau vẫn tạo được request.
  • Destroy xảy ra khi request active: inner bị unsubscribe ngay, không chỉ outer complete.
  • Inner không complete: trigger sau bị bỏ qua cho tới khi có giới hạn chờ hoặc teardown.

Checklist trước khi dùng

  • Trigger đến trong lúc bận có thể bỏ được, không cần queue hay replay.
  • Inner có điểm kết thúc rõ ràng; next không bị nhầm với complete.
  • Request được tạo trong projection, không khởi động trước ở map hoặc tap upstream.
  • Lỗi của từng tác vụ được xử lý bên trong nếu outer cần tiếp tục sống.
  • UI cho biết đang bận và phản ánh đúng thành công, lỗi hoặc hủy.
  • Teardown nằm sau flattening operator nếu cần hủy inner đang chạy.
  • Không có subscription thứ hai vô tình tạo thêm execution.
  • Backend có cơ chế chống trùng nếu thao tác không được phép ghi hai lần.

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

  1. Trong marble test, nếu inner đổi thành --x và không complete, d có được nhận không? Không. Output vẫn có x từ inner đầu, nhưng inner active mãi và outer complete không làm output complete.
  2. Nếu inner là of('ok') và outer là of(1, 2, 3), output có bao nhiêu value? Ba, vì mỗi inner complete đồng bộ trước outer value tiếp theo.
  3. Trong ví dụ form, nếu chuyển catchError(() => EMPTY) ra sau exhaustMap, sau request lỗi bạn còn submit được qua listener đó không? Không. Output chuyển sang EMPTY rồi complete, không resubscribe listener cũ.
  4. Nếu đổi exhaustMap thành concatMap, các trigger dư thừa lúc bận biến mất không? Không. Chúng được giữ để xử lý sau; chỉ nên đổi khi việc xử lý đủ mọi trigger là chủ ý.

Hãy sửa marble test để inner đầu lỗi còn inner tiếp theo thành công, rồi thêm catchError bên trong projection. Khi test cho thấy trigger sau vẫn được nhận, bạn đã kiểm chứng được chính sách phục hồi thay vì chỉ quan sát happy path.

Học tiếp

Nguồn tham khảo

On this page