如何搭多人会议系统(SFU 架构 / 多路流管理)

三个人以上的会议系统如何搭建?Mesh 还是 SFU?订阅哪几路流?大画面切谁?本文用 Claude Code 设计并实现一个可扩展的多人会议系统,从架构选型到多路流管理。

1、多人会议的架构困境

一对一通话简单:A 推流给 B,B 推流给 A。三个人呢?

方案 1:Mesh(全互联)          方案 2:SFU(选择性转发)
  A ↔ B                          A ──┐
  B ↔ C                          B ──┼──→ SFU ──→ 各端
  A ↔ C                          C ──┘
  ↑ 每人上行 2 路,下行 2 路        ↑ 每人上行 1 路,下行按需订阅
  • 🕸️ Mesh:每人推给其他所有人,N 个人上行 N-1 路。3 人还行,6 人就崩(手机上行带宽不够)
  • 🔀 SFU:每人只推 1 路到服务器,服务器转发。上行永远 1 路,下行按需订阅
  • 🎛️ MCU:服务器把所有流「合成」1 路再下发。省下行带宽,但服务器要解码重编码,成本极高

结论: 现代多人会议几乎都是 SFU(WebRTC 的默认架构)。上行固定 1 路,客户端只订阅「需要的」几路,服务器只做转发不重编码。

2、SFU 架构总览

┌─────────────────────────────────────────────┐
│              SFU 服务端                       │
│                                              │
│  ┌──────────┐   ┌──────────┐   ┌──────────┐ │
│  │ 房间管理   │   │ 流管理    │   │ 转发引擎  │ │
│  │ Room      │   │ Track    │   │ Forward  │ │
│  └────┬─────┘   └────┬─────┘   └────┬─────┘ │
│       │              │              │       │
│  ┌────▼──────────────▼──────────────▼────┐  │
│  │      信令 (WebSocket / Socket.io)      │  │
│  └───────────────────────────────────────┘  │
└─────────────────────────────────────────────┘
         ▲                    │
    WebSocket             RTP (SRTP)
         │                    ▼
┌────────┴──────────────────────────┐
│            客户端                    │
│  ┌────────┐  ┌────────┐  ┌──────┐ │
│  │ 推 1 路 │  │ 订阅 N 路│  │ 布局  │ │
│  └────────┘  └────────┘  └──────┘ │
└────────────────────────────────────┘

信令(WebSocket)管「谁在和谁说话」,媒体(RTP/SRTP)管「音视频数据」。两者分离是核心设计。

3、服务端:信令 + 房间 + 流管理

3.1、Prompt

帮我写一个 SFU 会议服务端的信令层(Node.js + Socket.io)。

【功能】
1. 房间管理:创建/加入/离开房间
2. 流管理:加入者 publish 自己的流,subscribe 别人的流
3. 转发管理:记录「谁订阅了谁」,为 SFU 转发提供依据
4. 状态同步:新加入者收到当前房间所有参与者列表
5. 级联预留:支持房间上限,超出自动分片

【关注点】
- 断线重连:用户网络抖动掉线,重连后恢复订阅关系
- 竞态:A 加入时 B 还没 publish,B 之后 publish 要广播给 A
- 房间清理:空房间自动回收

3.2、Node.js 信令服务

// signaling.js (Node.js + Socket.io)
// SFU 会议信令层

const { Server } = require('socket.io');

class ConferenceServer {
  constructor(server) {
    this.io = new Server(server, { cors: { origin: '*' } });

    // 房间 → { participantId: { sockets: Set, published: Set } }
    this.rooms = new Map();

    this.io.on('connection', (socket) => this.onConnection(socket));
  }

