Học RxJS
Observable & Subscription

Tự tạo Observable

Đóng gói producer tùy biến bằng constructor Observable.

Một SDK chỉ đưa cho bạn callback, một browser API trả về ID để huỷ, hoặc một resource cần mở khi có consumer và đóng ngay khi consumer rời đi. Đây là lúc new Observable(...) hữu ích: bạn tự định nghĩa cách producer khởi động, phát notification và teardown. Đổi lại, bạn cũng phải giữ đúng lifecycle contract mà các creation function có sẵn đã xử lý giúp mình.

Phạm vi phiên bản

Bài này dùng public API của RxJS 7.x và đối chiếu hành vi với RxJS 7.8.2. Ví dụ viết bằng TypeScript; các API thử nghiệm của RxJS 8 không nằm trong phạm vi bài. Bạn nên nắm Observer và subscribe trước khi đọc. Ví dụ timer chạy được trong Node hoặc browser; ví dụ Geolocation cần browser và DOM types.

Mục lục

Khi nào cần tự tạo Observable

Mặc định, mình sẽ tìm creation function có sẵn trước. of, from, defer, fromEvent, interval và timer diễn đạt ý định rõ hơn, đồng thời đã xử lý nhiều chi tiết lifecycle mà custom code dễ bỏ sót. Constructor chỉ đáng dùng khi nguồn dữ liệu có contract riêng mà RxJS chưa có adapter phù hợp.

Nguồn đang cóLựa chọn mặc địnhCó cần new Observable không?
Array, iterable, Promise, ReadableStream hoặc Observable-compatible inputfrom(...)Không
Một nhóm giá trị tĩnhof(...)Không
DOM EventTarget hoặc event API tương thíchfromEvent(...)Thường không
Giá trị phải được tạo lại tại thời điểm subscribedefer(...)Thường không
Timer theo chu kỳ hoặc một lầninterval(...), timer(...)Không
API đăng ký và gỡ event bằng hai function riêngfromEventPattern(...)Thường không
Callback một lần theo kiểu Node (error, result)bindNodeCallback(...)Thường không, nhưng phải kiểm tra khả năng huỷ của API gốc
SDK dùng callback và trả hàm huỷ riêngnew Observable(...)Có thể phù hợp
Resource cần open khi subscribe và close khi dừngnew Observable(...)Có thể phù hợp

Constructor là API cấp thấp

Nếu một creation function hoặc tổ hợp operator đã diễn đạt được nguồn, hãy dùng chúng. Tự tạo Observable không làm stream “RxJS hơn”; nó chỉ chuyển trách nhiệm setup, error và teardown sang code của bạn.

Giải phẫu constructor Observable

Dạng rút gọn của constructor như sau:

import { Observable } from 'rxjs';

type Value = { id: string };

const source$ = new Observable<Value>((subscriber) => {
  // 1. Setup producer.
  // 2. Gọi subscriber.next(...), error(...) hoặc complete().
  // 3. Trả teardown logic nếu đã cấp phát resource.

  return () => {
    // Cleanup producer.
  };
});

Luồng đời có thể đọc bằng sơ đồ nhỏ này:

new Observable(...)      subscribe()              terminal hoặc huỷ
      │                       │                           │
      ▼                       ▼                           ▼
chỉ tạo mô tả ───────► chạy subscribe function ──► đóng Subscription
                              │                           │
                              ├─ next(value)             ▼
                              └─ giữ resource        chạy teardown

Subscribe function chạy cho từng subscription

Callback truyền vào constructor thường được gọi là subscribe function. Tạo source$ chưa chạy callback này; mỗi lời gọi subscribe() mới tạo một execution.

import { Observable } from 'rxjs';

