API 参考

rx4u 操作符、创建函数和工具函数的完整文档。

from

创建
from<T>(create: () => T): Sub<T>

从同步工厂函数创建流。

参数:

  • create: () => T - 返回要发出值的工厂函数

示例:

JavaScript
const lazy$ = from(() => Math.random());
const greeting$ = from(() => 'hello');

fromPromise

创建
fromPromise<T>(create: () => Promise<T>): Sub<T>

从 Promise 工厂函数创建流。

参数:

  • create: () => Promise<T> - 返回 Promise 的工厂函数

示例:

JavaScript
const request$ = fromPromise(() => fetch('/api/data'));
request$(console.log);

of

创建
of<T>(...values: T[]): Sub<T>

创建按顺序发出给定值的流。

参数:

  • values: T[] - 要发出的值

示例:

JavaScript
const stream$ = of(1, 2, 3);
stream$(console.log); // Emits: 1, 2, 3

interval

创建
interval(period: number): Sub<number>

创建按指定间隔发出连续数字的流。

参数:

  • period: number - 以毫秒为单位的时间间隔

示例:

JavaScript
const timer$ = interval(1000);
timer$(console.log); // Emits: 0, 1, 2, ... every second

fromEvent

创建
fromEvent<T extends Event = Event>(target: EventTargetLike, type: string, options?): Sub<T>

从类似 DOM 的事件目标创建事件流。

参数:

  • target: EventTargetLike - 提供 addEventListener 和 removeEventListener 的对象
  • type: string - 要监听的事件类型
  • options: AddEventListenerOptions | EventListenerOptions | boolean - 可选的事件监听器选项

示例:

JavaScript
const clicks$ = fromEvent<MouseEvent>(button, 'click');
clicks$(event => console.log(event.clientX));

createState

创建
createState<T>(initialValue: T): [Sub<T>, StateSetter<T>, () => T, () => void]

使用初始值创建响应式状态。

参数:

  • initialValue: T - 初始状态值

示例:

JavaScript
const [count$, setCount, getCount, destroy] = createState(0);
count$(console.log);
setCount(value => value + 1);

createLazyState

创建
createLazyState<T>(): [Sub<T>, StateSetter<T>, () => T | undefined, () => void]

创建首次更新前不会发出值的响应式状态。

示例:

JavaScript
const [name$, setName] = createLazyState<string>();
name$(console.log);
setName('Ada');

createSubject

创建
createSubject<T>(): [Sub<T>, (value: T) => void]

创建多播流和用于广播值的函数。

示例:

JavaScript
const [messages$, nextMessage] = createSubject<string>();
messages$(console.log);
nextMessage('hello');

pipeFrom

核心
pipeFrom<T>(source: Sub<T>, ...operators: Operator[]): Sub<unknown>

从左到右将操作符应用到源流。

参数:

  • source: Sub<T> - 源流
  • operators: Operator[] - 按顺序应用的操作符

示例:

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

firstValueFrom

核心
firstValueFrom<T>(source: Sub<T>): Promise<T>

解析第一个值;当流报错或空完成时拒绝。

参数:

  • source: Sub<T> - 源流

示例:

JavaScript
const first = await firstValueFrom(of(1, 2, 3));

map

转换
map<T, R>(mapFn: (value: T) => R): Operator<T, R>

使用投射函数转换每个发出的值。

参数:

  • mapFn: (value: T) => R - 转换每个值的函数

示例:

JavaScript
const doubled$ = pipeFrom(
  of(1, 2, 3),
  map(x => x * 2)
); // Emits: 2, 4, 6

switchMap

转换
switchMap<T, R>(project: (value: T, index: number) => Sub<R> | Promise<R>): Operator<T, R>

将每个值投射为新流,并切换到最新投射。

参数:

  • project: (value: T, index: number) => Sub<R> | Promise<R> - 将每个值投射为流或 Promise 的函数

示例:

JavaScript
const searchResults$ = pipeFrom(
  searchQuery$,
  debounceTime(300),
  switchMap(query => fromPromise(() => fetch(`/search?q=${query}`)))
);

mergeMap

转换
mergeMap<T, R>(project: (value: T, index: number) => Sub<R>, concurrent?: number): Operator<T, R>

将每个值投射为流,并合并所有并发流。

参数:

  • project: (value: T, index: number) => Sub<R> - 将每个值投射为流的函数
  • concurrent: number - 可选的最大活跃内部流数量

示例:

