Lỗi & hoàn tất
Thiết kế luồng lỗi, retry và cleanup rõ ràng trong RxJS.
Một request lỗi thì UI hiện thông báo, nhưng lần bấm tiếp theo không còn làm gì. Hoặc request đã bị hủy mà spinner vẫn quay. Hai triệu chứng này thường có cùng gốc: pipeline chưa phân biệt rõ lỗi, hoàn tất và hủy subscription, nên recovery hoặc cleanup được đặt sai chỗ.
Nhóm bài này giúp bạn quyết định một thất bại nên kết thúc phần nào của stream, khi nào được thử lại, và tài nguyên cần được dọn ở boundary nào. Đây là trang định hướng; các bài con sẽ đi sâu vào từng operator.
Trước khi đọc
Ví dụ dùng TypeScript và public API của RxJS 7.x, đối chiếu với RxJS 7.8.2. Bạn nên biết pipe(), subscribe() và Subscription, teardown. Các khoảng thời gian trong ví dụ chỉ để minh họa, không phải cấu hình production được khuyến nghị.
Mục lục
- Ba cách kết thúc cần tách riêng
- Chọn operator theo quyết định cần làm
- Recovery nằm ở boundary nào
- Một pipeline có retry và cleanup rõ ràng
- Retry không phải lúc nào cũng an toàn
- Cleanup theo từng lần thử và toàn tác vụ
- Lộ trình học trong nhóm
- Checklist trước khi đưa vào ứng dụng
- Bài tập tự kiểm tra
- Học tiếp
- Nguồn tham khảo
Ba cách kết thúc cần tách riêng
Bạn có thể hình dung một execution như một cuộc gọi: bên kia gửi thông tin, báo đã xong, hoặc báo lỗi; bạn cũng có thể tự cúp máy. Nhưng trong RxJS, “tự cúp máy” không tạo notification complete. Đó là giới hạn của phép so sánh và cũng là khác biệt cần giữ khi viết cleanup.
| Đường kết thúc | Ý nghĩa | Observer ở output nhận gì? | Teardown và finalize đã đăng ký chạy? |
|---|---|---|---|
complete | Output không còn giá trị để phát | Callback complete, không có error | Có |
error | Output kết thúc do lỗi chưa được recovery | Callback error, không có complete | Có |
unsubscribe | Consumer hoặc operator hủy subscription đang xét | Không tự gửi error hay complete | Có |
Bảng mô tả subscription ở boundary bạn đang quan sát. Operator có thể làm upstream và downstream kết thúc khác nhau: take(1) unsubscribe upstream sau giá trị đầu tiên nhưng gửi complete cho downstream. Vì vậy, câu “stream đã dừng” chưa đủ; bạn cần nói rõ đoạn nào đã dừng và do tín hiệu nào.
Contract của một execution
Một execution: next* (error | complete)?
Dữ liệu rồi hoàn tất: --a--b--|
Dữ liệu rồi lỗi: --a--b--X
Consumer hủy sớm: --a--b--! (không phải notification)
Nguồn còn sống: --a--b--c--...
| = complete, X = error, ! = unsubscribeMột execution có thể không phát giá trị nào, và cũng có thể không tự kết thúc. Nhưng khi đã gửi error hoặc complete, nó không gửi thêm next và không gửi terminal notification thứ hai.
Vậy retry có “hồi sinh” execution vừa lỗi không? Không. Nó subscribe lại vào upstream, tạo một lần thực thi mới nếu source hỗ trợ điều đó. Tương tự, catchError chuyển sang Observable thay thế; nó không bỏ qua phần tử lỗi rồi tiếp tục execution cũ.
Callback error trong subscribe() chỉ xử lý thông báo ở consumer boundary. Việc bạn log lỗi ở đó không làm source tiếp tục chạy, và try/catch quanh subscribe() không thay thế cơ chế xử lý error notification của RxJS.
Hoàn tất không đồng nghĩa thành công nghiệp vụ
of({ status: 'failed' }).subscribe(...) phát một giá trị mô tả thất bại rồi complete bình thường. EMPTY complete mà không phát giá trị nào. Nếu catchError trả về hai nguồn này, downstream có thể complete dù thao tác ban đầu thất bại.
Vì vậy, đừng dùng callback complete làm bằng chứng “đã lưu đơn hàng thành công”. Nó chỉ nói rằng output đã kết thúc bình thường. Kết quả nghiệp vụ nên nằm trong value hoặc một trạng thái có kiểu rõ ràng.
finalize không cho biết kết quả nghiệp vụ
finalize(() => ...) chạy cả khi complete, error lẫn unsubscribe và không nhận tham số lý do kết thúc. Dùng nó để dọn resource hoặc kết thúc trạng thái loading, không để hiện thông báo “thành công”.
Chọn operator theo quyết định cần làm
Mình thường bắt đầu bằng câu hỏi “sau thất bại, consumer cần thấy gì?” rồi mới chọn operator. Thêm catchError(() => EMPTY) vào mọi pipeline có thể làm log yên hơn, nhưng cũng khiến UI không biết phân biệt “không có dữ liệu” với “không lấy được dữ liệu”.
| Quyết định | Công cụ | Điều phải nhớ |
|---|---|---|
| Dừng và báo lỗi cho tầng gọi | Để lỗi truyền đi, hoặc throwError(() => reason) | Consumer cần có error handler phù hợp; không có recovery tự động. |
| Đổi sang giá trị hay source thay thế | catchError | Phải trả một ObservableInput; replacement có thể complete, lỗi hoặc sống lâu. |
| Thử lại sau lỗi | retry({ count, delay }) | Subscribe lại upstream; side effect và giá trị đã phát có thể lặp. |
| Đọc hoặc duy trì logic notifier trong pipeline cũ | retryWhen | Notifier điều khiển thời điểm retry; cần hiểu cả emission, completion và error của nó. |
| Chạy lại sau hoàn tất bình thường | repeat({ count, delay }) | Không bắt error; nguồn chưa complete thì chưa có lượt lặp tiếp theo. |
| Giới hạn thời gian chờ emission | timeout({ first, each }) | Mặc định báo TimeoutError; có thể chuyển source bằng with. |
| Dọn khi subscription kết thúc theo bất kỳ đường nào | finalize | Scope phụ thuộc vị trí đặt; không phải error handler hay recovery. |
Với code mới trên RxJS 7.8.2, mình ưu tiên retry({ count, delay }) cho retry hữu hạn có delay hoặc backoff vì chính sách nằm ngay trong cấu hình. Bạn vẫn cần hiểu retryWhen để đọc các pipeline dùng notifier; bài retry và retryWhen giải thích sâu hơn.
Có một khác biệt nhỏ nhưng dễ gây lệch số request: retry(2) cho tối đa ba lần subscribe — một lần đầu và hai lần thử lại. repeat(2) cho tổng cộng hai lần subscribe nếu source complete bình thường. Hai số 2 không cùng nghĩa.
Recovery nằm ở boundary nào
Một request độc lập có thể kết thúc khi lỗi. Một stream click hay search thường cần sống tiếp để nhận thao tác sau. Vì vậy, điều quan trọng hơn việc “đã có catchError chưa” là nó đang bắt lỗi của request con hay của toàn luồng sự kiện.
Giữ outer stream sống khi request con lỗi
Sơ đồ dưới đây mô tả recovery cục bộ trong một higher-order pipeline:
Sự kiện A ──► request A ──error──► catchError ──► kết quả failed
Sự kiện B ──► request B ──next───► kết quả ok
▲
└── outer stream vẫn còn được subscribeVí dụ độc lập này mô phỏng hai lần chọn mã hàng. Request đầu lỗi, nhưng lần chọn sau vẫn được xử lý:
import { catchError, of, Subject, switchMap, throwError } from 'rxjs';
type Result =
| { status: 'ok'; id: number }
| { status: 'failed'; id: number };
const selectedId$ = new Subject<number>();
function load$(id: number) {
return id === 1
? throwError(() => new Error('Không tải được mã 1'))
: of<Result>({ status: 'ok', id });
}
const subscription = selectedId$.pipe(
switchMap((id) =>
load$(id).pipe(
catchError(() => of<Result>({ status: 'failed', id })),
),
),
).subscribe({
next: (result) => console.log(result.status, result.id),
error: (reason: unknown) => console.error('Lỗi chưa được xử lý:', reason),
complete: () => console.log('complete'),
});
selectedId$.next(1);
selectedId$.next(2);
console.log('closed:', subscription.closed);
selectedId$.complete();Kết quả:
failed 1
ok 2
closed: false
completeỞ đây catchError nằm trong switchMap, nên nó thay request lỗi bằng một kết quả failed. Khi replacement complete, outer stream vẫn chờ mã hàng tiếp theo. Ví dụ bắt mọi lỗi chỉ để minh họa boundary; trong ứng dụng, hãy phân loại lỗi trước khi quyết định lỗi nào được chuyển thành value.
Nếu chuyển catchError ra sau switchMap, lỗi của request đầu sẽ kết thúc subscription vào selectedId$, rồi output chuyển sang replacement. Replacement hữu hạn complete xong thì pipeline không còn nghe lần chọn thứ hai. Bắt lỗi ngoài không “sửa” lại outer subscription đã đóng.
switchMap cũng hủy inner cũ khi có sự kiện mới. Nếu cần cleanup cho từng request, đặt finalize trong inner pipeline để cả request thành công, lỗi và request bị thay thế đều đi qua boundary đó.
Thứ tự operator thay đổi chính sách
Không có một thứ tự đúng cho mọi ứng dụng. Với tác vụ đọc có fallback sau khi thử lại hết lượt, hình dạng thường là:
source → timeout → retry → catchError → finalize → consumertimeout nằm trước retry để mỗi lần subscribe mới có đồng hồ chờ riêng. catchError nằm sau retry để fallback chỉ được chọn khi lỗi đã vượt chính sách thử lại. finalize ở cuối bao phủ toàn tác vụ, kể cả thời gian chờ retry và Observable fallback.
| Pipeline | Hành vi cần dự đoán |
|---|---|
source → retry → catchError | Thử lại trước, sau đó mới recovery nếu vẫn lỗi. |
source → catchError(() => of(fallback)) → retry | Lỗi đã thành value rồi complete, nên retry không thấy lỗi đó. |
source → catchError(() => EMPTY) → repeat(3) | Mỗi lỗi đã được đổi thành complete có thể kích hoạt repeat; nếu source luôn lỗi, nó có thể được gọi tổng cộng ba lần. |
source → finalize(perAttempt) → retry → finalize(perTask) | Cleanup upstream theo mỗi lần thử; cleanup cuối theo toàn subscription. |
Nếu catchError chuyển tiếp lỗi bằng throwError, retry phía sau vẫn có thể nhìn thấy lỗi. Chính kết quả trả về của recovery, chứ không chỉ sự hiện diện của catchError, quyết định điều này.
timeout({ first: 500 }) chỉ chờ giá trị đầu tiên; sau đó nó không kiểm tra những giá trị tiếp theo. each kiểm tra khoảng chờ giữa các emission và cả lần đầu nếu không có first. Nó không mặc định là deadline cho toàn bộ một tác vụ nhiều bước. Muốn giới hạn cả retry, delay và fallback, bạn phải thiết kế budget cho toàn tác vụ riêng.
Một pipeline có retry và cleanup rõ ràng
Giả sử bạn có tác vụ đọc dữ liệu: lỗi tạm thời được thử lại tối đa hai lần, lỗi hết lượt trở thành trạng thái unavailable, còn lỗi không thuộc nhóm dự kiến phải đi tới error handler. Ví dụ dùng nguồn mô phỏng thay vì HTTP để bạn nhìn thấy số lần subscribe mà không cần server.
Ví dụ chạy độc lập
import {
catchError,
defer,
finalize,
map,
mergeMap,
of,
retry,
throwError,
timer,
timeout,
TimeoutError,
} from 'rxjs';
type Data = { name: string };
type Result =
| { status: 'ok'; data: Data }
| { status: 'unavailable'; message: string };
class TemporaryError extends Error {}
let attempts = 0;
const failUntil = 2; // Mô phỏng hai lần đầu lỗi; mỗi lần chạy lại script sẽ reset.
const request$ = defer(() => {
const attempt = ++attempts;
console.log('attempt:', attempt);
return timer(100).pipe(
mergeMap(() =>
attempt <= failUntil
? throwError(() => new TemporaryError('Dịch vụ tạm thời chưa sẵn sàng'))
: of<Data>({ name: 'RxJS' }),
),
finalize(() => console.log('cleanup attempt:', attempt)),
);
});
function isRecoverable(reason: unknown): reason is Error {
return reason instanceof TemporaryError || reason instanceof TimeoutError;
}
const task$ = request$.pipe(
timeout({ first: 500 }),
retry({
count: 2,
delay: (reason: unknown) =>
isRecoverable(reason)
? timer(200)
: throwError(() => reason),
}),
map((data): Result => ({ status: 'ok', data })),
catchError((reason: unknown) =>
isRecoverable(reason)
? of<Result>({ status: 'unavailable', message: reason.message })
: throwError(() => reason),
),
finalize(() => console.log('cleanup task')),
);
const subscription = task$.subscribe({
next: (result) => console.log('result:', result),
error: (reason: unknown) => console.error('error:', reason),
complete: () => console.log('complete'),
});Đọc kết quả theo từng lớp
Với failUntil = 2, log cho thấy ba lần bắt đầu (attempt: 1, 2, 3), mỗi attempt được cleanup một lần. Consumer chỉ nhận một kết quả ok từ lần thứ ba, rồi complete; cleanup toàn task chạy một lần. Hai lỗi trước đó không tới callback error của consumer vì retry đã xử lý chúng.
defer chạy factory mỗi lần subscribe. Đây là điều kiện để retry khởi động lại công việc thay vì chỉ quan sát lại một kết quả cũ. Biến attempts trong ví dụ là bộ đếm mô phỏng dùng cho một lần chạy script, không phải state nên dùng chung cho các request production.
Nếu đổi failUntil thành 3, cả ba lần đều lỗi. retry hết lượt, catchError phát unavailable, rồi output complete. Nó không được xem là thành công tải dữ liệu; UI phải đọc trường status.
Nếu thay lỗi mô phỏng bằng một Error không thuộc nhóm recoverable, delay factory trả throwError, nên không có lần thử tiếp theo và lỗi được chuyển đến consumer. Cleanup vẫn chạy. Không dùng EMPTY ở đây: một retry notifier complete mà không emit sẽ khiến output complete, thay vì báo lỗi như bạn có thể mong đợi.
Để quan sát hủy, thêm dòng sau ngay sau subscribe():
setTimeout(() => subscription.unsubscribe(), 150);Trong kịch bản mặc định, lúc đó task thường đang chờ retry sau lỗi đầu. Việc unsubscribe hủy cả delay và task; không có callback complete hay error do hành động hủy này, nhưng cleanup task vẫn chạy. Thời điểm timer thực tế phụ thuộc event loop; điều cần kiểm tra là không có attempt mới sau khi subscription đã bị hủy.
Retry không phải lúc nào cũng an toàn
Retry là một lần subscribe mới, không phải transaction rollback. Nếu server đã tạo đơn hàng nhưng response bị mất, gửi lại request có thể tạo đơn hàng thứ hai. Vì vậy, mình chỉ bật retry mặc định cho thao tác đọc hoặc thao tác ghi đã có cơ chế idempotency đáng tin cậy.
Trước khi thử lại, bạn cần xác định:
- Lỗi nào có thể hồi phục? Lỗi mạng tạm thời khác với dữ liệu đầu vào sai, quyền truy cập bị từ chối hoặc bug trong code. Việc phân loại HTTP status phải theo contract của API.
- Công việc có thực sự chạy lại không?
from(existingPromise)chỉ quan sát lại Promise đã tồn tại. Đưa phần tạo Promise vàodefermới tạo công việc mới cho mỗi lần subscribe. - Giá trị cũ có bị lặp không? Những
nextđã đi downstream trước lỗi không được thu hồi. Source pháta, b, lỗi, rồi retry pháta, b, cthì consumer có thể nhậna, b, a, b, c. - Giới hạn là gì? Số lượt hữu hạn, delay, backoff và deadline phục vụ các mục tiêu khác nhau. Retry vô hạn không delay có thể gây vòng lặp nhanh với nguồn lỗi đồng bộ.
- Ai có quyền hủy? Consumer phải hủy được cả request đang chạy lẫn thời gian chờ lần thử kế tiếp.
Subscribe lại cũng không đảm bảo producer được reset. Một Subject đã error không trở lại trạng thái hoạt động chỉ vì có subscriber mới. Với nguồn hot, shared hoặc có cache, hãy kiểm tra semantics của chính source trước khi thêm retry.
Hủy subscription không hoàn tác công việc bên ngoài
timeout, switchMap và unsubscribe() có thể ngừng quan sát và chạy teardown, nhưng chỉ abort request thật nếu source tích hợp cancellation tương ứng. Chúng không tự hủy Promise, không rollback thay đổi trên server và không chứng minh server chưa xử lý request.
Cleanup theo từng lần thử và toàn tác vụ
Một task có ba lần thử thì có ít nhất hai scope hữu ích: resource của từng attempt và trạng thái của toàn task. Đặt chung mọi cleanup ở một chỗ thường khiến log, spinner hoặc resource ownership sai.
Toàn task: [ attempt 1 ] --delay-- [ attempt 2 ] --delay-- [ attempt 3 ]
Cleanup: ↑ ↑ ↑
↑
cleanup toàn taskfinalize trước retry nằm trong phần upstream được subscribe lại. Nó thích hợp để đóng resource của từng attempt hoặc ghi nhận attempt đã kết thúc. finalize sau retry phù hợp cho loading của toàn task, vì nó không chạy sau mỗi lỗi đang được thử lại.
Nếu còn fallback qua catchError, đặt cleanup toàn task sau cả operator này khi bạn muốn loading bao phủ fallback. Nếu fallback sống vô hạn, cleanup ở boundary cuối sẽ chỉ chạy khi nó kết thúc hoặc consumer unsubscribe; đây là hệ quả của lifecycle, không phải finalize bị mất.
Callback complete hoặc tap({ complete: ... }) không đủ cho cleanup vì nó không chạy trên đường error hay hủy. Với custom Observable, phần tạo timer, listener, socket hoặc request phải cung cấp teardown thực sự; thêm một finalize chỉ log ở cuối không tự dọn các resource đó.
Một boolean loading cũng không đại diện đúng cho nhiều request chạy đồng thời: request đầu kết thúc và đặt false trong khi request thứ hai còn chạy. Khi có concurrency, dùng state theo request hoặc bộ đếm tác vụ đang active thay vì để các finalizer tranh nhau ghi một biến chung.
Lộ trình học trong nhóm
Đọc theo thứ tự sau để hiểu contract trước, rồi mới thiết kế recovery và chính sách vận hành:
1. Error channel
Lỗi kết thúc execution nào, truyền qua operator ra sao và khác exception ở consumer thế nào?
2. catchError
Thay source, chuyển tiếp lỗi và chọn boundary để không làm chết outer stream.
3. retry và retryWhen
Thiết kế retry hữu hạn, delay, backoff và hiểu notifier điều khiển resubscription.
4. repeat
Subscribe lại sau complete và tránh nhầm polling với retry sau lỗi.
5. finalize
Đặt cleanup đúng scope trên cả complete, error lẫn cancellation.
6. Timeout và fallback
Giới hạn thời gian chờ emission và chọn đường lui không che mất trạng thái lỗi.
Nếu đang sửa lỗi “bấm lần hai không chạy”, bắt đầu với Error channel và catchError. Nếu UI kẹt loading, đọc finalize rồi kiểm tra owner của subscription. Nếu request bị gửi nhiều lần, kiểm tra retry và repeat cùng vị trí side effect trước khi thêm bất kỳ operator mới nào.
Checklist trước khi đưa vào ứng dụng
- Mỗi pipeline có boundary rõ ràng: một request, một phiên kết nối hay một stream event sống lâu?
- Lỗi thuộc loại nào: thất bại dự kiến, lỗi tạm thời hay lỗi cần truyền cho tầng gọi?
-
catchErrorcó trả source đúng kiểu và đặt đúng trong/ngoài inner pipeline không? - Fallback có giúp consumer phân biệt thất bại với dữ liệu hợp lệ không?
- Retry có số lượt, delay và điều kiện cụ thể; side effect lặp có an toàn không?
- Công việc mới có được tạo mỗi lần subscribe, thay vì dùng lại một Promise hoặc nguồn đã lỗi không?
- Timeout đang đo first emission, khoảng giữa emission hay budget toàn task?
- Cleanup cho từng attempt và toàn task đã tách đúng scope chưa?
- Consumer hủy được cả source đang chạy và retry delay; producer có teardown thực sự chưa?
- Test có cả complete, lỗi hết lượt và hủy sớm, thay vì chỉ happy path không?
Bài tập tự kiểm tra
Dùng hai ví dụ trên, dự đoán trước rồi mới chạy:
- Đưa
catchErrorcủa ví dụ mã hàng ra sauswitchMap. Mã2còn tạo kết quả không, vàsubscription.closedđổi thành gì? Mong đợi: không còn kết quả cho mã2, output đã complete vàclosedlàtrue. - Giữ retry
count: 2nhưng đổifailUntilthành3. Consumer nhận error hay value? Mong đợi: nhận valueunavailable, sau đó complete; source được subscribe tổng cộng ba lần. - Đổi lỗi mô phỏng thành
new Error('Lỗi không được retry'). Mong đợi: chỉ một attempt, callbackerrorchạy, callbackcompletekhông chạy; cleanup vẫn chạy. - Hủy task trước khi timer đầu phát. Mong đợi: không có
Result, không có terminal notification do hủy, cleanup attempt và cleanup task đều chạy. - Thử
of('done').pipe(repeat(2))rồi so với một source luôn lỗi quaretry(2). Mong đợi: repeat có hai execution; retry có ba execution rồi báo lỗi. Không dùngrepeat()vô hạn trên nguồn đồng bộ để thử điều này.
Sau khi quan sát, hãy chỉ ra operator nào đã thay đổi error thành value, operator nào đã tạo subscription mới, và finalizer nào thuộc từng scope. Nếu giải thích được ba điểm đó, bạn đã sẵn sàng áp dụng cùng cấu trúc vào HTTP hoặc stream UI thật.
Học tiếp
Vòng đời của stream
Củng cố khác biệt giữa notification, cancellation và finalization.
Higher-order Observables
Hiểu inner/outer subscription trước khi chọn recovery boundary cho UI.
Testing
Kiểm tra timing, số lần subscribe và đường hủy bằng test có thể lặp lại.
Nguồn tham khảo
- RxJS — Observable guide: contract của notification và disposal.
- RxJS 7.8.2 — catchError: chuyển sang Observable thay thế sau lỗi.
- RxJS 7.8.2 — retry: số lần resubscribe, delay notifier và các emission được chuyển tiếp.
- RxJS 7.8.2 — repeat: lặp sau complete và semantics của
count. - RxJS 7.8.2 — finalize: cleanup trên complete, error và unsubscribe.
- RxJS 7.8.2 — timeout:
first,each,withvàTimeoutError.