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

如何用RxJS为多层数组的API请求加延迟,避免请求超限错误

MongoDB + RxJS 批量API请求频率控制方案

问题背景

需求

  • 从MongoDB获取含160个元素的主数组
  • 遍历主数组元素,取出FIELDTOUSE字段(包含900个子数组的数组结构)
  • 遍历每个子数组(约200个元素),为每个元素调用API
  • 需添加延迟控制:每10个API请求后延迟5秒,避免触发Too Many Requests错误

数据结构层级

  • 主数组:[X,Y,Z,...](160个元素)
  • 元素X:{FIELD1, FIELD2, FIELDTOUSE,...}
  • FIELDTOUSE:[EL1, EL2,...](900个子数组)
  • 子数组EL:[A,B,C,...](约200个元素)

请求量计算

单个FIELDTOUSE需发起 200*900=180,000 次请求,主数组总计 160*180,000=28,800,000 次请求,必须严格控制请求频率。

当前代码

function getAPIdata(res) {
    // 自定义逻辑
    if(/* 条件判断 */){
      return axios.post(
        urlOTP,
        stringLL,
        { headers: headers }
      ) 
    }else{
      return of(null).pipe(delay(1000))
    }
}

// 调用MongoDB集合
XModel.find({}).lean().exec((err, ELEMENTS) => {
    // 变量声明
    // 160个元素
    from(ELEMENTS).pipe(concatMap(el => {
        // 900个元素
        return from(el.x).pipe(concatMap(el_ => {
            _id = el._id;
            // 希望每10个元素延迟5秒
            return getAPIdata(el_) // <-------------------- 此处需添加延迟控制
        }))
    }),concatMap(g => g.data.hasOwnProperty("results") ? of(g.data.results).pipe(delay(1000)) : of(null).pipe(delay(1000)))).subscribe(r => {
        // 更新数据库逻辑
        XModel.updateOne({ _id: _id }, {$set:set}, (e, done) => {
            // 后续逻辑
        })
    });
});

解决方案

核心思路

利用RxJS的bufferCount(批量分组)、concatMap(串行处理组)和delay(组间延迟)操作符,实现每10个请求后延迟5秒的频率控制,同时修复原代码中_id被覆盖的问题。

修改后的完整代码

import { from, of, map } from 'rxjs';
import { concatMap, delay, bufferCount, mergeMap, tap, catchError } from 'rxjs/operators';
import axios from 'axios';

function getAPIdata(res) {
    // 自定义逻辑
    if(/* 条件判断 */){
      return axios.post(
        urlOTP,
        stringLL,
        { headers: headers }
      ).pipe(
        catchError(err => {
          console.error('API请求失败:', err);
          return of(null);
        })
      );
    }else{
      return of(null).pipe(delay(1000));
    }
}

// 调用MongoDB集合
XModel.find({}).lean().exec((err, ELEMENTS) => {
    if (err) {
        console.error('MongoDB查询失败:', err);
        return;
    }

    // 扁平化所有请求元素,绑定对应的_id
    const allRequests$ = from(ELEMENTS).pipe(
        concatMap(el => 
          from(el.x).pipe(
            concatMap(subArray => 
              from(subArray).pipe(
                map(item => ({ item, _id: el._id }))
              )
            )
          )
        )
    );

    allRequests$.pipe(
        // 每10个请求分为一组
        bufferCount(10),
        // 串行处理每组,避免并发过高
        concatMap(group => 
          from(group).pipe(
            // 并行处理组内10个请求(可替换为concatMap改为串行)
            mergeMap(({ item, _id }) => 
              getAPIdata(item).pipe(
                tap(apiRes => {
                  // 处理API结果并更新数据库
                  if (apiRes?.data?.results) {
                    const set = { /* 构造需要更新的字段 */ };
                    XModel.updateOne({ _id }, { $set: set }).exec(updateErr => {
                      if (updateErr) console.error('数据库更新失败:', updateErr);
                    });
                  }
                })
              )
            ),
            // 每组处理完成后延迟5秒
            delay(5000)
          )
        )
    ).subscribe({
        complete: () => console.log('所有请求处理完成'),
        error: err => console.error('流程执行出错:', err)
    });
});

关键改进点

  1. 修复_id覆盖问题:通过扁平化流时绑定每个元素对应的_id,确保数据库更新时关联正确的文档。
  2. 批量频率控制:bufferCount(10)将请求按10个一组打包,每组处理完成后延迟5秒,严格控制请求频率。
  3. 错误处理增强:为API请求添加catchError,避免单个请求失败中断整个处理流;同时增加MongoDB查询错误判断。
  4. 逻辑解耦:将请求扁平化、批量处理、结果更新拆分为独立环节,代码更易维护。

可选优化

  • 调整bufferCount数值(如20)和延迟时间(如3000),适配不同API的频率限制。
  • 添加retry(2)操作符,对失败的API请求进行重试。
  • 增加进度统计,通过tap记录已处理的请求数量,方便监控执行进度。

内容的提问来源于stack exchange,提问作者Teshtek

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.22 11:24:37