You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.23 17:15:07