Observable S1
提案速览
该提案向 ECMAScript 标准库引入了 Observable 类型,用于建模基于推送的数据源,如 DOM 事件、定时器和套接字。它具有可组合性和惰性,允许组合 observables 并且仅在订阅后发射。该提案定义了核心的 Observable 接口、subscribe 方法和静态方法 of/from,以及 Observer 和 SubscriptionObserver 类型,并注明 filter 和 map 等组合子尚未包含在内。
Note
以下 README 来自上游仓库,其中的阶段或状态标注可能滞后;当前信息以提案概览为准。
ECMAScript Observable
本提案向 ECMAScript 标准库引入 Observable 类型。
Observable 类型可用于建模基于推送的数据源,如 DOM 事件、定时器间隔和套接字。此外,Observables 具有以下特性:
- 可组合:Observables 可以通过高阶组合子进行组合。
- 惰性:Observables 在 观察者 订阅之前不会开始发出数据。
示例:观察键盘事件
使用 Observable 构造函数,我们可以创建一个函数,该函数为任意 DOM 元素和事件类型返回一个可观察的事件流。
function listen(element, eventName) {
return new Observable(observer => {
// 创建一个事件处理器,将数据发送到接收器
let handler = event => observer.next(event);
// 附加事件处理器
element.addEventListener(eventName, handler, true);
// 返回一个清理函数,用于取消事件流
return () => {
// 从元素上移除事件处理器
element.removeEventListener(eventName, handler, true);
};
});
}
然后我们可以使用标准组合子来过滤和映射流中的事件,就像处理数组一样。
// 返回一个特殊按键命令的 Observable
function commandKeys(element) {
let keyCommands = { "38": "up", "40": "down" };
return listen(element, "keydown")
.filter(event => event.keyCode in keyCommands)
.map(event => keyCommands[event.keyCode])
}
注意:本提案不包括“filter”和“map”方法。它们可能会在本规范的未来版本中添加。
当我们需要消费事件流时,我们使用 observer 进行订阅。
let subscription = commandKeys(inputElement).subscribe({
next(val) { console.log("Received key command: " + val) },
error(err) { console.log("Received an error: " + err) },
complete() { console.log("Stream complete") },
});
subscribe 方法返回的对象将允许我们随时取消订阅。
取消时,将执行 Observable 的清理函数。
// 调用此函数后,将不再发送事件
subscription.unsubscribe();
动机
Observable 类型代表了处理异步数据流的基本协议之一。它在建模源自环境并推送到应用程序的数据流(例如用户界面事件)方面特别有效。通过将 Observable 作为 ECMAScript 标准库的组件,我们允许平台和应用程序共享一个通用的基于推送的流协议。
实现
运行测试
要运行单元测试,请将 es-observable-tests 包安装到您的项目中。
npm install es-observable-tests
然后使用您要测试的构造函数调用导出的 runTests 函数。
require("es-observable-tests").runTests(MyObservable);
API
Observable
Observable 表示一个可被观察的值序列。
interface Observable {
constructor(subscriber : SubscriberFunction);
// 使用一个 observer 订阅序列
subscribe(observer : Observer) : Subscription;
// 使用回调函数订阅序列
subscribe(onNext : Function,
onError? : Function,
onComplete? : Function) : Subscription;
// 返回自身
[Symbol.observable]() : Observable;
// 将项目转换为 Observable
static of(...items) : Observable;
// 将 observable 或 iterable 转换为 Observable
static from(observable) : Observable;
}
interface Subscription {
// 取消订阅
unsubscribe() : void;
// 指示订阅是否已关闭的布尔值
get closed() : Boolean;
}
function SubscriberFunction(observer: SubscriptionObserver) : (void => void)|Subscription;
Observable.of
Observable.of 创建一个 Observable,其值为提供的参数。当调用 subscribe 时,这些值会同步传递。
Observable.of("red", "green", "blue").subscribe({
next(color) {
console.log(color);
}
});
/*
> "red"
> "green"
> "blue"
*/
Observable.from
Observable.from 将其参数转换为 Observable。
- 如果参数具有
Symbol.observable 方法,则返回调用该方法的结果。如果结果对象不是 Observable 的实例,则将其包装在一个 Observable 中,该 Observable 将委托订阅。
- 否则,假设参数是可迭代的,并在调用
subscribe 时同步传递迭代值。
从支持 Symbol.observable 的对象转换为 Observable:
Observable.from({
[Symbol.observable]() {
return new Observable(observer => {
setTimeout(() => {
observer.next("hello");
observer.next("world");
observer.complete();
}, 2000);
});
}
}).subscribe({
next(value) {
console.log(value);
}
});
/*
> "hello"
> "world"
*/
let observable = new Observable(observer => {});
Observable.from(observable) === observable; // true
从可迭代对象转换为 Observable:
Observable.from(["mercury", "venus", "earth"]).subscribe({
next(value) {
console.log(value);
}
});
/*
> "mercury"
> "venus"
> "earth"
*/
Observer
Observer 用于接收来自 Observable 的数据,并作为参数提供给 subscribe。
所有方法都是可选的。
interface Observer {
// 当调用 `subscribe` 时接收订阅对象
start(subscription : Subscription);
// 接收序列中的下一个值
next(value);
// 接收序列错误
error(errorValue);
// 接收完成通知
complete();
}
SubscriptionObserver
SubscriptionObserver 是一个规范化的 Observer,它包装了提供给 subscribe 的 observer 对象。
interface SubscriptionObserver {
// 发送序列中的下一个值
next(value);
// 发送序列错误
error(errorValue);
// 发送完成通知
complete();
// 指示订阅是否已关闭的布尔值
get closed() : Boolean;
}