Học RxJS
Higher-order Observables

mergeMap

Chạy song song và kiểm soát concurrency.

Bạn có một danh sách ID cần tải dữ liệu: chờ từng request thì chậm, nhưng chỉ giữ request mới nhất lại làm mất kết quả của những ID trước. mergeMap phù hợp khi các công việc độc lập, đều cần được xử lý và có thể chạy đồng thời. Bạn quyết định bao nhiêu công việc được phép chạy; operator quản lý inner subscription và đưa kết quả về cùng một stream.

Phạm vi và kiến thức nền

Bài này dùng API của RxJS 7.8.2, TypeScript và import operator từ rxjs. Bạn nên biết pipe, subscribe và khái niệm outer stream, inner stream. Các ví dụ timer chạy được trong playground có RxJS; ví dụ HTTP dùng endpoint minh họa của ứng dụng, không phải API công khai.

Mục lục

mergeMap giải quyết bài toán gì

Giả sử màn hình quản trị cần tải thông tin của nhiều sản phẩm. Request cho sản phẩm A không phụ thuộc request cho B; kết quả nào về trước có thể hiển thị trước. Nếu dùng concatMap, B phải đợi A complete. Nếu dùng switchMap, ID mới làm subscription của request trước bị hủy. Cả hai đều không khớp yêu cầu này.

mergeMap gọi một projection để đổi từng outer value thành inner Observable, subscribe các inner theo giới hạn concurrency rồi chuyển tiếp mọi inner emission xuống downstream. Với Promise và các kiểu ObservableInput khác, RxJS cũng chuyển chúng thành luồng để subscribe.

Mình chọn mergeMap khi công việc độc lập và không cần giữ thứ tự output. Với batch HTTP, mình thường đặt giới hạn concurrency rõ ràng thay vì dùng mặc định Infinity; giá trị cụ thể phải dựa trên giới hạn API và tải đo được, không phải một con số dùng chung cho mọi hệ thống.

Chạy đồng thời không có nghĩa là thêm thread

mergeMap cho phép nhiều inner subscription cùng hoạt động. Nó không tạo worker hay chuyển callback JavaScript sang thread khác. Nhiều request I/O có thể chờ đồng thời, nhưng một projection tính toán CPU nặng vẫn có thể chặn event loop.

Cơ chế outer và inner subscription

Giá trị mới không hủy inner cũ

Mỗi outer value được nhận khi còn slot sẽ đi qua projection và tạo một inner subscription. Những subscription trước đó vẫn sống; mergeMap không hủy chúng chỉ vì outer vừa phát value mới.

outer value A ──► project(A) ──► inner A ──┐
                                         ├──► output chung
outer value B ──► project(B) ──► inner B ──┘

A và B có thể cùng hoạt động nếu concurrency cho phép.

Hình dung quầy có hai nhân viên nhận việc: mỗi người giữ một việc cho tới khi xong, ai xong trước thì trả kết quả trước rồi nhận việc tiếp theo. Tuy nhiên, các inner không phải hai thread JavaScript; đây chỉ là cách hình dung số công việc đang hoạt động.

Một inner có thể phát không value, một value hoặc nhiều value. mergeMap chuyển tiếp từng next, chứ không gom toàn bộ kết quả của inner thành một item. Slot chỉ được trả khi inner complete, không phải ngay khi inner phát value đầu tiên.

Output theo thời điểm phát chứ không theo input

Ví dụ này dùng thời gian giả định để làm rõ thứ tự, không phải benchmark:

import { from, map, mergeMap, timer } from 'rxjs';

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

from(jobs).pipe(
  mergeMap((job) =>
    timer(job.durationMs).pipe(map(() => job.id)),
  ),
).subscribe({
  next: (id) => console.log(id),
  complete: () => console.log('complete'),
});

// Xấp xỉ:
// 100 ms: B
// 200 ms: C
// 300 ms: A
//         complete

from(jobs) phát cả ba input đồng bộ. Vì không đặt giới hạn, ba timer được subscribe gần như cùng lúc; B phát trước vì timer của B ngắn nhất. Outer đã complete từ lúc phát xong array, nhưng output còn chờ cả ba inner.

Mốc minh họa (ms)     0        100        200        300
outer                A,B,C,complete
inner A              subscribe ───────────────────► A,complete
inner B              subscribe ─► B,complete
inner C              subscribe ──────────► C,complete
output                         B          C          A,complete

