Học RxJS
Bắt đầu

Pipe và operator

Ghép các operator thành pipeline có thể đọc và tái sử dụng.

Bạn đã có một Observable, nhưng dữ liệu thô hiếm khi đúng ngay với thứ UI hoặc service cần. Một ô tìm kiếm phát ra mọi lần gõ phím; bạn lại chỉ muốn chuỗi đã bỏ khoảng trắng, đủ dài và không trùng với lần trước. pipe là nơi ghép các bước xử lý đó thành một pipeline có thứ tự rõ ràng.

Phiên bản sử dụng

Các API và import trong bài được kiểm chứng với RxJS 7.8.x, thuộc nhánh RxJS 7.x ổn định. Ví dụ dùng public exports trực tiếp từ rxjs của phiên bản này; bài không suy đoán API của phiên bản tương lai.

Mục lục

Bài toán pipe giải quyết

Không có pipe, code xử lý event rất dễ biến thành một callback dài: đọc giá trị, kiểm tra điều kiện, chuyển kiểu, ghi log rồi cập nhật UI. Các bước dính vào nhau, khó thử riêng và khó nhìn ra bước nào làm thay đổi dữ liệu.

Với RxJS, mỗi bước có thể được biểu diễn bằng một pipeable operator. Bạn truyền các operator vào source$.pipe(...); output của operator trước trở thành input của operator sau. Kết quả là một Observable mới để tiếp tục ghép hoặc subscribe.

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

const source$ = of(1, 2, 3, 4);

const result$ = source$.pipe(
  filter((value) => value % 2 === 0),
  map((value) => value * 10),
);

result$.subscribe((value) => console.log(value));

// Output:
// 20
// 40

source$ vẫn mô tả dãy 1, 2, 3, 4; pipe không sửa nó. Pipeline mới chỉ giữ công thức “lọc số chẵn rồi nhân 10”. Gọi pipe cũng chưa tự subscribe và chưa làm nguồn chạy.

Mental model của pipeline

Ba vai trò khác nhau

Ba khái niệm dưới đây thường nằm trên cùng vài dòng code nên rất dễ bị gọi lẫn tên:

Khái niệmHình dạng rút gọnVai trò
Observable<T>nguồn phát các notification kiểu TMô tả nguồn dữ liệu hoặc event theo thời gian
OperatorFunction<T, R>(source$: Observable<T>) => Observable<R>Nhận một Observable và trả về Observable khác
subscribe(...)Observable<T> -> SubscriptionGắn consumer vào pipeline và trả về handle để teardown

Ví dụ, map là operator factory. Lệnh map((n) => n * 2) chưa xử lý số nào; nó tạo một OperatorFunction<number, number>. source$.pipe(...) ghép các function này. Chỉ khi có subscriber, chuỗi subscription mới được thiết lập về phía nguồn.

Bạn có thể hình dung pipeline như dây chuyền kiểm tra hành lý: mỗi trạm chỉ biết input của mình và chuyển kết quả sang trạm kế tiếp. Phép so sánh dừng ở đây vì Observable còn có ba loại notification (next, error, complete) và có teardown, không chỉ có các “món đồ” đi qua.

Creation function không phải pipeable operator

of, from, interval và fromEvent tạo Observable nên thường đứng trước .pipe(...). Bạn không truyền of(...) vào giữa danh sách operator. Ngược lại, map, filter, take và tap nhận Observable nguồn một cách gián tiếp qua pipe.

Dữ liệu và teardown đi theo hướng nào

Giả sử code là source$.pipe(filter(...), map(...), tap(...), take(2)). Có ba chuyển động cần tách biệt:

Thiết lập subscription:
Observer ─► take ─► tap ─► map ─► filter ─► source$

Notification từ nguồn:
source$ ─► filter ─► map ─► tap ─► take ─► Observer
          next / error / complete đi từ trái sang phải trong pipe

Teardown khi unsubscribe:
Observer ─► take ─► tap ─► map ─► filter ─► teardown của source$

Vì vậy, “đọc pipeline từ trái sang phải” là cách đúng để theo dõi giá trị phát ra. Ở chiều ngược lại, subscriber cuối cùng kéo cả chuỗi subscription về nguồn; khi hủy, teardown truyền ngược lên để nguồn giải phóng timer, event listener hoặc tài nguyên khác.

