Học RxJS
Operators

Transformation operators

Biến đổi cấu trúc và giá trị của emission.

API thường không trả đúng hình dạng mà UI cần. Một stream có thể phát từng OrderEvent, trong khi màn hình lại cần tổng tiền, trạng thái hiện tại và nhãn để hiển thị. Transformation operator giúp bạn đổi từng emission, tích lũy state hoặc gom nhiều emission thành một cấu trúc mới mà vẫn giữ mọi thứ trong pipeline.

Phạm vi phiên bản

Bài này dùng public API của RxJS 7.8.2 và import trực tiếp từ rxjs. Cách phân nhóm bên dưới dựa trên hành vi để bạn dễ chọn operator; đây không phải danh sách đầy đủ mọi API được RxJS xếp vào nhóm transformation.

Mục lục

Mental model của transformation operator

Một transformation operator nhận Observable<T> và trả về một Observable mới. Source không bị sửa; operator chỉ mô tả cách subscriber của output sẽ nhận notification từ source.

Observable<T> ──► transformation operator ──► Observable<R>

next(T)       ──► đổi / tích lũy / gom       ──► next(R)
error         ────────────────────────────────► error
complete      ──► có thể kích hoạt output cuối ─► complete

Dòng error và complete rất đáng để để ý. map thường chuyển tiếp complete ngay, nhưng reduce và toArray phải chờ complete mới có kết quả. Vì vậy, hai pipeline cùng đổi T thành R vẫn có lifecycle hoàn toàn khác nhau.

Ba chiều có thể thay đổi

Khi nói “biến đổi”, hãy tách ba câu hỏi:

  1. Giá trị có đổi không? map đổi từng T thành R; pairwise đổi một value thành cặp [previous, current].
  2. Số emission có đổi không? scan thường phát một state cho mỗi input; reduce nén cả source thành tối đa một output; bufferCount nén nhiều input thành từng array.
  3. Cấu trúc thời gian có đổi không? bufferTime tạo ranh giới theo clock; higher-order mapping quản lý nhiều inner Observable và cả cancellation.

Transformation không đồng nghĩa với “chỉ đổi object”. Có operator giữ nguyên số emission, có operator trì hoãn output, và có operator tạo Observable lồng. Đọc operator theo cả data shape lẫn lifecycle sẽ ít bất ngờ hơn.

Bản đồ chọn operator

Bạn cần gì?Chọn mặc địnhOutput xuất hiện khi nào?
Đổi từng value hoặc chọn vài fieldmapSau mỗi source emission
Giữ state tích lũy và phát state mới liên tụcscanSau mỗi source emission
Tính một kết quả cuốireduceKhi source complete
Thu toàn bộ value thành arraytoArrayKhi source complete
So sánh value hiện tại với value ngay trướcpairwiseTừ source emission thứ hai
Gom theo số lượngbufferCountKhi batch đủ; batch dở dang còn lại phát lúc complete
Gom theo khoảng thời gianbufferTimeKhi cửa sổ thời gian đóng
Chia value theo keygroupByMỗi key tạo một GroupedObservable
Chạy công việc trả về ObservableconcatMap, mergeMap, switchMap, exhaustMapTùy chiến lược flatten
Chia stream thành các Observable conwindow, windowCount, windowTimeKhi một window được mở

Mặc định, mình bắt đầu bằng map nếu mỗi input độc lập. Chỉ chuyển sang scan, buffer hoặc flattening operator khi yêu cầu thật sự có state, batching hoặc công việc bất đồng bộ. Chọn operator phức tạp hơn quá sớm thường làm lifecycle khó đọc mà không đem lại lợi ích.

map biến đổi từng emission

Một input tạo một output

map(project) gọi projection cho từng source emission rồi phát giá trị được trả về. Thứ tự được giữ nguyên, và map không tự thêm concurrency hay delay.

import { from, map } from 'rxjs';

type ApiUser = {
  id: number;
  first_name: string;
  last_name: string;
};

type UserOption = {
  value: number;
  label: string;
};

