【发布时间】:2018-11-23 14:56:30
【问题描述】:
我想为 observable 创建一个速率限制器,有条件地改变它限制值的方式。该用例适用于不断接收要下载的新 URL 的下载器。我希望能够对传入的 url 进行排队并在两种排队方法之间切换。这两种方法是通过限制速率(例如每 2 秒不超过 10 个请求)和通过并发请求数(例如一次可以触发不超过 10 个请求)。
可以像下面这样实现一个简单的速率限制器(借用here):
const rateLimit = (limit, rate, scheduler = asyncScheduler) => {
let tokens = limit
const tokenChanged = new BehaviorSubject(tokens)
const consumeToken = () => tokenChanged.next(--tokens)
const renewToken = () => tokenChanged.next(++tokens)
const availableTokens = tokenChanged.pipe(filter(() => tokens > 0))
return source =>
source.pipe(
mergeMap(val =>
availableTokens.pipe(
take(1),
map(() => {
consumeToken()
timer(rate, scheduler).subscribe(renewToken)
return val
})
)
)
)
}
const o = urlSource.pipe(
rateLimit(10, 2000),
mergeMap(downloadUrl)
)
o.toPromise()
这是一个简单的并发限制器:
const o = urlSource.pipe(
mergeMap(downloadUrl, maxConcurrent)
)
o.toPromise()
最后,我可以创建这个组合切换器来选择使用哪种类型的限制器:
const toggleableLimiter = (func, limit, rate, concurrent, toggleObservable) => {
let useRateLimiter = true
toggleObservable.subscribe(() => (useRateLimiter = !useRateLimiter))
const rateLimiter = rateLimit(limit, rate)
return source => {
const operators = useRateLimiter
? [rateLimiter, mergeMap(func)]
: [mergeMap(func, concurrent)]
return source.pipe(...operators)
}
}
const e = new EventEmitter()
const toggler = fromEvent(e, 'toggle')
const o = urlSource.pipe(
toggleableLimiter(downloadUrl, 2, 1000, 2, toggler)
)
o.toPromise()
// using rate limiter
e.emit('toggle')
// incoming values now use concurrent limiter
这一切都可以很好地解决我的问题。我可以使用事件发射器在两种方法之间切换。然而,问题是,在事件发出之前,任何已传递给toggleableLimiter 的东西都必须遵守该限制器运算符。我想知道的是,我是否可以有条件地将值保留在队列中,并随心所欲地选择如何限制排队的值。
【问题讨论】:
标签: javascript node.js rxjs rxjs6