【发布时间】:2018-05-24 17:49:41
【问题描述】:
我有一个消息队列处理器,可以将消息提供给服务...
q.on("message", (m) => {
service.create(m)
.then(() => m.ack())
.catch(() => n.nack())
})
该服务使用 RxJS Observable 并订阅 debounceTime() 这些请求。
class Service {
constructor() {
this.subject = new Subject()
this.subject.debounceTime(1000)
.subscribe(({ req, resolve, reject }) =>
someOtherService.doWork(req)
.then(() => resolve())
.catch(() => reject())
)
}
create(req) {
return new Promise((resolve, reject) =>
this.subject.next({
req,
resolve,
reject
})
)
}
}
问题是只有去抖动的请求才会被确认/取消。如何确保订阅也解决/拒绝其他请求? bufferTime() 让我参与其中,但它不会重置每次调用 next() 的超时持续时间。
【问题讨论】:
-
您可以将
buffer与从debounceTime构建的关闭通知一起使用。有关基本机制,请参阅this answer。这将在谴责期限内为您提供所有排放量,您可以随心所欲地处理它们。 -
由于该解决方案合并了两个流,您将如何在此处合并该方法?
-
您可以将合并排除在外。通知缓冲区很常见——通知程序使用
debounceTime而不是auditTime。我可以尽快给你写一个答案;之前在移动设备上。 -
我明白你现在的意思了......尝试了这种方法,似乎有效。