const users: ApiUser[] = [
  { id: 1, first_name: 'An', last_name: 'Nguyễn' },
  { id: 2, first_name: 'Bình', last_name: 'Trần' },
];

from(users)
  .pipe(
    map(
      (user): UserOption => ({
        value: user.id,
        label: `${user.last_name} ${user.first_name}`,
      }),
    ),
  )
  .subscribe((option) => console.log(option));

// Output:
// { value: 1, label: 'Nguyễn An' }
// { value: 2, label: 'Trần Bình' }

Hai input tạo đúng hai output. Nếu mục tiêu là bỏ một số user, hãy thêm filter; trả undefined trong map không loại emission mà chỉ biến nó thành một emission có giá trị undefined.

Projection còn nhận index bắt đầu từ 0 cho từng subscription:

import { map, of } from 'rxjs';

of('alpha', 'beta').pipe(
  map((value, index) => `${index + 1}. ${value}`),
).subscribe(console.log);

// 1. alpha
// 2. beta

index thuộc execution của operator, nên một subscription mới sẽ đếm lại từ 0.

Giữ phép biến đổi thuần và rõ kiểu

Một projection dễ đọc thường chỉ tính output từ input, không sửa input và không làm side effect. Với object, ưu tiên tạo object mới:

import { map, of } from 'rxjs';

type Product = {
  id: string;
  price: number;
};

type PricedProduct = Product & {
  priceWithTax: number;
};

of<Product>({ id: 'book-01', price: 100_000 }).pipe(
  map(
    (product): PricedProduct => ({
      ...product,
      priceWithTax: product.price * 1.08,
    }),
  ),
).subscribe(console.log);

Nếu callback cần log, metric hoặc gọi một API imperative, đó là side effect. Dùng tap cho việc quan sát, hoặc higher-order mapping cho công việc trả về Observable; đừng giấu chúng trong map vì người đọc sẽ kỳ vọng map chỉ quyết định output.

map không sửa source

map tạo một Observable mới. Nó có thể trả lại cùng object reference nếu callback của bạn làm vậy, nhưng operator không tự mutate source value. Việc mutate hay tạo object mới nằm trong projection do bạn viết.

Lỗi trong projection đi vào error channel

Nếu projection ném exception, RxJS chuyển exception đó thành error và đóng subscription hiện tại:

import { map, of } from 'rxjs';

of('{"id":1}', '{broken}', '{"id":3}').pipe(
  map((text) => JSON.parse(text) as { id: number }),
).subscribe({
  next: (value) => console.log('next:', value.id),
  error: (error: unknown) => {
    const message = error instanceof Error ? error.message : String(error);
    console.log('error:', message);
  },
  complete: () => console.log('complete'),
});

// next: 1
// error: ...

Value thứ ba không được xử lý và callback complete không chạy. Nếu dữ liệu lỗi là trường hợp nghiệp vụ có thể bỏ qua, hãy parse an toàn rồi filter theo kết quả. Nếu đó là lỗi thật, giữ nó trong error channel và xử lý bằng catchError ở đúng cấp.

scan phát state sau mỗi emission

scan(accumulator, seed) giống một reducer chạy theo thời gian: mỗi input cập nhật accumulator, và state mới được phát ngay xuống downstream. Đây là lựa chọn tự nhiên cho counter, state machine nhỏ, giỏ hàng hoặc dữ liệu được xây dần từ event.

import { from, scan } from 'rxjs';

type CartEvent =
  | { type: 'add'; quantity: number }
  | { type: 'remove'; quantity: number };

type CartState = {
  itemCount: number;
};

const events: CartEvent[] = [
  { type: 'add', quantity: 2 },
  { type: 'add', quantity: 1 },
  { type: 'remove', quantity: 1 },
];

from(events).pipe(
  scan(
    (state: CartState, event): CartState => ({
      itemCount:
        event.type === 'add'
          ? state.itemCount + event.quantity
          : Math.max(0, state.itemCount - event.quantity),
    }),
    { itemCount: 0 },
  ),
).subscribe((state) => console.log(state));

// { itemCount: 2 }
// { itemCount: 3 }
// { itemCount: 2 }

