Skip to content
⚠️ This article was written in 2019. Some content may be outdated.

RxJS Basics: Introduction to Observable

RxJS is a reactive programming library for JavaScript built on the Observable pattern, providing a powerful toolkit for working with asynchronous data streams. The Angular ecosystem depends on RxJS heavily, but even in React projects RxJS can handle complex async scenarios elegantly. Starting from the basic concept of an Observable, this article systematically introduces RxJS's core usage.

What Is an Observable ​

An Observable represents a subscribable data stream. Unlike a Promise, which resolves only once, an Observable can emit multiple values:

js
import { Observable } from 'rxjs';

// 创建一个 Observable
const observable = new Observable(subscriber => {
  subscriber.next(1);
  subscriber.next(2);
  subscriber.next(3);

  setTimeout(() => {
    subscriber.next(4);
    subscriber.complete();
  }, 1000);
});

// 订阅
console.log('开始订阅');
observable.subscribe({
  next: value => console.log('收到值:', value),
  error: err => console.error('错误:', err),
  complete: () => console.log('完成')
});

// 输出:
// 开始订阅
// 收到值: 1
// 收到值: 2
// 收到值: 3
// (1秒后)
// 收到值: 4
// 完成

Observable Creation Factory Functions ​

RxJS ships a rich set of creation functions, so you rarely need to write new Observable by hand:

js
import {
  of,
  from,
  interval,
  timer,
  fromEvent,
  throwError,
  EMPTY,
  combineLatest,
  merge,
  zip
} from 'rxjs';

// of: 发出一系列值后完成
of(1, 2, 3).subscribe(console.log);
// 1, 2, 3

// from: 将数组、Promise、迭代器转换为 Observable
from([10, 20, 30]).subscribe(console.log);
// 10, 20, 30

from(fetch('/api/data').then(r => r.json())).subscribe(console.log);

// interval: 每隔一段时间发出递增的数字
interval(1000).subscribe(x => console.log(x));
// 0, 1, 2, 3 ... (每秒)

// timer: 延迟后开始发出
timer(3000, 1000).subscribe(x => console.log(x));
// 3秒后: 0, 1, 2 ... (每秒)

// fromEvent: 从 DOM 事件创建
fromEvent(document, 'click').subscribe(event => {
  console.log('点击位置:', event.clientX, event.clientY);
});

// EMPTY: 立即完成的空 Observable
EMPTY.subscribe({
  next: () => console.log('不会执行'),
  complete: () => console.log('立即完成')
});

// throwError: 立即发出错误
throwError('出错了').subscribe({
  error: err => console.error(err)
});

Operators ​

Operators are the heart of RxJS—they transform, filter, and combine Observables:

js
import { of, interval } from 'rxjs';
import {
  map,
  filter,
  take,
  debounceTime,
  distinctUntilChanged,
  switchMap,
  mergeMap,
  catchError,
  tap,
  reduce,
  scan
} from 'rxjs/operators';

// 链式调用
of(1, 2, 3, 4, 5)
  .pipe(
    filter(x => x % 2 === 0),    // 过滤偶数
    map(x => x * 10),             // 乘以 10
    tap(x => console.log('中间值:', x))  // 副作用
  )
  .subscribe(result => console.log('最终结果:', result));
// 中间值: 20
// 最终结果: 20
// 中间值: 40
// 最终结果: 40

Transformation Operators ​

js
import { from } from 'rxjs';
import { map, mapTo, pluck, scan, reduce } from 'rxjs/operators';

// map: 转换每个值
of(1, 2, 3).pipe(
  map(x => x * x)
).subscribe(console.log);
// 1, 4, 9

// pluck: 提取嵌套属性
from([
  { user: { name: '张三' } },
  { user: { name: '李四' } }
]).pipe(
  pluck('user', 'name')
).subscribe(console.log);
// '张三', '李四'

// scan: 累加器(类似 reduce,但每个中间值都会发出)
of(1, 2, 3, 4).pipe(
  scan((acc, val) => acc + val, 0)
).subscribe(console.log);
// 1, 3, 6, 10

Filtering Operators ​

js
import { interval, fromEvent, Subject } from 'rxjs';
import {
  filter,
  take,
  takeUntil,
  debounceTime,
  throttleTime,
  distinctUntilChanged,
  first,
  last
} from 'rxjs/operators';

// take: 只取前 N 个值
interval(1000).pipe(
  take(3)
).subscribe(console.log);
// 0, 1, 2 然后完成

// takeUntil: 直到另一个 Observable 发出值
const stop$ = new Subject();

interval(1000).pipe(
  takeUntil(stop$)
).subscribe(console.log);

setTimeout(() => stop$.next(), 5000);
// 5秒后停止: 0, 1, 2, 3, 4

// debounceTime: 防抖
fromEvent(document.querySelector('#search'), 'input').pipe(
  debounceTime(300),
  pluck('target', 'value'),
  distinctUntilChanged()
).subscribe(query => {
  console.log('搜索:', query);
});