Một error hoặc complete là tín hiệu kết thúc và không đi tiếp dưới dạng giá trị next. Sau khi subscription đóng, pipeline đó không phát thêm giá trị. Muốn tìm hiểu kỹ ba channel này, xem Subscribe và Observer.

Pipeline đầu tiên

Đọc code từ trái sang phải

Ví dụ sau cố ý dùng nguồn đồng bộ để bạn thấy chính xác thời điểm từng bước chạy:

import { filter, finalize, map, of, take, tap } from 'rxjs';

const result$ = of(1, 2, 3, 4, 5).pipe(
  filter((value) => value % 2 === 1),
  map((value) => value * 10),
  tap((value) => console.log('trong pipeline:', value)),
  take(2),
  finalize(() => console.log('đã teardown')),
);

console.log('trước subscribe');

result$.subscribe({
  next: (value) => console.log('next:', value),
  error: (error: unknown) => console.error('error:', error),
  complete: () => console.log('complete'),
});

console.log('sau subscribe');

Luồng xử lý là:

  1. filter chỉ cho 1, 3, 5 đi tiếp.
  2. map đổi chúng thành 10, 30, 50.
  3. tap quan sát giá trị mà không chủ ý biến đổi nó.
  4. take(2) chỉ nhận hai giá trị đầu rồi chủ động hoàn tất.
  5. finalize đăng ký callback cleanup cho lúc complete, error hoặc unsubscribe.

Đọc output và lifecycle

trước subscribe
trong pipeline: 10
next: 10
trong pipeline: 30
next: 30
complete
đã teardown
sau subscribe

of phát đồng bộ, nên toàn bộ next, complete và finalize xảy ra ngay bên trong lời gọi subscribe. Sau giá trị 30, take(2) gửi complete xuống Observer và hủy phần subscription phía trên; of không cần xử lý tiếp 4 và 5. Callback của finalize chạy trong teardown, rồi chương trình mới in sau subscribe.

Điểm cần nhớ không phải mọi Observable đều đồng bộ. interval, event DOM và HTTP phát ở thời điểm khác; lúc đó dòng sau subscribe thường chạy trước emission đầu tiên. Hãy đọc Marble diagram để nhìn trục thời gian của stream thay vì đoán từ thứ tự dòng code.

Thứ tự operator thay đổi kết quả

pipe áp dụng operator đúng theo thứ tự bạn viết. Đổi thứ tự không phải thao tác “format code”; nó có thể đổi cả số lượng emission lẫn thời điểm hoàn tất.

import { filter, of, take, toArray } from 'rxjs';

const source$ = of(1, 2, 3, 4);

source$
  .pipe(
    filter((value) => value % 2 === 0),
    take(2),
    toArray(),
  )
  .subscribe((values) => console.log('lọc trước:', values));

source$
  .pipe(
    take(2),
    filter((value) => value % 2 === 0),
    toArray(),
  )
  .subscribe((values) => console.log('take trước:', values));

// Output:
// lọc trước: [2, 4]
// take trước: [2]

Pipeline đầu lấy hai số chẵn. Pipeline sau chỉ nhìn hai giá trị đầu của nguồn, rồi mới lọc số chẵn. Khi thiết kế pipeline, hãy nói thành câu trước: “lọc rồi lấy hai” hay “lấy hai rồi lọc”. Câu nào đúng với yêu cầu thì thứ tự code nên phản ánh đúng câu đó.

Với operator tốn CPU, mình thường lọc dữ liệu không cần thiết càng sớm càng tốt. Nhưng semantics vẫn đứng trước tối ưu: đừng đẩy filter lên trên nếu việc đó đổi ý nghĩa nghiệp vụ.

Tái sử dụng pipeline

Dùng hàm pipe độc lập

Ngoài method source$.pipe(...), RxJS 7 còn export hàm pipe(...) độc lập. Hàm này ghép các unary function từ trái sang phải và trả về một function mới. Khi các function đó là pipeable operator, kết quả chính là một operator có thể tái sử dụng.

import { filter, from, map, type OperatorFunction, pipe } from 'rxjs';

type User = {
  id: number;
  name: string;
  active: boolean;
};

const toActiveNames: OperatorFunction<User, string> = pipe(
  filter((user: User) => user.active),
  map((user) => user.name.trim()),
);

