跳转到内容

Dispatcher 分发架构

Dispatcher 是 Monibuca V6 中负责将帧数据从 Publisher 分发到所有 Subscriber 的核心组件。它的设计目标是消除重复读取,实现真正的 O(1) 读取 + O(N) 广播。

在传统的流媒体服务器中,每个订阅者独立从缓冲区读取数据:

flowchart LR
  RB["RingBuffer"]
  R1["Reader 1"] --> S1["Subscriber 1"]
  R2["Reader 2"] --> S2["Subscriber 2"]
  R3["Reader 3"] --> S3["Subscriber 3"]
  RN["Reader N"] --> SN["Subscriber N"]
  RB --> R1 & R2 & R3 & RN

传统方案问题:N 个订阅者 = N 次读取同一帧;每 Reader 各自追踪状态;高并发原子操作竞争大。

flowchart TB
  Pub["Publisher"] --> RB["RingBuffer"]
  RB --> Disp["Dispatcher<br/>单线程读 1 次 · Arc 零拷贝 · try_send"]
  Disp --> Q1["Queue 1 bounded"] & Q2["Queue 2 bounded"] & Q3["Queue 3 bounded"] & QN["Queue N bounded"]
  Q1 --> W1["Writer 1 task"]
  Q2 --> W2["Writer 2 task"]
  Q3 --> W3["Writer 3 task"]
  QN --> WN["Writer N task"]

核心优势:

指标传统方案Dispatcher 方案
RingBuffer 读取次数N 次/帧1 次/帧
读锁竞争O(N)O(1)
帧数据拷贝0(Arc 共享)0(Arc 共享)
慢订阅者影响可能阻塞丢帧,不阻塞
pub struct Dispatcher {
stream_path: String,
// 订阅者列表 - ArcSwap 无锁读取 (COW 模式)
subscribers: ArcSwap<Vec<DispatchSubscriber>>,
next_id: AtomicU64,
running: AtomicBool,
queue_capacity: usize, // 默认 150
total_dispatched: AtomicU64,
}

每个 Stream 对应一个 Dispatcher 实例。

每个订阅者拥有一个 bounded channel(默认容量 150):

Queue 容量 = 150 帧(≈ 1.8 秒 @ 80fps 混合帧率),可容忍网络抖动。

当队列满时,新帧会被 丢弃try_send 返回 Full),而不是阻塞其他订阅者。丢帧计数器会自动递增,供监控使用。

Dispatcher 通过 channel 发送以下帧类型:

pub enum DispatchFrame {
Video(Arc<AVFrame>), // 视频帧
Audio(Arc<AVFrame>), // 音频帧
VideoSeqHeader(Bytes), // 视频序列头(AVC/HEVC 解码器配置)
AudioSeqHeader(Bytes), // 音频序列头(AAC 解码器配置)
Eos, // 流结束信号
}
flowchart LR
  Pub["Publisher 30fps"] --> Disp["Dispatcher"]
  Disp -->|try_send OK| Qok["Queue 部分满<br/>Subscriber 正常"]
  Disp -->|try_send FULL| Qfull["Queue 满<br/>丢帧 · frames_dropped++ · 不阻塞"]

设计原则: 慢速订阅者只影响自己,不会拖慢整个系统。

  • 网络波动导致的短暂队列积压可以被缓冲吸收
  • 持续的慢速消费会导致帧丢弃,客户端通常可以自行恢复
  • 每个订阅者的丢帧计数可通过 API 监控

订阅者列表使用 ArcSwap<Vec<DispatchSubscriber>> 实现无锁管理:

// Clone-on-Write: 不阻塞正在进行的分发
self.subscribers.rcu(|old| {
let mut new = (**old).clone();
new.push(subscriber);
new
});
// 原子加载 - 无锁
let subs = self.subscribers.load();
for sub in subs.iter() {
if sub.receive_video && !sub.is_closed() {
sub.try_send(DispatchFrame::Video(frame.clone()));
}
}

在 1000 个订阅者 @ 60fps 的场景下,分发操作完全无锁竞争。

当服务器需要处理大量并发流(>100)时,可以启用 DispatcherPool 模式。

flowchart TB
  Pool["DispatcherPool<br/>管理 N 个 Worker"]
  W0["Worker 0<br/>Stream A / B"]
  W1["Worker 1<br/>Stream C / D"]
  WN["Worker N-1<br/>Stream E / F"]
  Pool --> W0 & W1 & WN

流通过 一致性哈希 分配到 Worker:

fn worker_index(&self, stream_path: &str) -> usize {
let mut hasher = DefaultHasher::new();
stream_path.hash(&mut hasher);
(hasher.finish() as usize) % self.num_workers
}

同一个流路径始终被分配到相同的 Worker,保证了流的处理连续性。

通过 dispatcher_workers 配置项控制:

模式说明
0Per-Stream每个流一个独立的 Dispatcher task(默认)
NPoolN 个 Worker,每个处理多个流
# 配置示例
stream:
dispatcher_workers: 0 # Per-Stream 模式(默认)
dispatcher_workers: 4 # 4 个 Worker 的池模式
dispatcher_workers: 8 # 8 个 Worker 的池模式

选择建议:

  • 少量高并发流(< 100 路): 使用 0(Per-Stream),每路流独享一个 task
  • 大量流(> 100 路): 使用 Pool 模式,减少 task 开销
  • Worker 数量建议设置为 CPU 核心数的 1~2 倍
sequenceDiagram
  participant Pub as Publisher
  participant Track as VideoTrack.buffer
  participant Notify as frame_notify watch
  participant Disp as Dispatcher
  participant Sub as Subscriber
  participant Enc as Protocol encoder

  Pub->>Track: write_video(frame)
  Track->>Notify: send(count)
  Notify->>Disp: changed().await
  Disp->>Track: read_next() once → Arc AVFrame
  Disp->>Sub: try_send Video frame.clone
  Sub->>Sub: recv().await
  Sub->>Enc: RTMP / FLV / HLS / WebRTC

Dispatcher 每 100 个帧分发周期执行一次清理,移除已关闭的订阅者:

cleanup_counter += 1;
if cleanup_counter % 100 == 0 {
self.cleanup_closed(); // COW 模式移除已关闭的订阅者
}

当流结束时,Dispatcher 向所有订阅者发送 Eos(End of Stream)信号。

联系我们

微信公众号:不卡科技 微信公众号二维码
腾讯频道:流媒体技术 腾讯频道二维码
QQ 频道:p0qq0crz08 QQ 频道二维码
QQ 群:751639168 QQ 群二维码