Học RxJS
Nền tảng

Vòng đời của stream

Các trạng thái next, error, complete và quy tắc kết thúc stream.

Một request đã trả dữ liệu nhưng loading vẫn quay, hoặc một màn hình đã đóng mà timer vẫn chạy: cả hai thường bắt nguồn từ việc chưa phân biệt rõ phát giá trị, kết thúc và hủy đăng ký. Với RxJS, bạn nên theo dõi vòng đời của từng lần subscribe(), thay vì xem Observable như một chiếc hộp có một trạng thái chung.

Phạm vi phiên bản

Các API và hành vi trong bài áp dụng cho RxJS 7.x. Ví dụ dùng cách import công khai từ rxjs, không dựa vào API thử nghiệm của RxJS 8.

Mục lục

Mental model: vòng đời thuộc về từng subscription

Observable mô tả cách tạo và gửi dữ liệu. Mỗi lần gọi subscribe(), RxJS tạo một Observable execution và trả về một Subscription đại diện cho execution đó. Với cold Observable, hai lần subscribe thường chạy producer hai lần; với nguồn được share, nhiều subscriber có thể dùng chung producer. Dù nguồn thuộc loại nào, mỗi subscription vẫn có quyền được đóng riêng.

Hãy hình dung Observable là công thức, còn một lần subscribe() là một lần bắt đầu nấu. Công thức không “hoàn tất”; lần nấu cụ thể mới hoàn tất, gặp lỗi hoặc bị dừng. Phép so sánh này dừng ở đó: hot Observable có thể phát từ một producer đã chạy sẵn, nên không phải nguồn nào cũng đợi subscriber đầu tiên mới bắt đầu.

Một execution có ba đường kết thúc

Sơ đồ dưới đây mô tả một subscription. Repo chưa cấu hình Mermaid, nên bài dùng ASCII để sơ đồ hiển thị đúng mà không cần thêm dependency.

                         next(value)
                      ┌───────────────┐
                      │               │
subscribe() ─────► [ ACTIVE ] ◄───────┘
                      │
          ┌───────────┼──────────────────┐
          │           │                  │
       error(err)   complete()      unsubscribe()
          │           │                  │
          ▼           ▼                  ▼
      [ CLOSED ]   [ CLOSED ]        [ CLOSED ]
       do lỗi       bình thường        do hủy
          │           │                  │
          └───────────┴──────────────────┘
                      ▼
              chạy teardown đã đăng ký

Ba đường đều làm Subscription.closed trở thành true và kích hoạt cleanup đã đăng ký. Khác biệt quan trọng là error và complete là notification gửi đến Observer, còn unsubscribe() là hành động hủy từ phía consumer và không gửi notification complete.

Một nguồn có thể không bao giờ tự kết thúc. interval(), DOM event hay WebSocket thường cứ active cho đến khi consumer hủy, một operator giới hạn vòng đời, hoặc nguồn thực sự đóng.

Observable contract

Contract của một execution có thể viết gọn như sau:

next* (error | complete)?

Điều đó có nghĩa:

  • next có thể xuất hiện từ 0 đến nhiều lần;
  • error hoặc complete có thể xuất hiện tối đa một lần;
  • error và complete loại trừ nhau trong cùng một execution;
  • sau terminal notification, mọi notification tiếp theo đều không được chuyển đến Observer;
  • terminal notification là tùy chọn: một stream vô hạn hợp lệ có thể không error cũng không complete.

Contract này giúp bạn suy luận pipeline mà không cần biết producer bên trong. Nếu đã thấy error, đừng chờ complete; nếu đã thấy complete, đừng chờ thêm next.

Ba loại notification

next: gửi dữ liệu

next(value) chuyển một giá trị đến callback next của Observer. Nó không nói rằng công việc đã xong và cũng không đóng subscription. Một stream có thể gửi nhiều kiểu giá trị theo thời gian, nhưng TypeScript thường giữ chúng trong cùng kiểu generic Observable<T>.

import { of } from 'rxjs';

of('đang chờ', 'đang xử lý', 'đã lưu').subscribe({
  next: (status) => console.log('next:', status),
});

Kết quả:

next: đang chờ
next: đang xử lý
next: đã lưu