const users: User[] = [
  { id: 1, name: ' An ', active: true },
  { id: 2, name: ' Bình ', active: false },
  { id: 3, name: ' Chi ', active: true },
];

from(users)
  .pipe(toActiveNames)
  .subscribe((name) => console.log(name));

// Output:
// An
// Chi

Hai dạng pipe khác nhau ở điểm bắt đầu:

  • source$.pipe(opA, opB) áp dụng chuỗi operator lên một source cụ thể và trả về Observable.
  • pipe(opA, opB) chưa có source; nó trả về function đã ghép để dùng với nhiều source.

Đây là lựa chọn tốt khi nhiều màn hình dùng cùng quy tắc nghiệp vụ. Đặt tên theo kết quả như toActiveNames giúp caller hiểu ý định tốt hơn một dãy operator lặp lại.

Viết operator factory nhỏ

Vì operator chỉ là function nhận Observable và trả Observable, bạn cũng có thể bọc cấu hình thành factory:

import { map, of, type OperatorFunction } from 'rxjs';

function multiplyBy(factor: number): OperatorFunction<number, number> {
  return map((value) => value * factor);
}

of(2, 4, 6)
  .pipe(multiplyBy(10))
  .subscribe((value) => console.log(value));

// Output:
// 20
// 40
// 60

Mặc định, hãy ghép operator có sẵn trước khi tự viết operator cấp thấp bằng new Observable(...). Cách trên giữ nguyên lifecycle và teardown của map, ít chỗ để tạo lỗi hơn.

Ví dụ thực tế với sự kiện DOM

Giả sử ô tìm kiếm chỉ nên gửi từ khóa có ít nhất hai ký tự và bỏ qua lần gõ trùng. fromEvent tạo nguồn; pipeline chuẩn hóa và lọc event trước khi consumer nhận nó.

import {
  distinctUntilChanged,
  filter,
  finalize,
  fromEvent,
  map,
  tap,
} from 'rxjs';

const searchInput = document.querySelector<HTMLInputElement>('#search');

if (!searchInput) {
  throw new Error('Không tìm thấy #search');
}

const searchQuery$ = fromEvent<InputEvent>(searchInput, 'input').pipe(
  map((event) =>
    (event.currentTarget as HTMLInputElement).value.trim(),
  ),
  filter((query) => query.length >= 2),
  distinctUntilChanged(),
  tap((query) => console.log('gửi truy vấn:', query)),
  finalize(() => console.log('search pipeline đã teardown')),
);

const subscription = searchQuery$.subscribe({
  next: (query) => {
    // Gọi lớp search hoặc cập nhật state tại ranh giới side effect này.
    console.log('consumer nhận:', query);
  },
  error: (error: unknown) => console.error(error),
});

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

Nếu người dùng lần lượt nhập r, rx, rx và rxjs, output trước lúc hủy là:

gửi truy vấn: rx
consumer nhận: rx
gửi truy vấn: rxjs
consumer nhận: rxjs

r bị filter loại; lần rx thứ hai bị distinctUntilChanged loại. Khi gọi unsubscribe(), subscription của fromEvent tháo event listener và finalize in thông báo teardown. Callback complete của Observer không chạy chỉ vì bạn unsubscribe thủ công; finalize mới là hook phù hợp nếu cleanup cần chạy cho cả complete, error và explicit unsubscribe.

Trong ứng dụng thật, quy tắc hủy subscription phụ thuộc framework và tuổi thọ nguồn. Xem Subscription và teardown trước khi gắn stream sống lâu vào component. Nếu bước kế tiếp là gọi API và hủy request cũ khi query đổi, switchMap là phần tiếp theo phù hợp.

Chọn operator theo ý định

Đừng cố nhớ toàn bộ danh sách operator. Trước tiên hãy mô tả việc cần làm bằng một động từ, rồi chọn đúng nhóm:

Ý địnhOperator hoặc function thường gặpCâu hỏi kiểm tra
Tạo nguồnof, from, fromEvent, intervalDữ liệu bắt đầu từ đâu?
Biến đổi từng emissionmapInput và output có kiểu hoặc hình dạng gì?
Chọn hoặc giới hạnfilter, take, distinctUntilChangedGiá trị nào được phép đi tiếp?
Quan sát side effecttapCó cần log/metric mà không đổi notification không?
Gom nhiều emissionscan, reduce, toArrayCần kết quả sau mỗi bước hay chỉ khi complete?
Làm việc với Observable lồngswitchMap, concatMap, mergeMap, exhaustMapHủy, xếp hàng, chạy song song hay bỏ qua công việc mới?
Xử lý lỗicatchError, retryPhục hồi, thử lại hay chuyển lỗi xuống consumer?
Kết thúc và cleanuptake, takeUntil, finalizePipeline dừng bằng điều kiện nào và giải phóng gì?