Khác map, output hiện tại của scan phụ thuộc cả state trước đó lẫn input mới. Điều này rất hữu ích, nhưng cũng có nghĩa mỗi subscription có accumulator riêng. scan tự nó không tạo shared state cho toàn ứng dụng.

Seed không tự được phát

Trong ví dụ trên, { itemCount: 0 } là state ban đầu để xử lý event đầu tiên. scan không phát seed ngay lúc subscribe; nó chỉ phát sau khi source có value.

source$:  ─────(add 2)────(add 1)────(remove 1)────│
scan:           {2}─────────{3}──────────{2}──────│
seed {0}:  dùng để tính, không tự xuất hiện ở output

Nếu consumer phải nhận state ban đầu trước event đầu tiên, thêm startWith(initialState) sau scan, hoặc thiết kế source phát một event khởi tạo. Mặc định mình dùng startWith khi state ban đầu thật sự là một UI state cần render, không chỉ là chi tiết nội bộ của phép tính.

Không truyền seed cũng hợp lệ: value đầu tiên trở thành accumulator ban đầu. Tuy nhiên, type thường khó đọc hơn và source rỗng không có output; với state nghiệp vụ, seed rõ ràng thường an toàn hơn.

Tránh mutate accumulator

Đoạn dưới chạy được nhưng dễ tạo lỗi:

scan((state, event) => {
  state.itemCount += event.quantity;
  return state;
}, initialState);

Mọi emission trả cùng một object reference. distinctUntilChanged() mặc định sẽ coi chúng là cùng một value; UI dựa trên reference equality cũng có thể không render lại. Ngoài ra, initialState bên ngoài bị sửa, khiến subscription sau bắt đầu từ dữ liệu bẩn.

Hãy trả object mới như ví dụ trước. Khi state có nhiều nhánh, tách reducer thành một function thuần rồi unit test function đó trước khi đặt vào scan.

reduce và toArray chờ source complete

reduce cũng tích lũy state, nhưng không phát kết quả trung gian. Nó chỉ phát accumulator cuối cùng khi source complete:

import { of, reduce } from 'rxjs';

of(120, 80, 50).pipe(
  reduce((total, price) => total + price, 0),
).subscribe({
  next: (total) => console.log('tổng:', total),
  complete: () => console.log('complete'),
});

// tổng: 250
// complete

toArray() là trường hợp chuyên biệt: giữ mọi source emission trong một array rồi phát array đó khi complete.

import { interval, take, toArray } from 'rxjs';

interval(100).pipe(
  take(3),
  toArray(),
).subscribe(console.log);

// Sau khoảng 300 ms: [0, 1, 2]

Điểm quyết định là source có complete hay không:

OperatorPhát trung gian?Cần complete để có output?Source rỗng với seed 0 hoặc array
scan(..., 0)CóKhôngKhông phát gì
reduce(..., 0)KhôngCóPhát 0 khi complete
toArray()KhôngCóPhát [] khi complete

Đừng đặt toArray() trực tiếp sau fromEvent() hoặc một stream sống vô hạn rồi chờ kết quả: array sẽ tiếp tục giữ value và không được phát. Hãy tạo ranh giới bằng take, takeUntil, buffer..., hoặc chọn state liên tục bằng scan.

Completion là một phần của yêu cầu

reduce và toArray không phù hợp chỉ vì output type trông đúng. Bạn phải biết ai làm source complete. Nếu không có câu trả lời rõ ràng, pipeline có thể không phát gì và còn giữ dữ liệu trong bộ nhớ.

pairwise đặt giá trị hiện tại cạnh giá trị trước

pairwise() phát tuple [previous, current] từ source emission thứ hai. Nó hợp với việc tính delta, phát hiện chuyển trạng thái hoặc so sánh tọa độ liên tiếp.

import { from, map, pairwise } from 'rxjs';

from([100, 108, 103, 120]).pipe(
  pairwise(),
  map(([previous, current]) => ({
    current,
    delta: current - previous,
  })),
).subscribe(console.log);

// { current: 108, delta: 8 }
// { current: 103, delta: -5 }
// { current: 120, delta: 17 }