of() còn gửi complete sau giá trị cuối, nhưng vì Observer trên không khai báo callback complete, bạn không thấy dòng log tương ứng.

error: kết thúc do lỗi

error(err) vừa báo lỗi vừa đóng execution. Nó không phải một giá trị next đặc biệt mà pipeline có thể bỏ qua rồi chạy tiếp. Muốn chuyển sang nguồn dự phòng, thử lại hoặc biểu diễn lỗi như dữ liệu, bạn phải chọn operator và mô hình dữ liệu tương ứng.

Nên cung cấp callback error ở boundary nơi ứng dụng subscribe, hoặc xử lý lỗi có chủ đích trong pipeline. Nếu lỗi đi đến một subscription không có error handler, RxJS sẽ báo nó theo cơ chế unhandled error của runtime; try/catch đặt quanh một callback bất đồng bộ không phải cách xử lý đáng tin cậy.

complete: kết thúc bình thường

complete() báo rằng execution sẽ không gửi thêm dữ liệu. Notification này không mang payload. Nếu cần trả “kết quả cuối”, hãy gửi nó bằng next(result) trước rồi mới complete().

complete không đồng nghĩa với “thành công nghiệp vụ”. Ví dụ một stream tìm kiếm có thể next({ kind: 'not-found' }) rồi complete bình thường. Ngược lại, việc người dùng rời màn hình thường chỉ cần hủy subscription, không cần giả vờ rằng nguồn đã hoàn tất.

Theo dõi vòng đời qua ví dụ

Hoàn tất đồng bộ

of() là nguồn đồng bộ: toàn bộ next và complete xảy ra ngay trong lời gọi subscribe(). Vì vậy subscription đã đóng trước khi dòng kế tiếp chạy.

import { of } from 'rxjs';

const subscription = of(10, 20, 30).subscribe({
  next: (value) => console.log('next:', value),
  error: (error: unknown) => console.error('error:', error),
  complete: () => console.log('complete'),
});

console.log('closed:', subscription.closed);

Kết quả:

next: 10
next: 20
next: 30
complete
closed: true

Điểm dễ bỏ sót là Observable không mặc định bất đồng bộ. Scheduler và nguồn dữ liệu quyết định notification xuất hiện đồng bộ hay bất đồng bộ; vòng đời vẫn tuân theo cùng contract.

Lỗi dừng toàn bộ execution

Operator cũng phải giữ contract. Nếu hàm trong map() ném lỗi, RxJS chuyển lỗi đó sang error channel, đóng subscription và không xử lý phần tử kế tiếp.

import { map, of } from 'rxjs';

type User = { id: number };

const userTexts$ = of('{"id":1}', '{broken-json}', '{"id":2}');

userTexts$
  .pipe(map((text) => JSON.parse(text) as User))
  .subscribe({
    next: (user) => console.log('user:', user.id),
    error: (error: unknown) => {
      const message = error instanceof Error ? error.message : String(error);
      console.log('error:', message);
    },
    complete: () => console.log('complete'),
  });

Kết quả có dạng:

user: 1
error: ...

Thông điệp cụ thể của JSON.parse tùy JavaScript runtime, nhưng hai điều không đổi: user có id bằng 2 không được phát và callback complete không chạy. catchError có thể thay execution đã lỗi bằng một Observable khác; nó không làm execution cũ sống lại. Phần Error channel phân tích kỹ hơn ranh giới này.

Hủy đăng ký không phải complete

Ví dụ sau tạo một producer bằng setInterval. Consumer hủy sau khoảng 2,25 giây, nên timer dừng và teardown chạy, nhưng callback complete không chạy.

import { Observable } from 'rxjs';

const ticks$ = new Observable<number>((subscriber) => {
  let tick = 0;
  const intervalId = setInterval(() => {
    subscriber.next(tick++);
  }, 1_000);

  return () => {
    clearInterval(intervalId);
    console.log('teardown');
  };
});

const subscription = ticks$.subscribe({
  next: (value) => console.log('next:', value),
  complete: () => console.log('complete'),
});

setTimeout(() => subscription.unsubscribe(), 2_250);

Kết quả thường là:

next: 0
next: 1
teardown

