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

多个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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.26 18:48:25