Value 100 chưa có value trước để ghép nên chưa tạo output. Nếu cần so sánh emission đầu với một mốc ban đầu, thêm startWith(initialValue) trước pairwise().

pairwise chỉ nhớ đúng một value trước đó. Nếu bạn cần cửa sổ trượt ba hoặc nhiều value, dùng bufferCount(size, 1); tham số thứ hai 1 mở một buffer mới ở mỗi emission.

buffer gom emission thành batch

Buffer operator đổi nhiều next(T) thành từng next(T[]). Đây là cách tạo batch trước khi ghi log, gửi telemetry hoặc xử lý dữ liệu theo khối. Câu hỏi quan trọng không phải “có cần array không?” mà là điều gì đóng batch: đủ số lượng, hết thời gian, hay notifier phát tín hiệu.

Gom theo số lượng với bufferCount

bufferCount(bufferSize) đóng batch khi đủ số phần tử. Khi source complete, buffer còn dở nhưng không rỗng vẫn được phát:

import { bufferCount, of } from 'rxjs';

of(1, 2, 3, 4, 5, 6, 7).pipe(
  bufferCount(3),
).subscribe(console.log);

// [1, 2, 3]
// [4, 5, 6]
// [7]

Nếu downstream chỉ chấp nhận batch đủ ba phần tử, thêm filter((batch) => batch.length === 3). Đừng mặc định loại batch cuối: với log hoặc đơn hàng, đó có thể là dữ liệu thật cần được xử lý.

Tham số thứ hai điều khiển tần suất mở buffer:

import { bufferCount, from } from 'rxjs';

from([1, 2, 3, 4, 5]).pipe(
  bufferCount(3, 1),
).subscribe(console.log);

// [1, 2, 3]
// [2, 3, 4]
// [3, 4, 5]
// [4, 5]
// [5]

Đây là sliding window theo count. Vì nhiều buffer cùng mở, mỗi value có thể xuất hiện trong nhiều output.

Gom theo thời gian với bufferTime

bufferTime(1_000) đóng một batch xấp xỉ mỗi giây. Ví dụ dưới đếm click theo từng khoảng một giây và bỏ batch rỗng:

import { bufferTime, filter, fromEvent, map } from 'rxjs';

const button = document.querySelector<HTMLButtonElement>('#track-click');

if (!button) {
  throw new Error('Không tìm thấy #track-click');
}

const clickBatches$ = fromEvent<MouseEvent>(button, 'click').pipe(
  bufferTime(1_000),
  filter((batch) => batch.length > 0),
  map((batch) => ({
    count: batch.length,
    lastClickAt: batch.at(-1)?.timeStamp ?? 0,
  })),
);

const subscription = clickBatches$.subscribe(console.log);

// Khi view bị hủy:
subscription.unsubscribe();

Khoảng 1_000 ms là lịch theo scheduler, không phải deadline thời gian thực tuyệt đối; event loop có thể làm callback chạy muộn. Nếu hệ thống cần “đủ 100 item hoặc hết 1 giây thì gửi”, dùng overload có giới hạn kích thước của bufferTime và kiểm tra đúng signature của phiên bản đang chạy.

Với source nhanh, mọi buffer đều giữ value trong bộ nhớ cho tới khi đóng. Hãy chọn time span và giới hạn batch theo tải thực tế, thay vì đặt khoảng dài chỉ để giảm số lần gọi downstream.

Buffer và window khác nhau ở output

Hai nhóm cùng chia source thành đoạn, nhưng output khác nhau:

source$:     ──1──2──3──4──5──6──│

bufferCount(3)
output$:     ─────[1,2,3]────[4,5,6]│

windowCount(3)
output$:     ─────Observable<1,2,3>──Observable<4,5,6>│
  • buffer... materialize mỗi đoạn thành array. Dễ dùng, nhưng phải giữ toàn bộ batch trong bộ nhớ.
  • window... phát các Observable con. Downstream có thể xử lý từng value khi nó đến, nhưng bạn phải flatten và quản lý lifecycle của inner stream.