Thời điểm log có thể xê dịch theo event loop, nhưng sẽ không có complete. Nếu bạn đặt cleanup chỉ trong callback complete, đường hủy này sẽ bỏ qua cleanup đó.

Teardown và finalize

Teardown giải phóng tài nguyên của producer

Khi tự tạo Observable, hàm truyền vào constructor có thể trả về teardown function. Đây là nơi producer thu hồi thứ mà chính nó đã tạo: timer, event listener, socket, worker hoặc request có thể abort. RxJS gọi teardown khi subscription kết thúc bởi complete, error hoặc unsubscribe().

import { Observable } from 'rxjs';

const value$ = new Observable<number>((subscriber) => {
  console.log('producer: start');

  subscriber.next(42);
  subscriber.complete();

  return () => console.log('producer: teardown');
});

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

Kết quả:

producer: start
next: 42
complete
producer: teardown

Đừng dùng complete handler làm cleanup duy nhất

Callback complete chỉ phản ứng với terminal notification bình thường. Nó không chạy khi consumer gọi unsubscribe(). Tài nguyên của producer phải được giải phóng trong teardown; cleanup ở cấp pipeline nên dùng finalize().

finalize đặt cleanup ở pipeline

finalize(callback) chạy khi subscription kết thúc do complete, error hoặc bị unsubscribe rõ ràng. Nó phù hợp cho việc tắt loading, ghi log kết thúc một operation hoặc dọn state gắn với pipeline mà không sở hữu trực tiếp producer.

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

interval(500)
  .pipe(
    take(3),
    finalize(() => console.log('finalize')),
  )
  .subscribe({
    next: (value) => console.log('next:', value),
    complete: () => console.log('complete'),
  });

Sau khoảng 1,5 giây:

next: 0
next: 1
next: 2
complete
finalize

take(3) complete output sau ba giá trị và hủy subscription ngược lên interval, nên timer nội bộ được dọn. Nếu consumer unsubscribe sớm, finalize vẫn chạy nhưng callback complete không chạy. Xem bài finalize để xử lý các trường hợp cleanup phức tạp hơn.

Ví dụ thực tế: vòng đời event listener

DOM event là nguồn không tự complete. Khi màn hình gắn listener nhưng quên hủy lúc unmount, callback cũ có thể tiếp tục giữ tham chiếu đến state và DOM. fromEvent() đóng gói việc thêm listener khi subscribe và gỡ listener khi teardown.

Ví dụ dưới đây dùng destroy$ làm tín hiệu kết thúc vòng đời màn hình:

import { Subject, finalize, fromEvent, map, takeUntil } from 'rxjs';

const saveButton = document.querySelector<HTMLButtonElement>('#save-button');

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

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

const saveClicks$ = fromEvent<MouseEvent>(saveButton, 'click').pipe(
  map((event) => ({ x: event.clientX, y: event.clientY })),
  takeUntil(destroy$),
  finalize(() => console.log('Đã gỡ vòng đời saveClicks$')),
);

saveClicks$.subscribe({
  next: (position) => console.log('Lưu tại:', position),
  complete: () => console.log('Màn hình đã kết thúc'),
});

// Gọi khi màn hình bị tháo khỏi UI.
destroy$.next();
destroy$.complete();

Khi destroy$.next() phát tín hiệu, takeUntil() complete output, unsubscribe khỏi fromEvent() và listener được gỡ. Dòng destroy$.complete() sau đó đóng chính notifier để không giữ nó active nữa; chỉ gọi destroy$.complete() mà không next() sẽ không kích hoạt takeUntil().

Nếu framework đã có primitive quản lý vòng đời, hãy ưu tiên primitive đó thay vì tự tạo destroy$ cho mọi component. Mục tiêu không phải là dùng nhiều Subject, mà là làm cho thời điểm kết thúc subscription trở nên rõ ràng và kiểm soát được.

Chọn cách kết thúc phù hợp

