【问题标题】:in rx.js make source.subscribe await it's observer using async/await在 rx.js 中使用 async/await 让 source.subscribe 等待它的观察者
【发布时间】:2017-07-07 21:17:38
【问题描述】:

我有一个这样的观察者。

var source = rx.Observable.fromEvent(eventAppeared.emitter, 'event')
      .filter(mAndF.isValidStreamType)
      .map(mAndF.transformEvent)
      .share();

然后我与一些订阅者分享它。这些订阅者都接受事件并对它们执行一些异步操作。 所以我的订阅者就像

 source.subscribe(async function(x) {
  const func = handler[x.eventName];
  if (func) {
    await eventWorkflow(x, handler.handlerName, func.bind(handler));
  }
});

里面有一些额外的东西,但我认为意图很明确。 我需要每个处理这个特定事件的“处理程序”来处理它并阻止它直到它回来。然后处理下一个事件。

我发现上面的代码只是调用事件而不等待它,而我的处理程序正在踩到自己。

我已经阅读了相当多的帖子,但我真的不知道该怎么做。大多数人都在谈论让观察者等待。但这不是我需要的不是吗?看来我需要的是让观察者等待。我找不到任何东西,这通常意味着它要么超级容易,要么超级荒谬。我希望是前者。

如果您需要任何进一步的说明,请告诉我。

---更新---

我意识到我需要的是一个先进先出队列或缓冲区(先进先出),有时也称为背压。我需要按顺序处理所有消息,并且仅在前面的消息完成处理时才需要。

---结束更新---

起初我以为是因为我使用的是 rx 2.5.3,但我刚刚升级到 4.1.0,它仍然不同步。

【问题讨论】:

    标签: javascript async-await rxjs


    【解决方案1】:

    没有办法告诉可观察源在subscribe 中暂停事件,它只是让我们“观察”传入事件。异步事物应该通过 Rx 操作符来管理。

    例如,要让您的异步处理程序按顺序处理事件,您可以尝试使用concatMap 运算符:

    source
        .concatMap(x => {
            const func = handler[x.eventName];
            return func ?
                eventWorkflow(x, handler.handlerName, func.bind(handler)) :
                Rx.Observable.empty();
        })
        .subscribe();
    

    请注意,在上面的示例中,不需要await,因为concatMap 知道如何处理eventWorkflow 返回的promise:concatMap 将其转换为可观察对象并等到可观察对象完成后再继续下一个活动。

    【讨论】:

    • 您好,谢谢您,非常详细和有帮助。你会说 concatmap 是正确的模式吗?这是我的应用程序的基础。如果有更标准的方法,我愿意进行更深入的重构。
    • 恐怕这对我的目的不起作用。它立即处理所有事件,然后让异步工作并返回值。我需要的是按顺序一次处理每个事件
    • 嗨,雷夫。您能否详细说明到底出了什么问题?您是否需要某种背压,即在特定处理程序处理事件时“暂停”源 observable?还是您只需要按顺序处理源事件?如果是这样,那么concatMap 就允许这样做(这里是一个人为的例子:jsbin.com/mugeki/edit?js,console)。您还提到同一源流有多个订阅者,不同订阅者的处理程序是否也同步,这意味着即使对于不同的订阅者,也没有两个处理程序并行工作?
    • 嗨,你说得对,concatMap 会确保它们按顺序处理,但是我只需要在前面的消息完成后处理。我更新了我的问题。我也回答了我的问题。虽然我感谢您的帮助,但我必须将我的答案标记为解决了我的问题。但我会支持你的:)
    【解决方案2】:

    所以最终我发现我需要的更准确地描述为先进先出队列或缓冲区。我需要消息等到上一条消息处理完毕。

    我也很确定 rxjs 不提供此功能(有时称为背压)。所以我所做的只是导入一个先进先出队列并将其连接到每个订阅者。

    我正在使用并发队列,到目前为止似乎运行良好。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2021-12-11
      • 2016-03-08
      • 1970-01-01
      • 2019-08-20
      • 1970-01-01
      • 2019-04-17
      相关资源
      最近更新 更多