Với inner nhiều emission, output có thể xen kẽ, chẳng hạn A1, B1, B2, A2. Vì vậy, “kết quả nào xong trước về trước” chỉ là cách nói ngắn cho request phát một kết quả; quy tắc tổng quát là inner nào phát value thì value ấy đi xuống ngay.

Cú pháp và tham số

Dạng nên dùng trong code mới:

mergeMap(project, concurrent?)
Thành phầnÝ nghĩaLưu ý
project(value, index)Đổi outer value thành ObservableInputKhông trả Subscription hoặc undefined
valueValue nhận từ outer sourceCó thể được giữ trong buffer trước khi projection được gọi
indexChỉ số projection, bắt đầu từ 0 cho mỗi subscriptionLà thứ tự xử lý input, không phải thứ tự output
concurrentSố inner subscription hoạt động tối đaMặc định Infinity; dùng số nguyên dương hữu hạn khi cần giới hạn
OutputStream của value do các inner phátKhông tự gom array, không tự sắp lại theo input

Ví dụ mergeMap((id) => loadProduct$(id), 4) cho phép tối đa bốn inner đang hoạt động trong mỗi subscription vào pipeline. Nếu có hai subscriber độc lập vào một cold pipeline, giới hạn không tự trở thành bốn request cho cả ứng dụng.

Trong RxJS 7.8.2, concurrent không được validate để báo lỗi cho mọi input không hợp lệ. Đừng dùng 0, số âm hoặc NaN như công tắc pause: value có thể bị giữ trong buffer mà không có inner nào được bắt đầu, khiến stream không complete. Nếu nhận cấu hình bên ngoài, validate thành số nguyên dương trước khi tạo pipeline.

RxJS 7 còn có overload resultSelector, nhưng overload đó đã deprecated. Để ghép outer metadata với inner result, dùng map bên trong inner:

import { map, mergeMap, of } from 'rxjs';

of('p1', 'p2').pipe(
  mergeMap((id, inputIndex) =>
    of({ name: `Product ${id}` }).pipe(
      map((product) => ({ id, inputIndex, product })),
    ),
    2,
  ),
).subscribe(console.log);

Không gọi .subscribe() trong projection. subscribe() trả Subscription, không phải luồng kết quả; làm vậy còn tách inner khỏi ownership của pipeline và khiến cancellation khó kiểm soát.

Giới hạn concurrency

Khi hết slot thì value chờ trong buffer

Đặt concurrent = 2 cho ba job ở ví dụ trước:

import { defer, from, map, mergeMap, timer } from 'rxjs';

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

from(jobs).pipe(
  mergeMap(
    (job) => defer(() => {
      console.log('start', job.id);
      return timer(job.durationMs).pipe(map(() => job.id));
    }),
    2,
  ),
).subscribe((id) => console.log('result', id));

Ngay lúc subscribe, A và B bắt đầu; C chờ trong buffer. Khoảng 100 ms, B phát kết quả rồi complete, giải phóng một slot để C bắt đầu. A và C đều dự kiến phát gần mốc 300 ms, nên không nên dựa vào thứ tự tương đối của hai kết quả cùng mốc trong ví dụ timer thực.

0 ms      active: A, B       buffer: C
100 ms    B complete        C bắt đầu; active: A, C
~300 ms   A và C complete   buffer rỗng; output complete

Buffer giữ outer value, và projection cho value đang chờ chưa được gọi. Trong RxJS 7.8.2, các value chờ được lấy theo FIFO khi có slot. Thứ tự bắt đầu này không đảm bảo thứ tự phát kết quả.

mergeMap(project, 1) có cùng chiến lược tuần tự như concatMap(project). Mình dùng concatMap khi yêu cầu nghiệp vụ là tuần tự, vì tên operator thể hiện ý định rõ hơn. Nếu inner complete đồng bộ như of(...), một slot được trả ngay nên bạn cũng sẽ không thấy nhiều inner chồng lấp chỉ vì đã chọn mergeMap.

Concurrency không phải rate limit hay backpressure

Giới hạn bốn inner không có nghĩa là tối đa bốn request mỗi giây. Nếu mỗi request complete nhanh, cả bốn slot có thể liên tục nhận việc mới trong cùng một giây. Rate limit cần chính sách theo thời gian, quota của API hoặc cơ chế điều phối riêng.

Còn source có chậm lại khi hết slot không? Không. mergeMap vẫn nhận outer value và giữ value dư trong buffer. Với source liên tục nhanh hơn khả năng xử lý, buffer có thể tăng không giới hạn. Đây không phải backpressure làm producer giảm tốc.