Nhu cầuCách mặc định nên chọnVì sao
Producer đã phát xong hữu hạn dữ liệuĐể nguồn gọi complete()Observer nhận được terminal notification bình thường.
Producer gặp lỗi không thể tiếp tụcGọi error(err) hoặc để operator chuyển exception sang error channelContract đảm bảo không có dữ liệu lẫn lộn sau lỗi.
Consumer không còn cần dữ liệuunsubscribe() hoặc primitive lifecycle của frameworkHủy đúng execution mà consumer sở hữu; không giả mạo complete của nguồn.
Dừng sau một điều kiện có thể mô tả trong pipelinetake(), takeUntil() hoặc operator phù hợpVòng đời nằm trong pipeline, dễ đọc và compose hơn hủy thủ công.
Dọn tài nguyên do custom producer tạoTrả teardown function từ new Observable()Cleanup ở cùng nơi tài nguyên được cấp phát.
Chạy cleanup bất kể complete, error hay hủyfinalize()Bao phủ cả ba đường đóng subscription.
Chạy lại sau error hoặc completeretry() sau error, repeat() sau completeMỗi lần chạy lại là một subscription/execution mới, không hồi sinh execution đã đóng.

Mặc định, mình ưu tiên operator như take() hoặc takeUntil() khi điều kiện kết thúc là một phần của luồng dữ liệu, vì ý định nằm ngay cạnh pipeline. Dùng unsubscribe() trực tiếp khi lifecycle thực sự mang tính mệnh lệnh, chẳng hạn cleanup của một lớp hoặc màn hình không có primitive tích hợp.

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

  1. Chờ complete từ nguồn vô hạn. interval(), click và WebSocket có thể chạy mãi. Nếu UI phụ thuộc vào complete để tắt loading, UI cũng có thể chờ mãi. Hãy định nghĩa điều kiện kết thúc hoặc trạng thái loading riêng.
  2. Cho rằng unsubscribe sẽ gọi complete handler. Nó chỉ đóng subscription và chạy teardown/finalizer. Nếu cần cùng một cleanup cho cả ba đường, dùng finalize().
  3. Cố phát tiếp sau error hoặc complete. RxJS bỏ qua notification gửi đến subscriber đã đóng. Code producer vẫn có thể tiếp tục tiêu tốn CPU nếu bạn không dừng tài nguyên bên ngoài, nên teardown vẫn bắt buộc.
  4. Xử lý error rồi mong nguồn cũ tiếp tục. catchError() chuyển sang Observable thay thế; retry() subscribe lại. Cả hai tạo hướng đi mới sau khi execution lỗi đã kết thúc.
  5. Hủy thủ công bên trong nguồn đồng bộ. of() có thể phát trước khi biến nhận giá trị trả về từ subscribe(). Thay vì cố gọi subscription.unsubscribe() trong callback next, hãy dùng take(1) hoặc operator mô tả đúng điều kiện dừng.
  6. Quên rằng mỗi subscription có vòng đời riêng. Subscribe hai lần vào cold Observable có thể tạo hai timer hoặc hai HTTP request. Nếu muốn chia sẻ producer, đó là quyết định multicasting riêng, không phải thuộc tính mặc định của vòng đời.

Quy tắc đặt cleanup

Tài nguyên được tạo ở đâu thì đặt teardown gần đó. Trạng thái UI gắn với cả pipeline thì đặt trong finalize(). Cách chia này giúp bạn không phụ thuộc vào việc stream kết thúc bằng complete, error hay hủy.

Checklist khi đọc một stream

Trước khi đưa pipeline vào production, hãy trả lời được các câu hỏi sau:

  • Nguồn bắt đầu khi nào, và mỗi lần subscribe có tạo execution mới không?
  • Nguồn hữu hạn hay có thể chạy mãi?
  • Điều gì phát complete, điều gì phát error, và ai có quyền hủy?
  • Tài nguyên nào được cấp phát khi subscribe?
  • Teardown có chạy trên cả complete, error và unsubscribe không?
  • Error được xử lý ở đâu? Sau error có fallback hay một subscription mới không?
  • Có khả năng tạo nhiều subscription ngoài ý muốn không?

Nếu một câu chưa có đáp án, hãy làm rõ lifecycle trước khi thêm operator. Debug dữ liệu sai thường dễ hơn debug một execution không ai biết ai chịu trách nhiệm đóng.

Nguồn tham khảo

Học tiếp

On this page