多个fetch API并发分片请求时如何实现断点续传
大文件分片并发下载构造ReadableStream的异常处理方案
问题描述
需要对大文件发起分片请求,并基于响应构造自定义ReadableStream。
最初采用单分片串行请求方案,已经实现断点续传能力应对网络不稳定问题,初始实现代码如下:
let _this = this let responseList = [] let response = null let reader = null let offset = 0 let start = 0 let end = start + this.RANGE_SIZE return new ReadableStream({ async start(controller) { responseList.push(await fetch(_this.url, {responseType: 'blob', headers: {range: `bytes=${start}-${end}`}})) reader = responseList[0].body.getReader() }, async pull(controller) { try { var {done, value} = await reader.read() } catch (e) { let errorOffset = 0 responseList[0] = await fetch(_this.url, {responseType: 'blob', headers: {range: `bytes=${start}-${end}`}}) reader = responseList[0].body.getReader() while (true) { let {done, value} = await reader.read() if (errorOffset + value.length > offset){ value = value.slice(offset - errorOffset) controller.enqueue(value) return }else { errorOffset += value.length } } } if (done) { offset = 0 start = end + 1 end += _this.RANGE_SIZE responseList.shift() responseList.push(await fetch(_this.url, {responseType: 'blob', headers: {range: `bytes=${start}-${end}`}})) reader = responseList[0].body.getReader() const {done, value} = await reader.read() controller.enqueue(value) offset += value.length } else { controller.enqueue(value) offset += value.length } } })
为提升下载效率,调整为同时预加载2个分片,修改后的start方法如下:
async start(controller) { responseList.push(await fetch(_this.url, {responseType: 'blob', headers: {range: `bytes=${start}-${end}`}})) start = end + 1 end += _this.RANGE_SIZE responseList.push(await fetch(_this.url, {responseType: 'blob', headers: {range: `bytes=${start}-${end}`}})) reader = responseList[0].body.getReader() },
该实现存在缺陷:第二个分片的请求、读取错误无法被原有try...catch逻辑捕获,容易出现静默失败。
解决方案
核心逻辑是不要让预加载的分片请求脱离错误监听链路,具体调整如下:
- 改造分片存储结构,不再只存储response对象,给每个分片维护独立状态:分片起止范围、response实例、reader实例、错误信息、加载完成标记
- 封装统一的分片加载方法,内部直接捕获fetch、reader读取阶段的所有异常,存入对应分片的状态字段,从根源上避免未捕获的Promise reject
- 切换读取分片前,先校验对应分片的状态,如果预加载阶段已经记录了错误,直接复用原有断点续传的重试逻辑重新拉取分片
- 新增流取消回调,在流终止时主动释放所有预加载分片的reader资源,避免内存泄漏
- 所有分片的致命错误(重试后仍失败)统一调用
controller.error()抛出,保证流的错误状态可被外部感知
调整后的核心实现代码:
let _this = this const MAX_CONCURRENT_CHUNK = 2 // 并发预加载分片数 // 分片列表存储每个分片的完整状态 let chunkList = [] let currentReader = null let currentChunkIndex = 0 let offsetInCurrentChunk = 0 let nextChunkStart = 0 let nextChunkEnd = nextChunkStart + this.RANGE_SIZE // 统一封装分片加载逻辑,全链路捕获异常 const loadChunk = async (chunkStart, chunkEnd) => { const chunkMeta = { start: chunkStart, end: chunkEnd, response: null, reader: null, error: null, isLoaded: false } try { chunkMeta.response = await fetch(_this.url, { responseType: 'blob', headers: { range: `bytes=${chunkStart}-${chunkEnd}` } }) chunkMeta.reader = chunkMeta.response.body.getReader() chunkMeta.isLoaded = true } catch (err) { chunkMeta.error = err } return chunkMeta } return new ReadableStream({ async start(controller) { // 初始化预加载指定数量的分片,所有加载Promise都被托管,不会出现未捕获异常 for (let i = 0; i < MAX_CONCURRENT_CHUNK; i++) { const chunkStart = nextChunkStart + i * _this.RANGE_SIZE const chunkEnd = chunkStart + _this.RANGE_SIZE - 1 chunkList.push(loadChunk(chunkStart, chunkEnd)) } // 等待第一个分片加载完成,初始化当前读取器 chunkList[0] = await chunkList[0] if (chunkList[0].error) { controller.error(chunkList[0].error) return } currentReader = chunkList[0].reader // 更新下一个待预加载的分片位置 nextChunkStart = nextChunkEnd + 1 nextChunkEnd += _this.RANGE_SIZE }, async pull(controller) { let readResult try { readResult = await currentReader.read() } catch (e) { // 当前分片读取失败,走重试逻辑 const currentChunk = chunkList[currentChunkIndex] const retryChunk = await loadChunk(currentChunk.start, currentChunk.end) if (retryChunk.error) { controller.error(retryChunk.error) return } chunkList[currentChunkIndex] = retryChunk currentReader = retryChunk.reader // 定位到已读取的偏移位置,续传数据 let readOffsetInRetry = 0 while (true) { const { done, value } = await currentReader.read() if (readOffsetInRetry + value.length > offsetInCurrentChunk) { const validValue = value.slice(offsetInCurrentChunk - readOffsetInRetry) controller.enqueue(validValue) offsetInCurrentChunk += validValue.length return } readOffsetInRetry += value.length } } const { done, value } = readResult if (done) { // 当前分片读取完成,切换到下一个分片 currentChunkIndex++ offsetInCurrentChunk = 0 // 等待预加载的下一个分片就绪,此处可直接拿到预加载阶段捕获的错误 const nextChunk = await chunkList[currentChunkIndex] if (nextChunk.error) { // 预加载失败,重试该分片 const retryChunk = await loadChunk(nextChunk.start, nextChunk.end) if (retryChunk.error) { controller.error(retryChunk.error) return } chunkList[currentChunkIndex] = retryChunk } currentReader = chunkList[currentChunkIndex].reader // 追加预加载新的后续分片,保持固定并发数 chunkList.push(loadChunk(nextChunkStart, nextChunkEnd)) nextChunkStart = nextChunkEnd + 1 nextChunkEnd += _this.RANGE_SIZE // 读取新分片的第一个数据块入队 const firstChunkResult = await currentReader.read() controller.enqueue(firstChunkResult.value) offsetInCurrentChunk += firstChunkResult.value.length } else { controller.enqueue(value) offsetInCurrentChunk += value.length } }, cancel(reason) { // 流取消时主动释放所有reader资源 chunkList.forEach(chunk => { if (chunk?.reader) chunk.reader.cancel(reason) }) } })
可选优化项
- 增加分片加载重试次数上限,避免异常场景下无限重试导致流卡死
- 可根据实际网络质量动态调整并发预加载的分片数,弱网环境下适当降低并发减少请求失败概率
- 如果业务有分片校验逻辑,可在分片加载完成后提前校验数据完整性,发现损坏直接触发重传
内容的提问来源于stack exchange,提问作者Yue Zhao
相关产品推荐
相关产品推荐