Với array hữu hạn nhỏ, buffer thường là đánh đổi chấp nhận được. Với stream dài hạn, hãy quyết định rõ: bỏ bớt event nếu nghiệp vụ cho phép, gộp thành batch nếu API hỗ trợ, chia công việc thành từng đợt hữu hạn, hoặc chuyển sang queue bền vững có kiểm soát dung lượng. bufferCount có thể giúp batching, nhưng tự nó cũng không buộc producer chờ consumer.

mergeMap không phải durable job queue

Buffer nằm trong execution hiện tại, không được lưu bền. Unsubscribe, lỗi không được xử lý hoặc đóng ứng dụng có thể bỏ công việc đang chờ. Nếu mỗi job bắt buộc phải được hoàn tất, bạn cần persistence, retry policy và idempotency ở cấp hệ thống, không chỉ một flattening operator.

Ví dụ tải batch với lỗi riêng từng item

Giả sử backend có GET /api/products/:id, trả một object sản phẩm. Mục tiêu là tải tối đa ba sản phẩm đồng thời, hiển thị mỗi kết quả ngay khi có và không để một ID lỗi làm mất cả batch.

import {
  catchError,
  defer,
  from,
  map,
  mergeMap,
  of,
  timeout,
} from 'rxjs';
import { ajax } from 'rxjs/ajax';

type Product = {
  id: string;
  name: string;
};

type ProductResult =
  | { status: 'ok'; id: string; inputIndex: number; product: Product }
  | { status: 'error'; id: string; inputIndex: number; message: string };

const ids = ['p1', 'p2', 'p3', 'p4', 'p5'];

const results$ = from(ids).pipe(
  mergeMap(
    (id, inputIndex) => defer(() =>
      ajax.getJSON<Product>(`/api/products/${encodeURIComponent(id)}`),
    ).pipe(
      timeout({ first: 5_000 }),
      map((product): ProductResult => ({
        status: 'ok',
        id,
        inputIndex,
        product,
      })),
      catchError((error: unknown) => of<ProductResult>({
        status: 'error',
        id,
        inputIndex,
        message: error instanceof Error ? error.message : String(error),
      })),
    ),
    3,
  ),
);

const subscription = results$.subscribe({
  next: (result) => {
    if (result.status === 'ok') {
      console.log('loaded', result.id, result.product.name);
    } else {
      console.warn('failed', result.id, result.message);
    }
  },
  error: (error: unknown) => console.error('batch failed', error),
  complete: () => console.log('batch complete'),
});

// Khi owner của batch bị hủy:
// subscription.unsubscribe();

defer giữ cả thao tác tạo request trong lifecycle của inner: chỉ khi inner được subscribe thì factory mới chạy. Nếu factory ném lỗi đồng bộ, lỗi đó cũng vào inner error channel và được catchError bên trong xử lý. Với ajax.getJSON vốn đã lazy, defer không bắt buộc chỉ để trì hoãn network request, nhưng giúp thể hiện ranh giới này rõ ràng.

timeout đặt deadline cho response đầu tiên; 5_000 ms chỉ là cấu hình minh họa. Request bị timeout sẽ được thay bằng result có status: 'error'. Fallback of(...) phát một value rồi complete, nên slot được trả để batch tiếp tục xử lý ID kế tiếp.

Product chỉ là kiểu TypeScript, không validate JSON lúc runtime. Nếu API không đáng tin về schema, thêm validation trong inner trước khi tạo success result. Callback error cuối pipeline vẫn cần cho lỗi ngoài phạm vi recovery từng item, chẳng hạn lỗi từ outer source.

Giữ metadata và thứ tự khi gom kết quả

Với UI, gắn id vào result để cập nhật đúng dòng; đừng lấy vị trí emission làm vị trí sản phẩm. Nếu cần chờ hết batch rồi xuất array theo thứ tự input, nối đoạn dưới vào ví dụ trên:

import { map, toArray } from 'rxjs';

const orderedResults$ = results$.pipe(
  toArray(),
  map((results) =>
    [...results].sort((a, b) => a.inputIndex - b.inputIndex),
  ),
);

orderedResults$.subscribe((results) => console.log(results));

Đây là alternative consumer: khi thử orderedResults$, hãy bỏ subscription trực tiếp vào results$ nếu không muốn chạy batch lần thứ hai. results$ là cold pipeline, nên hai subscription độc lập tạo hai lượt request.