JavaScript
const allRequests$ = pipeFrom(
  urls$,
  mergeMap(url => fromPromise(() => fetch(url)))
);

concatMap

转换
concatMap<T, R>(project: (value: T, index: number) => Sub<R>): Operator<T, R>

将每个值投射为流,并按顺序连接。

参数:

  • project: (value: T, index: number) => Sub<R> - 将每个值投射为流的函数

示例:

JavaScript
const sequential$ = pipeFrom(
  of(1, 2, 3),
  concatMap(x => of(x, x * 2))
); // Emits: 1, 2, 2, 4, 3, 6

filter

过滤
filter<T>(predicate: (value: T) => boolean): Operator<T, T>

只发出满足给定谓词函数的值。

参数:

  • predicate: (value: T) => boolean - 测试每个值的函数

示例:

JavaScript
const evens$ = pipeFrom(
  of(1, 2, 3, 4, 5),
  filter(x => x % 2 === 0)
); // Emits: 2, 4

take

过滤
take<T>(count: number): Operator<T, T>

只发出源流的前 N 个值。

参数:

  • count: number - 要获取的值数量

示例:

JavaScript
const first3$ = pipeFrom(
  interval(1000),
  take(3)
); // Emits: 0, 1, 2, then completes

takeWhile

过滤
takeWhile<T>(predicate: (value: T, index: number) => boolean, inclusive?: boolean): Operator<T, T>

谓词通过时持续发出值,随后完成。

参数:

  • predicate: (value: T, index: number) => boolean - 测试每个值的函数
  • inclusive: boolean - 是否发出第一个未通过谓词的值

示例:

JavaScript
const values$ = pipeFrom(
  of(1, 2, 3, 4),
  takeWhile(value => value < 3, true)
); // Emits: 1, 2, 3

takeUntil

过滤
takeUntil<T>(notifier: Sub<unknown>): Operator<T, T>

在通知流发出值前持续发出源流的值。

参数:

  • notifier: Sub<unknown> - 其第一个值会停止源流的流

示例:

JavaScript
const values$ = pipeFrom(
  interval(1000),
  takeUntil(stop)
);

distinctUntilChanged

过滤
distinctUntilChanged<T, K = T>(compareFn?: (a: K, b: K) => boolean, keySelector?: (value: T) => K): Operator<T, T>

仅在当前值不同于前一个值时发出。

参数:

  • compareFn: (a: K, b: K) => boolean - 可选的比较函数
  • keySelector: (value: T) => K - 选择比较键的可选函数

示例:

JavaScript
const unique$ = pipeFrom(
  of(1, 1, 2, 2, 3),
  distinctUntilChanged()
); // Emits: 1, 2, 3

merge

组合
merge<T>(...sources: Sub<T>[]): Sub<T>

将多个流合并为单个流。

参数:

  • sources: Sub<T>[] - 要合并的流

示例:

JavaScript
const combined$ = merge(
  of(1, 2, 3),
  of(4, 5, 6)
); // Emits: 1, 2, 3, 4, 5, 6

combineLatest

组合
combineLatest<T extends readonly unknown[]>(...sources: {[K in keyof T]: Sub<T[K]>}): Sub<T>

任意源流发出值时组合所有流的最新值。

参数:

  • sources: Sub<T>[] - 要组合的流

示例:

JavaScript
const latest$ = combineLatest(
  of(1, 2),
  of('a', 'b')
); // Emits: [2, 'a'], [2, 'b']

zip

组合
zip<T extends readonly unknown[]>(...sources: {[K in keyof T]: Sub<T[K]>}): Sub<T>

按位置组合多个流中的值。

参数:

  • sources: Sub<T>[] - 要按位置配对的流

示例:

JavaScript
const paired$ = zip(
  of(1, 2),
  of('a', 'b')
); // Emits: [1, 'a'], [2, 'b']

withLatestFrom

组合
withLatestFrom<T>(...others: Sub<unknown>[]): Operator<T, [T, ...unknown[]]>

将源流的值与其他流的最新值组合。

参数:

  • others: Sub<unknown>[] - 提供最新伴随值的流

示例:

JavaScript
const values$ = pipeFrom(
  clicks$,
  withLatestFrom(state)
);

tap

工具
tap<T>(tapFn: (value: T) => void): Operator<T, T>

为源流通知执行副作用。

参数:

  • tapFn: (value: T) => void - 为每个值调用的函数

示例:

JavaScript
const logged$ = pipeFrom(
  of(1, 2, 3),
  tap(x => console.log('Value:', x))
);

delay

工具
delay<T = void>(ms: number): Operator<T, T>

将每个通知延迟指定的时间。

参数:

  • ms: number - 以毫秒为单位的延迟

示例:

JavaScript
const delayed$ = pipeFrom(of(1, 2, 3), delay(500));

debounceTime

工具
debounceTime<T>(dueTime: number): Operator<T, T>

在指定静默期后才发出值。

参数:

  • dueTime: number - 等待时间(毫秒)

示例:

JavaScript
const debounced$ = pipeFrom(
  searchInput$,
  debounceTime(300)
); // Wait 300ms after last input

throttleTime

工具
throttleTime<T>(duration: number, options?: {leading?: boolean; trailing?: boolean}): Operator<T, T>

在每个时间段内最多发出一个值。

参数:

  • duration: number - 时间周期(毫秒)
  • options: {leading?: boolean; trailing?: boolean} - 可选的前沿和后沿发出行为

示例:

JavaScript
const throttled$ = pipeFrom(
  clicks$,
  throttleTime(1000)
); // Max 1 click per second

timeout

工具
timeout<T>(ms: number, errorFactory?: () => Error): Operator<T, T>

当源流在给定时间内未发出值时失败。

参数:

  • ms: number - 超时长度(毫秒)
  • errorFactory: () => Error - 可选的自定义错误工厂函数

示例:

JavaScript
const guarded$ = pipeFrom(request$, timeout(5000));

startWith

工具
startWith<T>(...values: T[]): Operator<T, T>

在源流值之前发出指定值。

参数:

  • values: T[] - 要前置发出的值

示例:

JavaScript
const values$ = pipeFrom(of(2, 3), startWith(1));

endWith

工具
endWith<T>(...values: T[]): Operator<T, T>

在源流完成后发出指定值。

参数:

  • values: T[] - 要追加发出的值

示例:

JavaScript
const values$ = pipeFrom(of(1, 2), endWith(3));

share

工具
share<T>(): Operator<T, T>

在订阅者之间共享一次源流订阅。

示例:

JavaScript
const shared$ = pipeFrom(
  expensiveOperation$,
  share()
); // Multiple subscribers share the same execution

catchError

错误处理
catchError<T, R = T>(selector: (error: unknown, caught: Sub<T>) => Sub<R>): Operator<T, T | R>

捕获错误并使用回退流继续执行。

参数:

  • selector: (error: unknown, caught: Sub<T>) => Sub<R> - 处理错误的函数

示例:

JavaScript
const resilient$ = pipeFrom(
  riskyOperation$,
  catchError(err => of('fallback'))
);

retry

错误处理
retry<T>(count?: number): Operator<T, T>

源流发生错误时重试。

参数:

  • count: number - 重试次数

示例:

JavaScript
const retried$ = pipeFrom(
  unreliableRequest$,
  retry(3)
); // Retry up to 3 times

retryWhen

错误处理
retryWhen<T>(notifier: (errors: Sub<unknown>) => Sub<unknown>): Operator<T, T>

根据自定义重试逻辑重试。

参数:

  • notifier: (errors: Sub<unknown>) => Sub<unknown> - 控制重试时机的函数

示例:

JavaScript
const smartRetry$ = pipeFrom(
  source$,
  retryWhen(errors => pipeFrom(
    errors,
    debounceTime(1000),
    take(3)
  ))
);

scan

数学
scan<T, R = T>(accumulator: (acc: R, value: T, index: number) => R, seed?: R): Operator<T, R>

应用累加器函数并发出每个中间结果。

参数:

  • accumulator: (acc: R, value: T, index: number) => R - 累加器函数
  • seed: R - 可选的初始累加器值

示例:

JavaScript
const running$ = pipeFrom(
  of(1, 2, 3, 4),
  scan((acc, x) => acc + x, 0)
); // Emits: 1, 3, 6, 10

reduce

数学
reduce<T, R = T>(accumulator: (acc: R, value: T, index: number) => R, seed?: R): Operator<T, R>

应用累加器函数,并在源流完成时发出最终结果。

参数:

  • accumulator: (acc: R, value: T, index: number) => R - 累加器函数
  • seed: R - 可选的初始累加器值

示例:

JavaScript
const sum$ = pipeFrom(
  of(1, 2, 3, 4),
  reduce((acc, x) => acc + x, 0)
); // Emits: 10