RxJS带时间、数量及token校验的HTTP缓冲器26次请求卡住问题
问题背景
你正在开发一款HTTP请求缓冲器,可对HTTP请求进行缓存,满足以下任一条件时将发送缓冲的请求:
- 缓冲的调用次数>25
- 距离首个进入缓冲区的调用已过去x毫秒,且token observer状态为'valid'
现有实现代码
this.batcherObservable = new Subject<BatchItem>(); this.batcherObservable.pipe( tap(req => console.log(req.request.url)), bufferWhen(() => { // Buffer has opened const bufferStopt = new Subject() const limiter = this.batcherObservable.pipe(takeUntil(bufferStopt), bufferCount(26), take(1)) const timer = this.batcherObservable.pipe(takeUntil(bufferStopt), take(1), delay(this.options.batchDelay)) const tokensValid = this.batcherObservable.pipe(takeUntil(bufferStopt), take(1), delayWhen(() => this.tokenStatus.pipe(filter(state => state === 'valid'), take(1)))) return race(limiter, forkJoin([timer, tokensValid])).pipe(tap( () => { bufferStopt.complete() })) }), filter(requests => requests.length > 0), tap(req => console.log(req.length)), delayWhen(() => this.tokenStatus.pipe(filter(state => state === 'valid'), take(1))), ).subscribe(async requests => { const useBatch = requests.length > 1 let request: Request = !useBatch ? requests[0].request : await this.#renderCallsToBatchCall(<[BatchItem, ...BatchItem[]]>requests) let response: globalThis.Response try { response = await firstValueFrom(this.#sendRequest(request)) } catch (error) { console.error('error in batch', error) if (error instanceof globalThis.Response) { response = await error.json() } else { console.error(error) return new Error('Error when requesting') } } if (!useBatch) { if(response.ok && (response.status >= 200 && response.status <= 299)) return requests[0].responseObserver.next(response) else return requests[0].responseObserver.error(response) } const batchBody: FlatResponse[]|CWHttpError = await response.json() // We have to flat out the batchcall and make our own responses requests.forEach((batchItem, index) => { if (Array.isArray(batchBody)) { const {body, headers, status, message} = batchBody[index] const url = requests[index].request.url response = new Response(JSON.stringify(body || {}), {headers, statusText: message, status: parseInt(status)}) // URL is private of Response. This is the only way I can set the url Object.defineProperty(response, 'url', { value: url}); if (response.ok) { batchItem.responseObserver.next(response) } else { batchItem.responseObserver.error(body) } batchItem.responseObserver.complete() } else { batchItem.responseObserver.error(response) } }) });
故障现象
测试发现25次调用时可正常触发请求发送,27次及以上数量的调用也运行正常,但刚好发起26次调用时缓冲区不会触发,请求一直卡在缓冲中。
故障原因
bufferWhen 操作符的工厂函数执行时机是第一个请求到达、缓冲区初始化完成之后,此时你才对this.batcherObservable订阅了limiter、timer、tokensValid三个流,因此这三个流永远收不到触发缓冲区初始化的第一个请求事件。
你的limiter用bufferCount(26)做计数判断,天然少统计了第一个请求:
- 缓冲区累计26个请求时,
limiter只收到了后25个,凑不够26个不会触发关闭 - 缓冲区累计27个请求时,
limiter刚好收到后26个,满足条件正常触发
这就是刚好26次调用卡住的核心原因。
修复方案
最简单的修复方式是将limiter的bufferCount(26)修改为bufferCount(25),刚好抵消第一个请求未被统计的差值,缓冲区累计到26个请求时就能正常触发关闭。
如果想要更严谨的实现,也可以改用scan操作符自行统计缓冲区请求数量,避免订阅时机带来的计数偏差:
const limiter = this.batcherObservable.pipe( takeUntil(bufferStopt), scan(count => count + 1, 0), // 从0开始计数后续进入的请求 filter(count => count >= 25), // 累计25个后续请求时,加上首个请求总共有26个,满足阈值 take(1) )
另外可以优化冗余订阅,将timer和tokensValid的触发逻辑合并,减少不必要的流订阅。
内容的提问来源于stack exchange,提问作者Martijn Bakker
相关产品推荐
相关产品推荐