const requestId$ = new Observable<string>((subscriber) => {
  const requestId = crypto.randomUUID();
  console.log('setup:', requestId);

  subscriber.next(requestId);
  subscriber.complete();

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

requestId$.subscribe((id) => console.log('A nhận:', id));
requestId$.subscribe((id) => console.log('B nhận:', id));

Hai subscriber nhận hai requestId khác nhau vì subscribe function chạy hai lần. Teardown cũng thuộc về từng execution, không phải thuộc về biến requestId$ nói chung.

Tên biến subscriber ở đây có chủ ý. Nó là Subscriber<T> do RxJS đưa cho phía producer, không phải Observer object mà code ứng dụng truyền vào subscribe().

Subscriber gửi ba loại notification

Producer giao tiếp với consumer qua ba phương thức:

  • subscriber.next(value) gửi một giá trị và có thể gọi nhiều lần;
  • subscriber.error(error) báo thất bại và đóng execution;
  • subscriber.complete() báo hoàn tất bình thường và đóng execution.

Contract có thể viết ngắn gọn là:

next* (error | complete)?

Sau error hoặc complete, những notification tiếp theo không được chuyển tới Observer. Nhưng có một chi tiết quan trọng: gọi terminal notification không tự dừng các câu lệnh JavaScript phía sau nó.

import { Observable } from 'rxjs';

const source$ = new Observable<number>((subscriber) => {
  subscriber.next(1);
  subscriber.complete();

  subscriber.next(2); // Observer không nhận giá trị này.
  console.log('producer vẫn chạy tới đây');
});

Vì vậy, khi một nhánh gọi error hoặc complete, hãy return nếu producer không còn việc gì cần làm. Với loop đồng bộ, kiểm tra thêm subscriber.closed để ngừng tính toán sớm.

Teardown đóng resource

Subscribe function có thể trả về teardown logic. Trường hợp dễ đọc nhất là một hàm không nhận tham số:

return () => {
  clearInterval(intervalId);
  socket.close();
  sdk.removeListener(listener);
};

RxJS chạy teardown khi execution đi tới complete, error, hoặc khi consumer gọi unsubscribe(). Setup và cleanup nên đối xứng: code nào tạo timer thì biết cách xoá timer; code nào đăng ký listener thì biết cách gỡ listener.

Ngừng notification chưa phải là cleanup

Sau unsubscribe, Subscriber đóng nên giá trị mới không tới Observer. Nhưng timer, listener hoặc socket vẫn có thể chạy nền nếu teardown không thật sự giải phóng chúng. “Không còn thấy log” chưa chứng minh resource đã được dọn.

Trong RxJS 7.8.2, TeardownLogic gồm void, một hàm cleanup, Subscription hoặc object có unsubscribe(). Dù vậy, trả một hàm cleanup thường rõ nhất khi bạn đang bọc API callback. Teardown là hành động cleanup, không phải callback complete và cũng không phải giá trị để phát bằng next.

Teardown khi source kết thúc đồng bộ

Ở ví dụ requestId$, complete() xảy ra trước khi subscribe function trả hàm teardown. Vậy cleanup có bị mất không? Không: sau khi subscribe function trả về, RxJS thêm kết quả đó vào Subscriber bằng add(). Nếu Subscriber đã đóng, add() chạy teardown ngay.

import { Observable } from 'rxjs';

const source$ = new Observable<number>((subscriber) => {
  console.log('setup');
  subscriber.next(1);
  subscriber.complete();
  console.log('sau complete trong producer');

  return () => console.log('cleanup');
});

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

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

Output theo đúng thứ tự:

setup
next: 1
complete
sau complete trong producer
cleanup
closed: true

Chi tiết này cũng áp dụng khi SDK gọi callback ngay trong lúc đăng ký, trước khi trả handle huỷ. Một hàm teardown trả về bình thường vẫn được chạy dù callback đã làm subscription đóng. Nhưng nếu setup throw trước khi tới return, RxJS chưa nhận được hàm đó; phần cleanup khi setup thất bại xử lý trường hợp khác này.

Ví dụ đầu tiên với countdown

Ví dụ sau tạo countdown phát giá trị đầu tiên ngay lập tức, sau đó giảm mỗi giây. Nó cố ý dùng constructor để lộ toàn bộ lifecycle; trong code thật, timer kết hợp operator thường ngắn hơn.

import { Observable } from 'rxjs';

function countdown(seconds: number): Observable<number> {
  return new Observable<number>((subscriber) => {
    if (!Number.isInteger(seconds) || seconds < 0) {
      subscriber.error(
        new RangeError('seconds phải là số nguyên không âm'),
      );
      return;
    }

    let remaining = seconds;
    subscriber.next(remaining);

    // take(1) có thể đã đóng upstream ngay trong next ở trên.
    if (subscriber.closed) {
      return;
    }

    if (remaining === 0) {
      subscriber.complete();
      return;
    }

    const intervalId = setInterval(() => {
      remaining -= 1;
      subscriber.next(remaining);

      if (remaining === 0) {
        subscriber.complete();
      }
    }, 1_000);

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

const subscription = countdown(5).subscribe({
  next: (value) => console.log('còn:', value),
  error: (error: unknown) => console.error('lỗi:', error),
  complete: () => console.log('hết giờ'),
});

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

Output thường dừng quanh đây:

còn: 5
còn: 4
còn: 3
teardown countdown

hết giờ không xuất hiện vì consumer huỷ trước khi producer complete. Dù vậy, teardown vẫn xoá interval. Nếu để countdown chạy tới 0, producer gọi complete() và RxJS cũng chạy chính teardown đó. Timer không đảm bảo thời điểm chính xác tuyệt đối; output trên giả sử event loop không bị chặn lâu.

Kiểm tra subscriber.closed sau emission đầu giúp tránh tạo interval nếu downstream đã dừng bằng take(1). Với countdown(0), source phát 0 rồi complete ngay và không tạo timer, nên cũng không có log teardown countdown.

Ví dụ còn cho thấy validation nằm bên trong subscribe function. Nhờ vậy, countdown(-1) vẫn chỉ tạo mô tả; lỗi được gửi qua error channel khi có consumer thật sự subscribe.

Bọc một callback API thực tế

Custom Observable phát huy giá trị nhất tại ranh giới giữa RxJS và một API imperative. Browser Geolocation là ví dụ phù hợp: watchPosition nhận callback, trả một watchId, còn clearWatch(watchId) dùng để dừng theo dõi.

Có API chưa chắc đã có quyền truy cập

Geolocation cần secure context, thường là HTTPS hoặc localhost, và quyền của người dùng; Permissions Policy cũng có thể chặn nó. Kiểm tra navigator.geolocation chỉ phát hiện API có tồn tại, không đảm bảo lấy được vị trí. Không dùng ví dụ này trong code chạy trên server để kỳ vọng nó trả toạ độ.

Chuyển Geolocation thành Observable

import { Observable } from 'rxjs';

function watchPosition$(
  options?: PositionOptions,
): Observable<GeolocationPosition> {
  return new Observable<GeolocationPosition>((subscriber) => {
    if (typeof navigator === 'undefined' || !navigator.geolocation) {
      subscriber.error(
        new Error('Môi trường này không hỗ trợ Geolocation'),
      );
      return;
    }

    const geolocation = navigator.geolocation;
    let watchId: number;

    try {
      watchId = geolocation.watchPosition(
        (position) => subscriber.next(position),
        (error) => subscriber.error(error),
        options,
      );
    } catch (error) {
      subscriber.error(error);
      return;
    }

    return () => {
      geolocation.clearWatch(watchId);
    };
  });
}

Consumer vẫn dùng API RxJS bình thường:

const position$ = watchPosition$({
  enableHighAccuracy: true,
  maximumAge: 10_000,
  timeout: 15_000,
});

const subscription = position$.subscribe({
  next: (position) => {
    const { latitude, longitude } = position.coords;
    console.log({ latitude, longitude });
  },
  error: (error: unknown) => {
    console.error('Không lấy được vị trí:', error);
  },
});

// Gọi hàm này trong lifecycle cleanup của màn hình hoặc component.
function stopTracking(): void {
  subscription.unsubscribe();
}

stopTracking() thể hiện quyền sở hữu chứ không được gọi ngay sau subscribe. Hãy gọi nó khi màn hình bị tháo, người dùng tắt tính năng định vị, hoặc một điều kiện nghiệp vụ kết thúc stream.

Đọc lifecycle của ví dụ

Mỗi lần subscribe đi qua bốn bước:

  1. kiểm tra môi trường có Geolocation hay không;
  2. đăng ký hai callback với watchPosition;
  3. chuyển vị trí thành next và lỗi từ browser thành error;
  4. trả teardown gọi clearWatch đúng với watchId vừa tạo.

Nếu browser báo lỗi qua callback, subscriber.error(error) đóng execution và teardown gỡ watch. Đây là quyết định của adapter: mọi lỗi Geolocation, kể cả timeout tạm thời, đều kết thúc stream này. Browser API không bắt buộc mọi lỗi phải là terminal; nếu sản phẩm cần tiếp tục watch sau lỗi tạm thời, hãy thiết kế kiểu dữ liệu biểu diễn cả vị trí và trạng thái lỗi rồi phát qua next.

Nếu consumer rời màn hình trước, unsubscribe() cũng chạy teardown nhưng không gọi callback complete. Source này không tự complete khi nhận vị trí đầu tiên; dùng take(1) nếu chỉ cần một vị trí. clearWatch dừng watch do adapter sở hữu, không cam kết tắt mọi hoạt động định vị của browser hoặc thu hồi quyền người dùng đã cấp.

Adapter chỉ làm nhiệm vụ chuyển contract: callback vào notification, handle huỷ vào teardown. Việc format toạ độ, lọc độ chính xác hoặc debounce nên nằm trong pipe(...), không nên nhồi vào constructor. Ranh giới hẹp như vậy giúp adapter dễ test và tái sử dụng.

Xử lý lỗi và dừng producer đúng cách

Lỗi đồng bộ và lỗi bất đồng bộ

Với cấu hình mặc định, RxJS 7 bắt lỗi đồng bộ bị throw ngay trong subscribe function và chuyển nó tới subscriber.error. Tuy nhiên, viết subscriber.error(...) rõ ràng vẫn tốt hơn cho lỗi dự kiến, vì người đọc thấy ngay đây là nhánh terminal của producer. Điều này không có nghĩa RxJS tự thu hồi mọi resource bạn vừa mở trước chỗ throw.

Lỗi xảy ra trong callback bất đồng bộ chạy trên một call stack khác. Constructor không thể bắt hộ lỗi đó; callback phải tự chuyển lỗi vào error channel. Ở consumer, hãy cung cấp handler error hoặc xử lý bằng catchError khi lỗi là điều có thể xảy ra; không truyền handler không có nghĩa lỗi tự biến mất, mà RxJS mặc định báo nó như unhandled error.

import { Observable } from 'rxjs';

function parseLater(input: string): Observable<unknown> {
  return new Observable((subscriber) => {
    const timeoutId = setTimeout(() => {
      try {
        subscriber.next(JSON.parse(input));
        subscriber.complete();
      } catch (error) {
        subscriber.error(error);
      }
    }, 0);

    return () => clearTimeout(timeoutId);
  });
}

Nếu JSON.parse thất bại, Observer nhận error; nếu consumer unsubscribe trước khi timer chạy, teardown xoá timer. Đừng trông chờ try/catch bọc bên ngoài subscribe() để bắt một lỗi async xảy ra sau đó.

Lỗi producer khác lỗi của Observer

Nếu handler next do consumer truyền vào tự throw, RxJS 7 mặc định báo lỗi đó qua cơ chế unhandled error, không chuyển nó thành error của cùng stream. try/catch quanh subscriber.next(...) không phải cách bắt lỗi code của consumer. Đặt phép biến đổi có thể thất bại trong operator như map, hoặc xử lý lỗi ngay tại handler nếu đó là side effect của ứng dụng.

Kiểm tra subscriber closed trong producer đồng bộ

Một producer đồng bộ có thể đang chạy loop dài trong khi operator downstream như take(3) đã đóng subscription. Kiểm tra subscriber.closed giúp producer dừng tính toán ngay.

import { Observable, take } from 'rxjs';

function countTo(max: number): Observable<number> {
  return new Observable<number>((subscriber) => {
    for (let value = 1; value <= max; value += 1) {
      if (subscriber.closed) {
        return;
      }

      subscriber.next(value);
    }

    subscriber.complete();
  });
}

countTo(1_000_000)
  .pipe(take(3))
  .subscribe((value) => console.log(value));

Kết quả chỉ là 1, 2, 3; producer không tiếp tục quay gần một triệu vòng. Ví dụ này dùng để học contract. Nếu chỉ cần một dãy số, hãy dùng creation function range() có sẵn.

Unsubscribe không phải complete

complete() là notification từ producer: “nguồn đã kết thúc bình thường”. unsubscribe() là hành động từ consumer: “tôi không muốn nhận nữa”. Cả hai đều đóng subscription và chạy teardown, nhưng unsubscribe không gọi callback complete của Observer.

Phân biệt này ảnh hưởng trực tiếp đến thiết kế:

  • đặt cleanup bắt buộc trong teardown, không đặt riêng trong handler complete;
  • chỉ gọi complete() khi producer thật sự đã xong;
  • dùng unsubscribe() hoặc operator như take, takeUntil khi consumer quyết định dừng; hai operator này có thể gửi complete tới downstream trong lúc unsubscribe upstream, khác với gọi trực tiếp subscription.unsubscribe();
  • nếu cần side effect chạy trên cả complete, error và unsubscribe ở trong pipeline, cân nhắc finalize().

Cleanup khi setup thất bại giữa chừng

Giả sử adapter mở một listener, rồi bước setup tiếp theo throw. Nếu chỉ khai báo cleanup ở return cuối hàm, execution không bao giờ tới đó: RxJS có thể gửi lỗi cho Observer, nhưng listener đã mở vẫn bị giữ lại. Khi có nhiều bước cấp phát, đăng ký cleanup bằng subscriber.add(...) ngay sau từng bước thành công.

Ví dụ browser dưới đây cố ý throw sau khi đăng ký listener để bạn quan sát cleanup:

import { Observable } from 'rxjs';

const clicks$ = new Observable<MouseEvent>((subscriber) => {
  const onClick = (event: MouseEvent) => subscriber.next(event);

  document.addEventListener('click', onClick);
  subscriber.add(() => {
    document.removeEventListener('click', onClick);
    console.log('listener đã được gỡ');
  });

  // Mô phỏng một bước setup tiếp theo thất bại.
  throw new Error('Setup bước hai thất bại');
});

clicks$.subscribe({
  error: (error: unknown) => console.error('setup lỗi:', error),
});

Observer nhận lỗi, sau đó teardown đã đăng ký gỡ listener. Trong adapter thật, mỗi resource cấp phát thành công cần cleanup tương ứng trước khi đi tới bước tiếp theo. Nếu chính API nguồn cấp phát resource rồi throw mà không trả handle huỷ, bạn phải xử lý theo contract của API đó; RxJS không thể tự đoán cách đóng nó.

Đã add(cleanup) thì không cần đồng thời return cleanup cho cùng resource. RxJS có thể chạy hai đăng ký hàm riêng, nên đăng ký một cleanup hai lần không phải cách “cho chắc”. Với adapter đơn giản chỉ mở một watch hoặc timer, trả một hàm teardown vẫn là lựa chọn dễ đọc hơn.

Mỗi subscription có producer riêng

Với cách viết thông thường, code bên trong constructor tạo resource cho từng lần subscribe. watchPosition$ vì thế là một cold Observable: hai subscriber đồng thời sẽ tạo hai watchId và hai teardown riêng. Bản thân new Observable không đảm bảo cold: nếu subscribe function chỉ gắn listener vào một producer đã chạy sẵn bên ngoài, nguồn đó vẫn là hot.

Nếu nhiều consumer cần dùng chung một watch, đừng chuyển watchPosition ra ngoài constructor để “tiết kiệm”. Làm vậy khiến resource khởi động eagerly và ownership trở nên mơ hồ. Hãy giữ adapter đúng lifecycle, rồi multicast ở lớp sử dụng:

import { share } from 'rxjs';

const sharedPosition$ = watchPosition$().pipe(
  share(),
);

const mapSubscription = sharedPosition$.subscribe((position) => {
  console.log('render map:', position.coords);
});

const analyticsSubscription = sharedPosition$.subscribe((position) => {
  console.log('record position:', position.timestamp);
});

function stopConsumers(): void {
  mapSubscription.unsubscribe();
  analyticsSubscription.unsubscribe();
}

Hai subscription chồng thời gian dùng chung producer qua share(). Khi stopConsumers() làm subscriber cuối cùng rời đi, ref count trở về 0 và teardown upstream chạy theo cấu hình mặc định của share trong RxJS 7.

Chia sẻ là một quyết định riêng

Constructor trả lời “tạo và dọn producer thế nào”. share trả lời “bao nhiêu consumer dùng chung producer ấy”. Đừng trộn hai quyết định bằng biến global hoặc Subject ẩn bên trong adapter.

Thiết kế type và ranh giới adapter

Generic Observable<T> mô tả kiểu của giá trị đi qua next. Error channel của RxJS 7 không có generic riêng, nên Observer vẫn nên nhận unknown và thu hẹp kiểu trước khi dùng.

position$.subscribe({
  next: (position) => console.log(position.coords.latitude),
  error: (error: unknown) => {
    const message = error instanceof Error
      ? error.message
      : String(error);

    console.error(message);
  },
});

Không phải mọi trạng thái không thành công đều nên đi qua error. Nếu “không có kết quả” là dữ liệu nghiệp vụ có thể xử lý rồi stream tiếp tục, hãy mô hình hoá nó trong T, chẳng hạn { kind: 'empty' }. Chỉ dùng error khi execution hiện tại không thể tiếp tục theo contract của source.

Một adapter tốt thường có ba đặc điểm:

  1. public API nhận đúng options cần thiết và trả Observable<T>;
  2. subscribe function chỉ setup source, chuyển callback thành notification và khai báo teardown;
  3. biến đổi nghiệp vụ nằm ở operator phía ngoài.

Cách tách này cũng giữ ownership rõ ràng: adapter sở hữu resource, còn caller sở hữu Subscription hoặc lifecycle operator.

Những lỗi thường gặp

Tạo side effect bên ngoài subscribe function

// Không nên: watch bắt đầu ngay khi module được import.
const watchId = navigator.geolocation.watchPosition(onPosition);

const position$ = new Observable((subscriber) => {
  // Không còn lifecycle rõ ràng để sở hữu watchId.
});

Setup nên diễn ra trong subscribe function, trừ khi bạn chủ ý bọc một hot producer đã tồn tại và đã thiết kế ownership riêng.

Quên trả teardown

Timer hoặc listener vẫn hoạt động sau unsubscribe dù Observer không nhận thêm notification. Hãy rà từng lệnh setInterval, addEventListener, open, connect hoặc watch... và tìm hành động cleanup đối xứng.

Gọi error hoặc complete nhưng quên return

Terminal notification đóng Subscriber, nhưng phần còn lại của hàm vẫn chạy. Điều này có thể tạo resource sau khi execution đã đóng hoặc làm công việc vô ích.

if (!isValid(config)) {
  subscriber.error(new Error('Config không hợp lệ'));
  return;
}

Tự bọc lại API đã có adapter

Viết new Observable quanh Promise, DOM event hoặc một Observable khác thường làm code dài hơn và dễ quên nối teardown. Dùng from, fromEvent, defer hoặc operator trước; chỉ hạ xuống constructor khi contract nguồn thật sự khác. Riêng from(promise) chỉ ngừng chuyển kết quả khi unsubscribe, không tự huỷ công việc tạo Promise. Với HTTP cần abort, xem adapter chuyên biệt như fromFetch thay vì giả định constructor sẽ tự có cancellation.

Dùng subscribe function async

Đừng truyền async (subscriber) => ... vào constructor. Subscribe function phải trả TeardownLogic ngay, trong khi async luôn trả Promise. TypeScript sẽ báo không khớp kiểu; nếu bỏ qua kiểm tra kiểu, Promise đó cũng không trở thành teardown hợp lệ.

Với tác vụ một lần không cần cancellation, defer(() => loadData()) phù hợp hơn nếu loadData trả Promise. Nếu API có khả năng abort, adapter phải tạo cơ chế abort và đăng ký teardown đồng bộ, rồi mới khởi động công việc async. Lỗi Promise rejection cũng phải được nối vào error channel; RxJS không chờ Promise trả từ subscribe function để làm điều đó hộ bạn.

Trộn logic nghiệp vụ vào producer

Adapter Geolocation không cần biết bản đồ sẽ render gì. Giữ constructor nhỏ, rồi dùng map, filter, distinctUntilChanged hoặc operator phù hợp ở pipeline.

Dùng Subject thay constructor chỉ để phát giá trị

Subject phù hợp khi bạn cần một bridge imperative hoặc multicast chủ động. Nó không tự gắn setup với subscribe và teardown với unsubscribe. Nếu mỗi consumer phải mở một resource riêng, constructor Observable thể hiện ownership chính xác hơn.

Test notification và cleanup riêng

Đừng chỉ assert rằng next dừng sau unsubscribe: Subscriber đã chặn notification ngay cả khi bạn quên dọn producer. Một test lifecycle nên kiểm tra cả resource thật sự đóng và số lần cleanup chạy.

Ví dụ độc lập dưới đây dùng một callback source giả, không cần quyền Geolocation hay chờ timer. Bạn có thể chạy nó trong TypeScript playground đã có RxJS 7.8.2:

import { Observable, take } from 'rxjs';

type Start = (
  next: (value: number) => void,
  fail: (error: unknown) => void,
) => () => void;

function fromCallbacks(start: Start): Observable<number> {
  return new Observable<number>((subscriber) => {
    return start(
      (value) => subscriber.next(value),
      (error) => subscriber.error(error),
    );
  });
}

function assert(condition: boolean, message: string): void {
  if (!condition) {
    throw new Error(message);
  }
}

let active = 0;
let stopped = 0;
let emit: (value: number) => void = () => {};

const source$ = fromCallbacks((next) => {
  active += 1;
  emit = next;
  return () => {
    active -= 1;
    stopped += 1;
  };
});

assert(active === 0, 'không setup trước subscribe');
const values: number[] = [];
let completed = false;
const subscription = source$.subscribe({
  next: (value) => values.push(value),
  complete: () => { completed = true; },
});

assert(active === 1, 'setup khi subscribe');
emit(10);
subscription.unsubscribe();
subscription.unsubscribe();
emit(20); // Mô phỏng callback tới muộn dù source đã bị huỷ.

assert(values.join(',') === '10', 'chặn callback tới muộn');
assert(!completed, 'unsubscribe không phát complete');
assert(active === 0 && stopped === 1, 'cleanup đúng một lần');

let syncStopped = 0;
fromCallbacks((next) => {
  next(1); // take(1) đóng upstream trước khi start trả teardown.
  return () => { syncStopped += 1; };
}).pipe(take(1)).subscribe();
assert(syncStopped === 1, 'cleanup sau callback đồng bộ');

let errorStopped = 0;
let receivedError: unknown;
const expectedError = new Error('SDK lỗi');
fromCallbacks((_next, fail) => {
  fail(expectedError);
  return () => { errorStopped += 1; };
}).subscribe({
  error: (error: unknown) => { receivedError = error; },
});
assert(receivedError === expectedError, 'giữ nguyên lỗi nguồn');
assert(errorStopped === 1, 'cleanup trên error');

console.log('lifecycle assertions passed');

Đây là source giả cho một subscription tại một thời điểm, không phải SDK dùng chung để đưa vào production. Khi test adapter thật, thay start bằng mock của SDK hoặc browser API và assert hàm dispose/clearWatch nhận đúng handle. Thêm case complete tự nhiên, subscribe hai lần với hai handle riêng, và setup lỗi sau khi resource đầu tiên đã được cấp phát.

Với countdown, dùng fake timer nếu test runner hiện tại đã hỗ trợ; không cần chờ vài giây thật trong mỗi test. Với Geolocation, kiểm tra cả lỗi permission và callback tới muộn sau clearWatch, vì hai case đó dễ bị bỏ qua khi chỉ thử happy path bằng browser.

Checklist review custom Observable

Trước khi merge một adapter dùng new Observable, hãy trả lời được các câu sau:

  • Có creation function RxJS nào thay thế được đoạn code này không?
  • Mọi side effect có bắt đầu bên trong subscribe function không?
  • Mỗi resource được cấp phát có cleanup đối xứng không, kể cả khi setup lỗi giữa chừng?
  • Teardown có an toàn nếu lifecycle kết thúc bởi complete, error hoặc unsubscribe không?
  • Nhánh gọi error hoặc complete có dừng phần producer không còn cần thiết không?
  • Callback async có chuyển lỗi vào subscriber.error không?
  • Loop đồng bộ dài có quan sát subscriber.closed không?
  • Subscribe hai lần tạo hai execution có đúng ý định không?
  • Nếu cần chia sẻ, multicast có được đặt ở pipeline thay vì giấu trong adapter không?
  • Type T có phân biệt rõ dữ liệu nghiệp vụ với lỗi terminal không?
  • Subscribe function có trả teardown đồng bộ, thay vì một Promise từ async không?
  • Test có kiểm tra hàm cleanup và handle bị huỷ, thay vì chỉ kiểm tra Observer im lặng không?

Nếu một câu chưa rõ, vấn đề thường nằm ở ownership: ai mở resource, ai được quyền dừng, và đoạn code nào thật sự đóng nó.

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

Lấy hàm countdown và thử ba thay đổi:

  1. trong playground riêng, bỏ clearInterval khỏi teardown, unsubscribe giữa chừng, rồi thêm log ngay trong callback timer để thấy producer vẫn chạy dù Observer im lặng; sau khi quan sát, khôi phục cleanup và reset playground hoặc dừng process;
  2. subscribe hai lần và thêm ID cho mỗi interval để xác nhận mỗi subscription có producer riêng;
  3. thêm share(), subscribe hai consumer chồng thời gian, rồi xác nhận chỉ còn một interval upstream.

Sau đó trả lời bốn câu: setup chạy lúc nào, notification nào kết thúc source, unsubscribe có gọi complete không, và cleanup thuộc về ai. Nếu câu trả lời khớp với log, bạn đã nắm phần khó nhất của custom Observable.

Học tiếp

Nguồn tham khảo

On this page