如何用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) }); });
关键改进点
- 修复
_id覆盖问题:通过扁平化流时绑定每个元素对应的_id,确保数据库更新时关联正确的文档。 - 批量频率控制:
bufferCount(10)将请求按10个一组打包,每组处理完成后延迟5秒,严格控制请求频率。 - 错误处理增强:为API请求添加
catchError,避免单个请求失败中断整个处理流;同时增加MongoDB查询错误判断。 - 逻辑解耦:将请求扁平化、批量处理、结果更新拆分为独立环节,代码更易维护。
可选优化
- 调整
bufferCount数值(如20)和延迟时间(如3000),适配不同API的频率限制。 - 添加
retry(2)操作符,对失败的API请求进行重试。 - 增加进度统计,通过
tap记录已处理的请求数量,方便监控执行进度。
内容的提问来源于stack exchange,提问作者Teshtek
相关产品推荐
相关产品推荐