Mặc định, dùng map cho phép biến đổi thuần và để tap cho logging, metric hoặc side effect quan sát. Sự phân vai này giúp pipeline dễ test: nhìn thấy map là biết giá trị đầu ra có thể đổi; nhìn thấy tap là biết luồng dữ liệu đáng lẽ vẫn giữ nguyên.

Các bài Transformation operators, Filtering operators và Utility operators đi sâu vào từng nhóm. Trang Operators là bản đồ để chọn bài tiếp theo.

Những lỗi dễ mắc

Gọi pipe nhưng quên subscribe

import { map, of } from 'rxjs';

const doubled$ = of(1, 2, 3).pipe(
  map((value) => value * 2),
);

// Chưa có consumer; không có output nào được in.

pipe mô tả biến đổi, không phải lệnh “chạy ngay”. Bạn cần subscribe hoặc một consumer do framework quản lý. Riêng hot source có thể hoạt động độc lập với subscriber, nhưng gọi pipe vẫn không tự đăng ký vào source đó.

Dùng map cho side effect

import { map } from 'rxjs';

// Không nên: callback trả về kết quả của console.log, tức undefined.
map((value) => console.log(value));

Nếu chỉ muốn log, dùng tap((value) => console.log(value)). Nếu muốn đổi dữ liệu, map phải trả về giá trị mới một cách rõ ràng. Cũng tránh mutate object trong tap; subscriber phía sau sẽ thấy object đã bị đổi dù type không hề báo hiệu điều đó.

Cho rằng nhiều subscription dùng chung một lần chạy

import { defer } from 'rxjs';

const requestLike$ = defer(() => {
  console.log('khởi tạo nguồn');
  return Promise.resolve('ok');
});

requestLike$.subscribe();
requestLike$.subscribe();

// Output:
// khởi tạo nguồn
// khởi tạo nguồn

Với cold Observable, mỗi subscription thường tạo một execution riêng. pipe không tự cache hay multicast. Nếu cần chia sẻ một execution, hãy học rõ lifecycle của share thay vì thêm subscription rồi hy vọng RxJS tự dùng chung.

Subscribe bên trong subscribe

Nested subscribe tách lifecycle thành nhiều nhánh, khiến error và teardown khó theo dõi. Khi emission bên ngoài khởi động một Observable khác, hãy compose bằng higher-order operator phù hợp. Bài Nested subscribe chỉ ra cách nhận diện và sửa mẫu này.

Bỏ qua error channel

Lỗi bị ném trong callback của operator sẽ đi vào error channel và kết thúc subscription hiện tại. Đừng biến lỗi thành một giá trị giả nếu caller cần phân biệt thất bại; ngược lại, nếu nghiệp vụ có fallback rõ ràng, đặt catchError ở đúng cấp. Xem catchError để hiểu phạm vi bắt lỗi.

Operator order cũng quyết định phạm vi lỗi và cleanup

Vị trí của catchError, retry, higher-order operator và finalize có thể thay đổi phần pipeline được thử lại, thay thế hoặc teardown. Khi lifecycle quan trọng, đừng chỉ nhìn danh sách operator; hãy vẽ source nào subscribe vào source nào.

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

Dùng ví dụ of(1, 2, 3, 4, 5) ở trên và thử ba thay đổi nhỏ:

  1. Chuyển take(2) lên trước filter. Dự đoán output trước khi chạy.
  2. Chuyển tap lên trước map. So sánh giá trị mà tap nhìn thấy.
  3. Thay of(...) bằng interval(500) và giữ take(2). Quan sát dòng nào chạy đồng bộ, dòng nào chạy sau timer và lúc finalize xuất hiện.

Nếu dự đoán đúng cả notification lẫn thời điểm teardown, bạn đã nắm phần quan trọng nhất của pipe: đây không chỉ là cú pháp nối hàm, mà là cách tổ chức cả data flow và lifecycle.

Đi tiếp từ đây

Nguồn kiểm chứng API

On this page