前端控制并发上传请求个数
在日常开发中我们经常会遇到批量上传, 并发请求等需求, 在大部分情况下服务端的承载能力是有限的, 因此需要前端控制每次并发数量, 减少服务端压力.
具体代码
javascript
export async function parallel<T>(
jobs: T[],
fn: (params: T) => Promise<any>,
limit = 2,
successBack: (args: { status: 'success' | 'error' }[]) => void = () => {}
) {
const ret: { status: 'success' | 'error' }[] = [];
let cursor = 0;
async function worker() {
let currentJob;
while (cursor < jobs.length) {
try {
currentJob = cursor;
cursor += 1;
await fn(jobs[currentJob])
.then(() => ret.push({ status: 'success' }))
.catch(() => {
ret.push({ status: 'error' });
})
.finally(() => successBack(ret));
} catch (e: any) {
throw new Error(e?.message ?? '请重试');
}
}
}
const workers = [];
for (let i = 0; i < limit && i < jobs.length; i += 1) {
workers.push(worker());
}
await Promise.allSettled(workers);
return Promise.resolve(ret);
}基础使用demo
场景: 前端选择了多个文件(记为 files), 随后需要调用上传方法将文件发送给服务端(记为upload), 由于服务端能力有限, 因此前端要控制同一时间最多有两个上传请求
伪代码
javascript
const files: File[] = [...]
const uploadFns = async (file: File ) => {
return await upload(file).catch(e => {
// 处理错误文件等业务相关逻辑(非必要)
setErrorFiles(pre => {
pre.push({
name: file.name,
error: e.message,
});
return pre;
});
});
};
parallel<File>(files, uploadFns, 2).finally(() => { // 上传完成后业务逻辑 })实现原理
- 初始将limit个请求塞进works中 (注意每一个promise请求都放在了while循环中)
- 当works中有某个请求结束了, 由于该请求处于while循环中, 就会继续往后走, 取下一个请求然后执行 (所有while循环都共享一个全局cursor)
- 当没有剩余的请求, 所有的while循环都会停止, 至此请求处理完毕
进一步优化
可以看到, 上述方法只是一个纯工具类函数, 只管发送请求, 不负责打断操作, 因此可以进一步完善.
tsx
/**
* 批量发送请求
*/
class ParallelUpload {
// 并发请求数
protected limit = 2;
// 存储请求序列
protected requestQueue: Promise<any>[] = []
// 当前请求下标
protected cursor = 0
// 需要发送的请求入参
protected jobs: string[] = []
// 外部传入的请求方法
protected fn: (job: string, abort?: AbortController) => Promise<any> = () => Promise.resolve()
// 判断是否正在上传
protected isUploading: boolean = false
// 控制请求打断
protected abortControllers: AbortController[] = []
constructor(jobs: string | string[], fn: (job: string, abort?: AbortController) => Promise<any>, limit?: number) {
this.jobs = Array.isArray(jobs) ? jobs : [jobs]
this.fn = fn
limit && (this.limit = limit);
this.requestQueue = new Array(jobs.length)
}
private async handle() {
let currentJob;
while (this.cursor < this.jobs.length) {
if (!this.isUploading) {
break;
}
try {
const abort = new AbortController()
currentJob = this.cursor;
this.cursor += 1;
this.abortControllers[currentJob] = abort
this.requestQueue[currentJob] = await this.fn(this.jobs[currentJob], abort);
} catch (e) {
console.log(`job: ${currentJob}`, e);
}
}
}
async start() {
this.isUploading = true
const works = []
for (let i = 0; i < this.limit ; i++) {
works.push(this.handle())
}
await Promise.all(works);
this.init()
return works
}
stop() {
this.abortControllers.forEach(controller => controller.abort())
this.init()
}
protected init() {
this.cursor = 0
this.requestQueue = []
this.abortControllers = []
this.isUploading = false
}
}测试demo
tsx
const allRequest = Array.from(Array(10), (v,k) => k.toString());
async function fetchFn(url: string, abort?: AbortController) {
const time = Math.random() * 2000
return new Promise(resolve => {
setTimeout(() => {
console.log(url, time)
fetch(`http://localhost:4000/about?index=${url}`, abort).then(response => {
if (!abort?.signal.aborted) {
console.log('请求成功,数据:');
} else {
console.log('请求被取消');
}
})
.catch(error => {
if (error.name === 'AbortError') {
console.log('请求被取消,错误:AbortError', error);
} else {
console.error('请求失败,错误:Other', error);
}
});
resolve(url)
}, time)
})
}
const parallel = new ParallelUpload(allRequest, fetchFn, 2)
parallel.start()
setTimeout(() => {
parallel.stop()
}, 1500)