  onConnection(socket) {
    let currentRoom = null;
    let participantId = null;

    // 加入房间
    socket.on('join', ({ roomId, userId }) => {
      currentRoom = roomId;
      participantId = userId;
      socket.join(roomId);

      // 初始化房间状态
      if (!this.rooms.has(roomId)) {
        this.rooms.set(roomId, new Map());
      }
      const room = this.rooms.get(roomId);
      if (!room.has(userId)) {
        room.set(userId, { sockets: new Set(), published: new Set() });
      }
      room.get(userId).sockets.add(socket.id);

      // 通知房间里其他人:有人加入
      socket.to(roomId).emit('peer-joined', { userId });

      // 回执:告诉新加入者当前所有参与者 + 他们已 publish 的流
      const peers = [...room.entries()]
        .filter(([id]) => id !== userId)
        .map(([id, state]) => ({ userId: id, streams: [...state.published] }));
      socket.emit('room-state', { peers });
    });

    // 发布流(客户端本地媒体就绪后)
    socket.on('publish', ({ streamId, kind }) => {
      if (!currentRoom || !participantId) return;
      const state = this.rooms.get(currentRoom).get(participantId);
      state.published.add(streamId);

      // 广播:新流可用,其他端自行决定是否订阅
      socket.to(currentRoom).emit('stream-published', {
        userId: participantId, streamId, kind
      });
    });

    // 订阅某路流
    socket.on('subscribe', ({ targetUserId, streamId }) => {
      // 记录订阅关系 → 通知 SFU 转发引擎建立转发
      socket.to(currentRoom).emit('forward-request', {
        from: participantId, to: targetUserId, streamId
      });
    });

    // 取消订阅
    socket.on('unsubscribe', ({ targetUserId, streamId }) => {
      socket.to(currentRoom).emit('forward-cancel', {
        from: participantId, to: targetUserId, streamId
      });
    });

    // 断开连接
    socket.on('disconnect', () => {
      if (!currentRoom || !participantId) return;
      const room = this.rooms.get(currentRoom);
      if (!room) return;

      const state = room.get(participantId);
      state.sockets.delete(socket.id);

      // 该用户所有 socket 都断了才算真离开
      if (state.sockets.size === 0) {
        room.delete(participantId);
        socket.to(currentRoom).emit('peer-left', { userId: participantId });

        // 空房间回收
        if (room.size === 0) {
          this.rooms.delete(currentRoom);
          console.log(`[Room] ${currentRoom} 已回收`);
        }
      }
    });
  }
}

4、客户端:多路流管理

客户端的难点:管理「我的 1 路上行」+「按需的 N 路下行」,以及大画面的动态切换。

4.1、多路流管理器

// StreamManager.swift (iOS)
// 管理本地上行流 + 多路远端下行流

import WebRTC

final class StreamManager: NSObject {

    // 本地流(上行,1 路)
    private(set) var localStream: RTCMediaStream?

    // 远端流(下行,N 路)userId -> stream
    private(set) var remoteStreams: [String: RTCMediaStream] = [:]

    // 大画面用户(当前显示谁)
    private(set) var activeSpeakerId: String?

    // 音频能量检测(用于自动切大画面)
    private var audioLevels: [String: Double] = [:]

    /// 发布本地流
    func publish(capturedStream: RTCMediaStream) {
        localStream = capturedStream
        signaling.publish(streamId: capturedStream.streamId, kind: "av")
    }

    /// 远端有新流 → 决定是否订阅
    func onStreamPublished(userId: String, streamId: String, kind: String) {
        // 策略:默认订阅所有人(小会议);大会议按需订阅
        if remoteStreams.count < 9 {  // 最多显示 9 宫格
            signaling.subscribe(targetUserId: userId, streamId: streamId)
        }
    }

    /// 远端流到达
    func onRemoteStream(_ stream: RTCMediaStream, from userId: String) {
        remoteStreams[userId] = stream
        render(stream: stream, to: tileView(for: userId))
    }

    /// 自动切大画面:根据音频能量选「谁在说话」
    func onAudioLevel(_ level: Double, from userId: String) {
        audioLevels[userId] = level
        let speaker = audioLevels.max(by: { $0.value < $1.value })?.key

        if speaker != activeSpeakerId, let speaker {
            activeSpeakerId = speaker
            switchActiveSpeaker(to: speaker)
        }
    }

    /// 切大画面(渲染到主窗口)
    private func switchActiveSpeaker(to userId: String) {
        guard let stream = remoteStreams[userId] else { return }
        mainView.render(stream: stream)
        // 说话的人放大到主窗口,其余保持小窗
    }
}

4.2、Android 端订阅策略

// StreamManager.kt (Android)
// 多路流订阅策略:按需订阅 + 分页

class StreamManager {

    private val remoteStreams = mutableMapOf<String, RTCMediaStream>()
    private val subscribedUsers = mutableSetOf<String>()

    // 订阅策略枚举
    enum class SubscriptionPolicy {
        SUBSCRIBE_ALL,      // 小会议:全订阅
        SUBSCRIBE_ACTIVE,   // 大会议:只订阅说话的人 + 最近 N 人
        SUBSCRIBE_PAGED     // 超大会议:分页,每页 9 人
    }

    var policy = SubscriptionPolicy.SUBSCRIBE_ALL
        set(value) {
            field = value
            applyPolicy()
        }