Sắp xếp này giữ thứ tự trình bày sau khi mọi công việc đã xong, không giữ thứ tự thực thi ở server. toArray() cũng giữ mọi result trong bộ nhớ và chỉ phát khi output complete; vì vậy chỉ dùng cách này cho batch hữu hạn có kích thước kiểm soát được. Nếu các lệnh ghi phải tác động tới cùng resource theo đúng thứ tự, đổi sang concatMap thay vì sắp lại response.

Đặt catchError ở đúng cấp

Mặc định, lỗi từ outer source, exception trong projection hoặc lỗi từ một inner sẽ làm output error. RxJS đóng subscription của pipeline, unsubscribe các inner còn hoạt động và không xử lý tiếp các outer value đang chờ.

So sánh hai vị trí bằng ví dụ không cần HTTP:

import { catchError, from, mergeMap, of, throwError } from 'rxjs';

function task$(id: number) {
  return id === 2
    ? throwError(() => new Error('task 2 failed'))
    : of(`ok:${id}`);
}

// Recovery từng task: outer vẫn tiếp tục.
from([1, 2, 3]).pipe(
  mergeMap((id) => task$(id).pipe(
    catchError(() => of(`failed:${id}`)),
  )),
).subscribe(console.log);
// ok:1
// failed:2
// ok:3

// Recovery cả pipeline: thay stream lỗi bằng fallback rồi kết thúc.
from([1, 2, 3]).pipe(
  mergeMap((id) => task$(id)),
  catchError(() => of('batch failed')),
).subscribe(console.log);
// ok:1
// batch failed

catchError bên ngoài không tiếp tục outer cũ sau khi fallback complete. Nó thay toàn bộ stream upstream đã lỗi bằng stream khác. Khi muốn một task thất bại nhưng batch vẫn chạy, mình đặt recovery bên trong inner và trả result có trạng thái lỗi để downstream không phải đoán item nào đã biến mất.

Nếu cố ý bỏ item lỗi, inner catchError(() => EMPTY) sẽ complete mà không phát result, cũng trả slot. Ngược lại, thay bằng NEVER giữ inner sống mãi và chiếm slot mãi. Điều đó có thể làm batch không complete hoặc làm tất cả job chờ bị kẹt nếu hết slot.

Đặt retry trong inner khi muốn retry riêng request lỗi. retry ngoài mergeMap có thể subscribe lại toàn bộ outer và chạy lại cả job đã thành công. Với thao tác ghi, hãy kiểm tra idempotency trước khi retry; client bị lỗi không chứng minh server chưa ghi dữ liệu.

Completion và cancellation

Outer complete chưa chắc output complete

Để mergeMap complete bình thường, cả ba điều kiện phải đúng:

  1. Outer source đã complete.
  2. Không còn outer value chờ trong buffer.
  3. Mọi inner đang hoạt động đã complete.

Vì vậy, of('A', 'B') complete ngay không làm hai request của A và B bị hủy. Ngược lại, một inner như interval(1_000) không có giới hạn sẽ giữ output sống, ngay cả khi outer là array hữu hạn. Nếu inner đại diện cho task hữu hạn, dùng điều kiện complete phù hợp, chẳng hạn take(1) hoặc deadline với timeout.

Sự kiệnInner đang chạyBufferDownstream
Outer completeTiếp tục chạyTiếp tục xử lý khi có slotComplete sau khi mọi việc hoàn tất
Lỗi chưa được recoveryBị unsubscribeKhông xử lý tiếpNhận error, không nhận complete
Consumer gọi unsubscribe()Bị unsubscribeKhông xử lý tiếpKhông nhận terminal notification

finalize phù hợp để cleanup vì chạy cả khi complete, error hoặc unsubscribe. Nhưng finalize không phải tín hiệu task thành công; nếu cần trạng thái nghiệp vụ, tạo success/error result rõ ràng như ví dụ batch.

Đặt takeUntil sau mergeMap để hủy toàn bộ

Nếu owner bị destroy và bạn muốn dừng cả outer lẫn mọi inner, đặt takeUntil sau flattening operator:

import {
  from,
  map,
  mergeMap,
  Subject,
  takeUntil,
  timer,
} from 'rxjs';

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

from(['A', 'B', 'C']).pipe(
  mergeMap((id) =>
    timer(1_000).pipe(map(() => id)),
    2,
  ),
  takeUntil(destroy$),
).subscribe({
  next: (id) => console.log(id),
  complete: () => console.log('owner stopped'),
});