Mặc định, dùng buffer khi batch nhỏ và cần gửi cả mảng. Dùng window khi mỗi đoạn vẫn cần pipeline riêng hoặc khi không muốn chờ materialize toàn bộ đoạn trước khi bắt đầu xử lý.

groupBy chia stream theo key

groupBy(keySelector) biến một stream thành stream của các GroupedObservable. Mỗi group có thuộc tính key và phát các value thuộc key đó. Vì output bị lồng, bạn thường cần mergeMap, concatMap hoặc một chiến lược flatten khác để tiêu thụ từng group.

import { from, groupBy, map, mergeMap, reduce, toArray } from 'rxjs';

type Payment = {
  userId: string;
  amount: number;
};

const payments: Payment[] = [
  { userId: 'u1', amount: 50 },
  { userId: 'u2', amount: 40 },
  { userId: 'u1', amount: 70 },
];

from(payments).pipe(
  groupBy((payment) => payment.userId),
  mergeMap((group$) =>
    group$.pipe(
      reduce((total, payment) => total + payment.amount, 0),
      map((total) => ({ userId: group$.key, total })),
    ),
  ),
  toArray(),
).subscribe(console.log);

// [
//   { userId: 'u1', total: 120 },
//   { userId: 'u2', total: 40 }
// ]

reduce bên trong mỗi group chỉ phát khi group complete; ở đây source hữu hạn complete sau ba payment nên mọi group cũng complete. Với event stream sống lâu, đoạn code tương tự có thể không bao giờ phát tổng.

groupBy cần chiến lược đóng group

Với source dài hạn và key không giới hạn như userId, sessionId hoặc URL, số group có thể tăng mãi. groupBy có tùy chọn duration để đóng group và connector để kiểm soát subject đứng sau group. Nếu bạn không định nghĩa được vòng đời của từng key, đừng dùng groupBy như một Map cache vô hạn.

Nếu chỉ cần cộng tổng cho một collection hữu hạn đã có sẵn, một reduce trên array có thể đơn giản hơn. groupBy đáng dùng khi từng group cần tiếp tục là một stream với operator và lifecycle riêng.

map không flatten Observable lồng

Giả sử mỗi search query tạo một Observable request. Nếu dùng map, output sẽ là Observable<Observable<Result>>:

const nested$ = query$.pipe(
  map((query) => search$(query)),
);

map đã làm đúng nhiệm vụ: mỗi string được đổi thành một Observable. Nó không tự quyết định subscribe inner nào, giữ bao nhiêu inner cùng lúc, hoặc hủy inner cũ hay không. Những quyết định đó thuộc về flattening operator.

Chọn chiến lược flatten theo lifecycle

OperatorKhi inner đang chạy mà outer phát value mớiHợp với
concatMapXếp value mới vào hàng đợiGhi tuần tự, cần giữ thứ tự
mergeMapSubscribe thêm inner, có thể giới hạn concurrencyCông việc độc lập chạy song song
switchMapUnsubscribe inner cũ, chuyển sang inner mớiSearch, route param, dữ liệu mới nhất thắng
exhaustMapBỏ qua value mới cho tới khi inner hiện tại completeChặn double-submit, login đang xử lý

Không có operator “nhanh nhất” cho mọi bài toán. Mặc định của mình là nói thành câu nghiệp vụ trước: xếp hàng, chạy song song, hủy cũ, hay bỏ mới. Câu trả lời chọn operator, không phải thói quen.

outer value mới đến khi inner cũ chưa xong

concatMap   ─► đợi
mergeMap    ─► chạy thêm
switchMap   ─► hủy inner cũ rồi chuyển
exhaustMap  ─► bỏ outer value mới

Các operator này vừa project value thành inner Observable, vừa flatten output. Chúng được giải thích sâu trong nhóm Higher-order Observables.

Cẩn thận với async callback

Callback async luôn trả Promise. Vì vậy, đoạn sau có type Observable<Promise<User>>, không phải Observable<User>:

const userPromise$ = userId$.pipe(
  map(async (id) => fetchUser(id)),
);

Nếu fetchUser đã trả Promise, async còn có thể che khuất hình dạng thật của pipeline. Dùng flattening operator để nhận kết quả:

import { from, switchMap } from 'rxjs';

const user$ = userId$.pipe(
  switchMap((id) => from(fetchUser(id))),
);

Tuy nhiên, unsubscribe khỏi from(promise) chỉ ngăn kết quả đi xuống subscriber; nó không tự hủy Promise hoặc request nền. Muốn cancellation thật, source phải tích hợp AbortController, fromFetch hoặc teardown tương ứng. switchMap quản lý subscription, không thể tự phát minh cơ chế hủy mà producer không cung cấp.

Ví dụ thực tế xây state đơn hàng

Giả sử backend hoặc WebSocket phát event của một đơn hàng. UI không muốn tự nối từng event; nó cần một OrderViewModel hoàn chỉnh sau mỗi thay đổi.

import { from, map, scan } from 'rxjs';

type OrderEvent =
  | { type: 'item-added'; price: number }
  | { type: 'item-removed'; price: number }
  | { type: 'submitted' };

type OrderState = {
  itemCount: number;
  total: number;
  status: 'draft' | 'submitted';
};

type OrderViewModel = OrderState & {
  canSubmit: boolean;
  totalLabel: string;
};

const initialState: OrderState = {
  itemCount: 0,
  total: 0,
  status: 'draft',
};

const events: OrderEvent[] = [
  { type: 'item-added', price: 120_000 },
  { type: 'item-added', price: 80_000 },
  { type: 'item-removed', price: 120_000 },
  { type: 'submitted' },
];

const orderViewModel$ = from(events).pipe(
  scan((state: OrderState, event): OrderState => {
    switch (event.type) {
      case 'item-added':
        return {
          ...state,
          itemCount: state.itemCount + 1,
          total: state.total + event.price,
        };
      case 'item-removed':
        return {
          ...state,
          itemCount: Math.max(0, state.itemCount - 1),
          total: Math.max(0, state.total - event.price),
        };
      case 'submitted':
        return { ...state, status: 'submitted' };
    }
  }, initialState),
  map(
    (state): OrderViewModel => ({
      ...state,
      canSubmit: state.status === 'draft' && state.itemCount > 0,
      totalLabel: `${state.total.toLocaleString('vi-VN')} ₫`,
    }),
  ),
);

orderViewModel$.subscribe(console.log);

Pipeline chia trách nhiệm thành hai lớp:

  1. scan diễn giải event và tạo domain state mới. Đây là nơi chứa quy tắc cộng, trừ và chuyển trạng thái.
  2. map tạo view model từ state, gồm canSubmit và chuỗi tiền tệ dành cho UI.

Cách tách này giúp reducer không phụ thuộc format hiển thị. Nếu sau này có màn hình khác cần cùng OrderState nhưng nhãn khác, bạn tái sử dụng state stream rồi đặt một map riêng ở ranh giới UI.

Có hai chi tiết cần quyết định trong ứng dụng thật. Thứ nhất, event item-removed nên mang itemId để reducer kiểm tra đúng item thay vì chỉ trừ một con số. Thứ hai, submitted có được phép tới khi giỏ rỗng hay không là validation nghiệp vụ; đừng để map hiển thị che đi một state không hợp lệ.

Các shortcut deprecated nên tránh

RxJS 7 vẫn có một số shortcut cũ, nhưng code mới nên dùng API tổng quát để dễ nâng cấp:

Tránh trong code mớiThay bằng
mapTo(value)map(() => value)
pluck('user', 'name')map((value) => value.user?.name)
concatMapTo(inner$)concatMap(() => inner$)
mergeMapTo(inner$)mergeMap(() => inner$)
switchMapTo(inner$)switchMap(() => inner$)

Dạng thay thế nói rõ projection và cho TypeScript suy luận tốt hơn. Với property lồng, optional chaining trong map cũng dễ nhìn thấy kiểu undefined hơn pluck.

Đừng nhầm với static creation function

mapTo bị deprecated không liên quan tới Map của JavaScript. Tương tự, các operator có hậu tố To ở bảng trên là convenience API cũ; operator gốc như map, concatMap, mergeMap và switchMap vẫn là API chính.

