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

如何通过同一PeerJS连接发送多条消息?大文件分片传输难题

解决PeerJS长连接下大量文件分片连续发送问题

问题根源分析

你的核心问题出在send函数的实现上:

  • 每次调用send都会给同一个连接重复绑定open事件,当连接真正打开时,所有绑定的监听都会触发,导致消息重复发送或队列混乱,最终只能成功发送少量分片;
  • 没有处理连接未打开时的消息缓存,直接调用send会触发连接未就绪错误;
  • 频繁创建/销毁连接会超出浏览器RTCPeerConnection的数量限制,必须复用长连接。

解决方案:连接状态管理+消息队列

重构PeerInstance类,给每个连接维护消息队列,仅在连接创建时绑定一次open事件,实现连接复用与消息可靠发送:

修改后的完整代码

export default class PeerInstance {
  instance: Peer | null
  uId: string
  dataListeners: DataListeners
  // 重构连接存储:每个连接包含DataConnection和消息队列
  connections: { 
    [key: string]: {
      conn: DataConnection | null,
      messageQueue: Array<{type: string, payload?: any}>
    }
  }

  constructor(uId: string, accessToken: string) {
    this.instance = null
    this.uId = uId
    // 初始化连接结构
    this.connections = {}

    this.dataListeners = {
      ['message']: new Set(),
      ['file:start']: new Set(),
      ['file:chunk:send']: new Set(),
      ['file:chunk:received']: new Set(),
      ['file:done']: new Set(),
      ['file:downloaded']: new Set(),
      ['file:accept']: new Set(),
      ['file:reject']: new Set(),
      ['status:check']: new Set(),
      ['status:response']: new Set(),
      ['status:disconnected']: new Set(),
    }

    this.init(accessToken)
  }

  async init(accessToken: string) {
    if (this.instance) return

    const { Peer } = await import('peerjs')

    this.instance = new Peer(this.uId, {
      token: accessToken,
      host: config.peerHost,
      port: config.peerPort,
      path: config.peerPath,
    })

    this.instance.on('connection', (conn) => {
      conn.on('data', this.onData)
      // 处理主动连接的对方消息,同时如果本地有该连接的队列,绑定open事件
      const targetId = conn.peer
      if (!this.connections[targetId]) {
        this.connections[targetId] = {
          conn: conn,
          messageQueue: []
        }
      } else {
        this.connections[targetId].conn = conn
      }
      // 绑定open事件处理队列
      this.bindConnectionEvents(targetId)
    })
  }

  destroy = () => {
    if (!this.instance) return
    // 销毁所有连接
    Object.values(this.connections).forEach(item => {
      if (item.conn) {
        item.conn.close()
      }
    })
    this.instance.destroy()
  }

  onData = (data: any) => {
    this.fireDataListeners(data.type, data)
  }

  addDataListener = <T>(e: keyof DataListeners, fn: DataListener<T>) => {
    if (this.dataListeners[e]) {
      this.dataListeners[e].add(fn)
    }
  }

  removeDataListener = <T>(e: keyof DataListeners, fn: DataListener<T>) => {
    if (this.dataListeners[e]) {
      this.dataListeners[e].delete(fn)
    }
  }

  fireDataListeners = <T>(e: keyof DataListeners, data: any) => {
    if (this.dataListeners[e]) {
      this.dataListeners[e].forEach((listener: DataListener<T>) => {
        listener(data)
      })
    }
  }

  // 绑定连接的open/error/close事件,仅执行一次
  bindConnectionEvents = (to: string) => {
    const item = this.connections[to]
    if (!item || !item.conn) return

    const conn = item.conn

    // 仅绑定一次open事件
    conn.once('open', () => {
      // 发送队列中所有缓存的消息
      while (item.messageQueue.length > 0) {
        const msg = item.messageQueue.shift()
        if (msg) {
          conn.send({
            type: msg.type,
            to,
            from: this.uId,
            payload: msg.payload,
          })
        }
      }
    })

    conn.on('error', (err) => {
      console.error(`连接${to}出错:`, err)
      // 可选:清空队列或重新连接
      item.messageQueue = []
    })

    conn.on('close', () => {
      console.log(`连接${to}已关闭`)
      // 标记连接为null,下次发送时重新创建
      item.conn = null
    })
  }

  createConnection = (to: string) => {
    if (!this.instance) return

    const conn = this.instance.connect(to)
    // 初始化连接的队列结构
    this.connections[to] = {
      conn: conn,
      messageQueue: []
    }
    // 绑定事件
    this.bindConnectionEvents(to)

    return conn
  }

  send = (type: string, to: string, payload?: any) => {
    if (!this.instance) return

    let item = this.connections[to]
    // 如果连接不存在,创建连接
    if (!item) {
      this.createConnection(to)
      item = this.connections[to]
    }

    const conn = item?.conn
    if (!conn) return

    // 判断连接状态:已打开直接发送,未打开则加入队列
    if (conn.open) {
      conn.send({
        type,
        to,
        from: this.uId,
        payload,
      })
    } else {
      item.messageQueue.push({type, payload})
    }
  }
}

关键改进点说明

  1. 连接结构重构:每个目标用户的连接存储DataConnection实例和messageQueue数组,用于缓存连接未就绪时的消息;
  2. 单次事件绑定:使用conn.once('open')替代conn.on('open'),确保每个连接仅处理一次打开事件,避免重复触发;
  3. 消息队列处理:连接打开后自动清空队列发送所有缓存消息,解决未就绪时的发送报错问题;
  4. 连接复用:不再频繁创建/销毁连接,仅在连接关闭后下次发送时重新创建,避免RTCPeerConnection数量超限;
  5. 异常处理:新增连接错误、关闭事件的处理,保证连接状态可控。

文件下载进度跟踪方案

在上述基础上,你可以通过以下方式实现进度跟踪:

  • 发送文件分片时,在payload中携带分片索引和总分片数;
  • 接收方收到分片后,发送file:chunk:received确认消息,携带已接收的分片索引;
  • 发送方维护已确认的分片计数,通过已确认数/总分片数计算实时发送进度;
  • 接收方维护已接收的分片集合,全部接收完成后发送file:done消息。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.02 14:24:53