setTimeout(() => {
  destroy$.next();
  destroy$.complete();
}, 100);

// Khoảng 100 ms: owner stopped
// Không có kết quả A/B; C đang chờ sẽ không được bắt đầu.

Notifier phải phát value: chỉ gọi destroy$.complete() không kích hoạt takeUntil. Trong ví dụ này, takeUntil complete downstream và unsubscribe upstream, nên các timer đang chạy được dọn.

Nếu đặt takeUntil(destroy$) trước mergeMap, tín hiệu chỉ làm outer complete. mergeMap vẫn có thể tiếp tục inner và xử lý buffer đã nhận; với outer đồng bộ đã complete như from(array), tín hiệu còn đến sau khi outer đã xong. Đó là hành vi có ích khi muốn ngừng nhận việc mới nhưng cho batch cũ chạy hết, không phải cách hủy toàn bộ.

Unsubscribe không bảo đảm server rollback

RxJS unsubscribe inner; công việc bên dưới chỉ dừng nếu producer có teardown tương ứng. ajax có thể abort request đang chờ, timer được dọn, nhưng Promise thông thường không bị cancel. Ngay cả khi client abort HTTP, server vẫn có thể đã nhận và xử lý lệnh ghi. Cancellation ở client không phải transaction rollback.

Promise và thời điểm bắt đầu công việc

mergeMap nhận Promise trực tiếp và phát giá trị resolve của Promise. Tuy nhiên, Promise thường đã bắt đầu công việc ngay khi được tạo, không chờ RxJS subscribe.

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

const urls = ['/api/products/p1', '/api/products/p2', '/api/products/p3'];

// Sai nếu kỳ vọng concurrency = 2 sẽ giới hạn số fetch bắt đầu.
const requests = urls.map((url) => fetch(url));
// Cả ba fetch đã bắt đầu ở đây.

const eagerResults$ = from(requests).pipe(
  mergeMap((request) => request, 2),
);

// Tạo fetch bên trong inner: chỉ bắt đầu khi có slot.
const lazyResults$ = from(urls).pipe(
  mergeMap((url) => defer(() => fetch(url)), 2),
);

Hai pipeline trên chỉ để so sánh cách khai báo, không nên chạy cả hai trong ứng dụng. lazyResults$ phải được subscribe mới bắt đầu request. Với eagerResults$, giới hạn chỉ quản lý việc subscribe Promise đã có, không thể quay ngược thời gian để trì hoãn fetch.

defer giúp lazy creation và chuyển exception từ factory vào error channel, không tự thêm cancellation cho Promise. Nếu cần hủy fetch, dùng wrapper có AbortController và teardown hoặc API Observable phù hợp như fromFetch với selector bao gồm phần đọc body. Đồng thời kiểm tra response.ok, vì fetch không reject chỉ vì HTTP trả 404 hay 500.

Một callback async bên trong mergeMap cũng trả Promise và được flatten. Nhưng đừng vì cú pháp tiện mà bỏ qua chính sách lỗi, timeout và teardown: lifecycle của công việc quan trọng hơn việc callback có từ khóa async hay không.

Chọn giữa các flattening operator

OperatorKhi outer có value mới và inner cũ chưa xongKhi nên chọn
mergeMapChạy thêm inner nếu có slot, nếu không thì xếp hàngBatch task độc lập, không cần thứ tự output
concatMapXếp hàng, đợi inner trước completeLệnh ghi phụ thuộc thứ tự, các bước cần chạy tuần tự
switchMapUnsubscribe inner cũ rồi chuyển sang inner mớiSearch hoặc dữ liệu UI chỉ cần query mới nhất
exhaustMapBỏ value mới trong lúc inner đang chạyChặn submit lặp trong cùng một lượt xử lý

Không chọn mergeMap cho search chỉ vì muốn response nhanh: request cũ có thể về sau request mới và ghi đè UI bằng dữ liệu cũ. Không chọn nó cho hai lệnh sửa cùng một resource nếu thứ tự tác động có ý nghĩa. Và nếu không muốn nhận thêm việc khi đang bận, buffer của mergeMap không giống chính sách bỏ event của exhaustMap.

Ngược lại, đừng chọn switchMap cho batch bắt buộc lấy kết quả của mọi ID. mergeMap giữ inner cũ khi có input mới, nhưng vẫn phải có owner quản lý cancellation và chính sách lỗi; không có operator nào tự bảo đảm công việc hoàn tất trong mọi tình huống.

