【问题标题】:RxJs avoid external state but still access previous valuesRxJs 避免外部状态,但仍然访问以前的值
【发布时间】:2016-04-14 01:57:41
【问题描述】:

我正在使用 RxJs 来收听 amqp 队列(不是很相关)。

我有一个函数createConnection,它返回一个Observable,它发出新的connection 对象。建立连接后,我想每 1000 毫秒通过它发送消息,并且在 10 条消息之后我想关闭连接。

我试图避免外部状态,但如果我不将连接存储在外部变量中,我该如何关闭它?看我从连接开始,然后flatMap 并推送消息,所以在几个链之后我不再拥有连接对象。

这不是我的流程,但想象一下这样的事情:

createConnection()
  .flatMap(connection => connection.createChannel())
  .flatMap(channel => channel.send(message))
  .do(console.log)
  .subscribe(connection => connection.close()) <--- obviously connection isn't here

现在我明白这样做很愚蠢,但现在我该如何访问连接?我当然可以从var connection = createConnection()开始

然后以某种方式加入。但是我该怎么做呢?我什至不知道如何正确地问这个问题。底线,我所拥有的是一个可观察的,它发出一个连接,在打开连接后我想要一个每 1000 毫秒发出消息的可观察(带有take(10)),然后关闭连接

【问题讨论】:

    标签: javascript node.js reactive-programming rxjs


    【解决方案1】:

    您的问题的直接答案是“您可以完成每个步骤”。例如,您可以替换此行

    .flatMap(connection => connection.createChannel())
    

    用这个:

    .flatMap(connection => ({ connection: connection, channel: connection.createChannel() }))
    

    并一直保持对连接的访问​​。

    但是还有另一种方法可以做你想做的事。假设您的 createConnection 和 createChannel 函数如下所示:

    function createConnection() {
      return Rx.Observable.create(observer => {
        console.log('creating connection');
        const connection = {
          createChannel: () => createChannel(),
          close: () => console.log('disposing connection')
        };
    
        observer.onNext(connection);
    
        return Rx.Disposable.create(() => connection.close());
      });
    }
    
    function createChannel() {
      return Rx.Observable.create(observer => {
        const channel = {
          send: x => console.log('sending message: ' + x)
        };
    
        observer.onNext(channel);
    
        // assuming no cleanup here, don't need to return disposable
      });
    }
    

    createConnection(和createChannel,但我们将关注前者)返回一个冷的可观察对象;每个订阅者都将获得自己的包含单个连接的连接流,当订阅到期时,将自动调用 dispose 逻辑。

    这允许你做这样的事情:

    const subscription = createConnection()
      .flatMap(connection => connection.createChannel())
      .flatMap(channel => Rx.Observable.interval(1000).map(i => ({ channel: channel, data: i })))
      .take(10)
      .subscribe(x => x.channel.send(x.data))
    ;
    

    实际上,您不必为了进行清理而处置订阅;满足take(10) 后,整个链条将完成并触发清理。您需要在订阅上显式调用 dispose 的唯一原因是,如果您想在 10 1000 毫秒间隔结束之前将其拆除。

    请注意,此解决方案还包含直接回答您的问题的一个实例:我们将频道推到最后一行,以便我们可以在传递给 subscribe 调用的 onNext lambda 中使用它(通常会出现此类代码) .

    这是整个工作:https://jsbin.com/korihe/3/edit?js,console,output

    【讨论】:

    • 很好的答案,谢谢马特。您能否详细说明如何处理连接?我猜是return Rx.Disposable.create(() =&gt; console.log('disposing connection'));,但我不太明白清理逻辑是如何触发的。
    • 对不起,我完全跳过了那部分!我已经更新了我的答案以使事情更清楚。
    • 其实我可能错过了你之前评论的重点;我跳过了接线(我现在已经更换了),但仔细阅读,听起来你在问如何触发处置本身。如果有,请告诉我,我会详细说明。
    • 是的,我就是这个意思,呵呵。我不知道一次性用品或实际情况如何。因为在observer.onNext() 之后你会返回它。另外,如果我想共享一个连接,我可能需要一个 hot observable 和一个 AsyncSubject?因为如果是这样,我认为我无法完成这个技巧。
    • 了解 rxjs 中各种 Disposable 类型的文档是个好主意,因为它们提供了很好的清理模式。通常,当订阅完成或显式处置时,将触发与 observable 关联的任何处置逻辑。 Observable.create 允许通过返回 Disposable 的实例来指定此类清理逻辑。有诸如“share”之类的调用允许您共享对底层冷 observable 的单个订阅,因此您仍然可以使用这种方法。
    【解决方案2】:

    这段代码给了我一个错误,因为 flatmap 等待 observable({ connection: connection, channel: connection.createChannel() })这是一个对象。

    .flatMap(connection =&gt; ({ connection: connection, channel: connection.createChannel() }))

    您可以使用 combineLatest 运算符

    .flatMap(connection => Observable.combineLatest( Observable.of(connection), connection.createChannel(), (connection, channel) => { ... code .... });

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2021-07-28
      • 2014-07-01
      • 2014-11-21
      • 1970-01-01
      • 2023-03-24
      相关资源
      最近更新 更多