    fun onStreamPublished(userId: String, kind: String) {
        when (policy) {
            SubscriptionPolicy.SUBSCRIBE_ALL -> subscribe(userId)
            SubscriptionPolicy.SUBSCRIBE_ACTIVE -> {
                // 只自动订阅前 9 个,后续靠说话触发订阅
                if (subscribedUsers.size < 9) subscribe(userId)
            }
            SubscriptionPolicy.SUBSCRIBE_PAGED -> {
                // 分页:用户主动翻页才订阅
            }
        }
    }

    /** 说话触发订阅:某人开始说话,如果没订阅就订阅 */
    fun onActiveSpeaker(userId: String) {
        if (policy == SubscriptionPolicy.SUBSCRIBE_ACTIVE && !subscribedUsers.contains(userId)) {
            subscribe(userId)
            // 订阅新的同时,取消订阅最不活跃的一个,控制下行带宽
            evictLeastActive()
        }
    }

    private fun subscribe(userId: String) {
        if (subscribedUsers.add(userId)) {
            signaling.subscribe(userId)
        }
    }

    private fun evictLeastActive() {
        if (subscribedUsers.size > 9) {
            val toEvict = subscribedUsers.first()  // 简化:踢掉最早的
            subscribedUsers.remove(toEvict)
            signaling.unsubscribe(toEvict)
        }
    }
}

5、大画面切换:三种模式

多人会议最考验产品的地方是「谁是大画面」。三种模式:

1. 说话切换模式(声控)
   谁说话谁放大 → 适合会议讨论

2. 固定模式(主讲人锁定)
   主持人锁定某人大画面 → 适合讲课/培训

3. 手动切换
   用户自己点谁谁放大 → 适合灵活场景

5.1、声控切换的防抖

// ActiveSpeakerDetector.swift (iOS)
// 声控切换:音频能量 + 防抖 + 最小说话时长

final class ActiveSpeakerDetector {

    private var energyHistory: [String: [Double]] = [:]
    private var currentSpeaker: String?
    private var speakerSince: Date?

    // 阈值参数
    let minSpeakDuration: TimeInterval = 0.5   // 至少说 0.5 秒才切
    let switchHysteresis: Double = 1.5          // 新说话人要高 50% 才切

    func onEnergy(_ energy: Double, userId: String) {
        // 滑动窗口保留最近 10 个能量样本
        var history = energyHistory[userId] ?? []
        history.append(energy)
        if history.count > 10 { history.removeFirst() }
        energyHistory[userId] = history

        let avg = history.reduce(0, +) / Double(history.count)

        if let speaker = currentSpeaker, speaker != userId {
            // 切换需要:新说话人能量 > 当前说话人 × hysteresis
            let speakerAvg = (energyHistory[speaker] ?? [0]).reduce(0, +) /
                Double((energyHistory[speaker] ?? [1]).count)
            guard avg > speakerAvg * switchHysteresis else { return }
        }

        // 防抖:持续说话 minSpeakDuration 才切换
        if speakerSince == nil { speakerSince = Date() }
        if Date().timeIntervalSince(speakerSince!) >= minSpeakDuration {
            currentSpeaker = userId
        }
    }
}

6、踩坑记录

#问题现象根因修复
1新人看不到老人后加入的看不到已发布流join 时没回传 room-statejoin 回执里带上所有已 publish 的流
2订阅关系错乱断线重连后订阅丢失disconnect 只清 socket 没清订阅重连时根据 room-state 重建订阅
3大画面疯狂切换画面每秒切好几次能量检测无防抖/迟滞加 minSpeakDuration + hysteresis
49 宫格卡爆6 人会议手机发烫默认全订阅 + 全渲染大会议改按需订阅 + 未显示的流不渲染
5房间泄漏服务器内存涨空房间不回收空房间自动 delete + 定时巡检
6回声严重多人会议回声多人用扬声器外放开启 AEC(回声消除),扬声器场景强制耳机

7、架构对比

维度MeshSFUMCU
上行带宽N-1 路1 路1 路
下行带宽N-1 路按需1 路
服务器成本转发(低)转码(高)
端上 CPU高(多路编解码)
扩展性差(≤4人)
延迟

结论: 3 人以下 Mesh 省成本,3 人以上 SFU 是标准答案。MCU 只在「端侧能力极弱」或「需要服务端录制合成画面」时用。

学习和提升音视频开发技术,欢迎你加入我们的知识星球

如何搭多人会议系统(SFU 架构 / 多路流管理)

版权声明:本文内容转自互联网,本文观点仅代表作者本人。本站仅提供信息存储空间服务,所有权归原作者所有。如发现本站有涉嫌抄袭侵权/违法违规的内容, 请发送邮件至1393616908@qq.com 举报,一经查实,本站将立刻删除。

(0)

相关推荐