// throttleTime: 节流
fromEvent(document, 'scroll').pipe(
  throttleTime(100)
).subscribe(() => {
  console.log('页面滚动');
});

Higher-Order Operators ​

Higher-order operators deal with "Observables of Observables":

js
import { interval, of, timer } from 'rxjs';
import {
  switchMap,
  mergeMap,
  concatMap,
  exhaustMap
} from 'rxjs/operators';

// switchMap: 切换到新的 Observable,取消之前的
// 典型场景:搜索框自动补全
fromEvent(searchInput, 'input').pipe(
  debounceTime(300),
  pluck('target', 'value'),
  switchMap(query =>
    from(fetch(`/api/search?q=${query}`).then(r => r.json()))
  )
).subscribe(results => {
  renderResults(results);
});

// mergeMap: 并行执行所有内部 Observable
// 典型场景:批量请求
of(1, 2, 3).pipe(
  mergeMap(id =>
    from(fetch(`/api/user/${id}`).then(r => r.json()))
  )
).subscribe(user => {
  console.log('用户:', user);
});
// 三个请求并行执行

// concatMap: 顺序执行,一个完成后再执行下一个
// 典型场景:有序操作队列
of(1, 2, 3).pipe(
  concatMap(id =>
    from(fetch(`/api/user/${id}`).then(r => r.json()))
  )
).subscribe(user => {
  console.log('用户:', user);
});
// 按顺序 1 → 2 → 3 依次执行

// exhaustMap: 忽略新值直到当前 Observable 完成
// 典型场景:防止重复提交
fromEvent(submitBtn, 'click').pipe(
  exhaustMap(() =>
    from(fetch('/api/submit', { method: 'POST' }))
  )
).subscribe(response => {
  console.log('提交成功');
});

Subject ​

A Subject is both an Observable and an Observer, and is used to share data among multiple subscribers:

js
import { Subject, BehaviorSubject, ReplaySubject } from 'rxjs';

// Subject: 普通主题
const subject = new Subject();

subject.subscribe(x => console.log('订阅者A:', x));
subject.next(1);

subject.subscribe(x => console.log('订阅者B:', x));
subject.next(2);
// 订阅者A: 2(能收到)
// 订阅者B: 2(能收到)

// BehaviorSubject: 保存最新值,新订阅者立即收到
const behavior = new BehaviorSubject('初始值');

behavior.subscribe(x => console.log('订阅者A:', x));
// 立即输出: 订阅者A: 初始值

behavior.next('新值');

behavior.subscribe(x => console.log('订阅者B:', x));
// 立即输出: 订阅者B: 新值(收到最新值)

// ReplaySubject: 重放最近 N 个值
const replay = new ReplaySubject(3);

replay.next(1);
replay.next(2);
replay.next(3);
replay.next(4);

replay.subscribe(x => console.log('新订阅者:', x));
// 重放: 2, 3, 4(最近3个值)

Error Handling ​

js
import { of, throwError, timer } from 'rxjs';
import { catchError, retry, retryWhen, delay, map } from 'rxjs/operators';

// catchError: 捕获错误并返回 fallback Observable
from(fetch('/api/data')).pipe(
  catchError(err => {
    console.error('请求失败:', err);
    return of({ data: '默认数据' });  // 返回默认值
  })
).subscribe(data => console.log(data));

// retry: 自动重试
from(fetch('/api/data')).pipe(
  retry(3)  // 失败后重试3次
).subscribe(data => console.log(data));

// retryWhen: 自定义重试策略
from(fetch('/api/data')).pipe(
  retryWhen(errors => errors.pipe(
    delay(1000),
    take(3),
    concat(throwError('重试次数已达上限'))
  ))
).subscribe(
  data => console.log(data),
  err => console.error(err)
);

Unsubscribing ​

js
import { Subscription, interval } from 'rxjs';

// 手动管理订阅
const sub: Subscription = interval(1000).subscribe(console.log);

// 取消订阅
sub.unsubscribe();

// 合并多个订阅
const subs = new Subscription();
subs.add(interval(1000).subscribe(x => console.log('A:', x)));
subs.add(interval(2000).subscribe(x => console.log('B:', x)));

// 一次性取消所有
subs.unsubscribe();

Summary ​

  • An Observable represents a subscribable data stream that can emit multiple values
  • of, from, interval, and fromEvent are the most commonly used creation functions
  • Operators are chained via pipe and fall into categories like transformation, filtering, and combination
  • switchMap fits search scenarios, mergeMap fits parallel requests, and concatMap fits ordered execution
  • A Subject is the core of multicasting; BehaviorSubject keeps the latest value and ReplaySubject replays past values
  • catchError and retry handle errors
  • Unsubscribe promptly to avoid memory leaks

MIT Licensed