Những bẫy thường gặp

  1. Dùng map để lọc. Trả null hoặc undefined vẫn tạo emission. Dùng filter khi muốn loại value.
  2. Dùng map cho side effect. Logging thuộc về tap; công việc async thuộc về flattening operator phù hợp.
  3. Mutate object trong map hoặc accumulator trong scan. Reference không đổi làm equality check sai và state có thể rò giữa các subscription.
  4. Chờ output từ reduce hoặc toArray trên source không complete. DOM event, WebSocket và interval thường cần ranh giới kết thúc rõ ràng.
  5. Quên rằng pairwise không phát ở value đầu. Thêm startWith nếu cần một mốc so sánh ban đầu.
  6. Bỏ batch cuối của bufferCount. Khi source complete, batch dở dang vẫn được phát; hãy xử lý nó có chủ đích.
  7. Tạo buffer quá lớn hoặc quá lâu. Value được giữ trong bộ nhớ cho tới khi buffer đóng.
  8. Dùng groupBy trên key không giới hạn mà không đóng group. Long-lived source có thể giữ ngày càng nhiều group.
  9. Dùng map(async ...) rồi tưởng downstream nhận value đã resolve. Downstream thật ra nhận Promise.
  10. Chọn switchMap chỉ vì phổ biến. Nó hủy inner cũ; với thao tác ghi bắt buộc hoàn tất, hành vi đó có thể làm mất công việc.
  11. Cho rằng transformation operator sửa source. Mỗi lời gọi pipe tạo Observable mới; source vẫn có thể được dùng trong pipeline khác.
  12. Nhầm data transformation với scheduling. map và scan không tự chuyển việc sang background thread hay làm source thành async.

Checklist chọn operator

Trước khi thêm operator, trả lời lần lượt:

  • Mỗi input tạo đúng một output, hay cần giữ/bỏ/gom nhiều input?
  • Output có cần state từ emission trước không?
  • Consumer cần kết quả liên tục hay chỉ kết quả cuối?
  • Source có chắc chắn complete không, và ai làm nó complete?
  • Batch đóng theo count, time hay một notifier khác?
  • Output có phải Observable hoặc Promise lồng không?
  • Nếu có công việc lồng, value mới phải xếp hàng, chạy song song, hủy cũ hay bị bỏ?
  • Accumulator hoặc buffer có thể tăng vô hạn không?
  • Projection có thuần không, hay đang giấu side effect?
  • Khi unsubscribe, resource bên ngoài có thật sự dừng không?

Nếu chưa rõ lifecycle, hãy phác một marble diagram trước. Vài dòng timeline thường làm lộ ngay việc source không complete, batch không đóng hoặc inner Observable bị hủy sai lúc.

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

  1. Dùng map chuyển from([{ price: 10, quantity: 2 }, { price: 5, quantity: 3 }]) thành từng subtotal. Sau đó dùng reduce để tính grand total và dự đoán thời điểm output xuất hiện.
  2. Viết counter bằng scan cho ba event increment, increment, decrement. Thêm startWith(0) và so sánh số emission trước và sau.
  3. Chạy bufferCount(3, 1) với source [1, 2, 3, 4]. Vẽ tất cả batch, kể cả batch dở dang lúc complete.
  4. Dùng pairwise để tính chênh lệch nhiệt độ từ [28, 30, 29, 32]. Giải thích vì sao chỉ có ba output.
  5. Thay map(async (id) => load(id)) bằng từng concatMap, mergeMap, switchMap và exhaustMap. Với hai ID đến gần nhau, viết trước expected lifecycle của request thứ nhất rồi mới chạy code.
  6. Tạo source không complete như interval(100) và đặt toArray() phía sau. Sau đó thêm take(3) để thấy chính xác điều gì làm array được phát.

Sau các bài này, lấy một pipeline thật trong dự án và ghi chú type sau từng operator. Nếu một bước biến Observable<T> thành Observable<Observable<R>>, Observable<Promise<R>> hoặc giữ collection không giới hạn, đó là chỗ cần xem lại trước tiên.

Học tiếp

Nguồn tham khảo

On this page