【发布时间】: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