Kiểm thử concurrency bằng marble test

Timer thực thích hợp để quan sát, nhưng unit test nên dùng virtual time để không phụ thuộc event loop. Test dưới kiểm tra cả output lẫn thời điểm subscribe inner C khi giới hạn là hai:

import { mergeMap } from 'rxjs';
import { TestScheduler } from 'rxjs/testing';

const scheduler = new TestScheduler((actual, expected) => {
  if (JSON.stringify(actual) !== JSON.stringify(expected)) {
    throw new Error('Marble assertion failed');
  }
});

scheduler.run(({ cold, expectObservable, expectSubscriptions }) => {
  const source$ = cold('(abc|)');
  const a$ = cold('------(x|)', { x: 'A' });
  const b$ = cold('--(y|)', { y: 'B' });
  const c$ = cold('--(z|)', { z: 'C' });
  const inners = { a: a$, b: b$, c: c$ };

  const result$ = source$.pipe(
    mergeMap((key) => inners[key as keyof typeof inners], 2),
  );

  expectObservable(result$).toBe('--y-z-(x|)', {
    x: 'A',
    y: 'B',
    z: 'C',
  });
  expectSubscriptions(a$.subscriptions).toBe('^-----!');
  expectSubscriptions(b$.subscriptions).toBe('^-!');
  expectSubscriptions(c$.subscriptions).toBe('--^-!');
});

Ở frame 0, outer phát A, B, C rồi complete. A và B được subscribe ngay; B complete ở frame 2, C mới được subscribe tại đó và phát ở frame 4. A phát và complete ở frame 6, làm output complete. expectSubscriptions chứng minh C thật sự chờ slot, thay vì chỉ có output tình cờ đúng thứ tự.

Khi áp dụng vào production code, thêm test cho inner lỗi, outer lỗi, teardown khi hủy và inner không complete. TestScheduler kiểm soát Observable dựa trên scheduler, không tự virtualize native Promise; dùng cold Observable làm test double thay vì gọi fetch thật trong marble test.

Checklist trước khi dùng

  • Các task có độc lập không, hay cùng sửa một resource theo thứ tự?
  • Mọi kết quả đều cần được giữ, hay chỉ kết quả mới nhất có giá trị?
  • Giới hạn concurrency là bao nhiêu, và nó áp dụng cho một subscription hay toàn hệ thống?
  • Công việc có được tạo lazy trong projection, hay Promise đã chạy từ trước?
  • Source hữu hạn hay có thể phát nhanh hơn consumer mãi mãi?
  • Mỗi inner có đường complete hoặc timeout rõ ràng không?
  • Một item lỗi phải dừng batch hay trở thành error result riêng?
  • Ai hủy pipeline, và teardown có dừng được tài nguyên bên dưới không?
  • Downstream có cần id hoặc inputIndex để tránh nhầm thứ tự không?
  • Nếu retry lệnh ghi, backend có bảo đảm idempotency không?

Nếu chưa trả lời được câu hỏi về buffer hoặc completion, hãy vẽ timeline trước khi tăng concurrency. Tăng số slot có thể che việc chờ trong một batch nhỏ, nhưng không sửa được source quá nhanh hay inner sống vô hạn.

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

  1. Đổi concurrency của ví dụ A/B/C từ mặc định sang 1, 2, 3. Dự đoán thời điểm C bắt đầu và thời điểm cả stream complete trước khi chạy.
  2. Trong ví dụ batch, chuyển catchError ra ngoài mergeMap. Cho một request lỗi và quan sát request khác cùng các ID đang chờ.
  3. Thay một inner bằng NEVER, đặt concurrency là 1. Giải thích vì sao task sau không bắt đầu dù outer đã complete.
  4. Di chuyển takeUntil lên trước mergeMap trong ví dụ cancellation. Kiểm tra A/B có phát không và C có được bắt đầu không.
  5. Sửa marble test thành concurrent = 1, rồi viết lại expected output và subscription của B/C.

Sau đó, lấy một batch request trong ứng dụng của bạn, gắn metadata vào từng result và thêm test subscription cho trường hợp hết slot. Đây là cách kiểm tra trực tiếp rằng concurrency giới hạn công việc bắt đầu, không chỉ làm output trông có vẻ tuần tự.

Học tiếp

Nguồn tham khảo

On this page