For AI agents: the complete documentation index is available at /tc39-atlas/llms.txt, the full documentation bundle is available at /tc39-atlas/llms-full.txt, and this page is available as Markdown at /tc39-atlas/proposals/proposal-observable.md.
  • 简体中文
  • 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;
    }