直播连麦与云直播服务
连麦技术原理
连麦架构
直播连麦的核心架构涉及三个角色:观众、主播和连麦方,以及CDN和合流转推服务。
┌─────────────────────────────────────┐
│ CDN / 源站 │
│ (低延迟拉流: HTTP-FLV / HLS / WebRTC) │
└──────┬──────────────────────┬────────┘
│ │
┌────────────┘ └────────────┐
↓ ↓
┌──────────┐ ┌──────────┐
│ 观众端 │ │ 观众端 │
│ (CDN拉流) │ │ (CDN拉流) │
└──────────┘ └──────────┘
┌──────────────────────────────────────────────────────────────────┐
│ 合流转推服务 │
│ (RTMP 合流转码 → 输出一路混合流 → 推送到 CDN 源站) │
└──────┬─────────────────────────────────────┬─────────────────────┘
│ │
↓ ↓
┌──────────┐ ┌──────────┐
│ 主播 A │ │ 连麦方 B │
│ (上行推流)│ │ (上行推流)│
└──────────┘ └──────────┘
↑ ↑
│ WebRTC / RTMP │ WebRTC / RTMP
│ (低延迟连麦通道) │ (低延迟连麦通道)
└────────────────┬────────────────────┘
│
┌──────────────┐
│ 连麦信令服务 │
│ (房间管理/信令)│
└──────────────┘数据流转路径:
- 观众端:从 CDN 拉取低延迟流(HTTP-FLV 或 WebRTC),延迟控制在 1-3 秒。
- 主播端:上行推流到 CDN 源站(RTMP 协议),同时通过 WebRTC 与连麦方建立低延迟通道。
- 连麦双方:通过 WebRTC 实现双向音视频通信,延迟控制在 300ms 以内;同时将各自的音视频流推送到合流转推服务。
- 合流转推服务:将多路输入流混合为一路输出流,再通过 RTMP 推送到 CDN 源站,供所有观众拉流观看。
合流转推方案
合流转推是指将多路主播的音视频流在服务端合并为一路流,再推送到 CDN。核心涉及两种模式:MCU(Multipoint Control Unit)和 SFU(Selective Forwarding Unit)。
MCU 与 SFU 模式对比
| 特性 | MCU(服务端混流) | SFU(服务端转发) |
|---|---|---|
| 工作原理 | 服务端解码多路流,合成为一路新流再编码输出 | 服务端不解码,直接转发各路流到订阅方 |
| 服务端负载 | 高(需解码+编码+合成) | 低(仅转发,无需编解码) |
| 延迟 | 较高(编解码引入延迟) | 低(转发延迟极小) |
| 带宽消耗 | 低(观众只需拉一路流) | 高(每路流独立传输) |
| 客户端复杂度 | 低(播放器播放单路流) | 高(需管理多路流渲染) |
| 灵活性 | 固定布局(服务端决定画面排列) | 灵活(客户端自行布局) |
| 适用场景 | 直播连麦 PK、多人连麦转 CDN | 视频会议、在线教育 |
混合架构方案:
生产环境中常采用 MCU + SFU 混合架构:连麦双方通过 SFU 实现低延迟互通,同时 MCU 服务将连麦画面合成为一路流推送到 CDN,观众端通过 CDN 拉流。
┌─────────────┐
│ 主播 A │
└──────┬──────┘
│ WebRTC
↓
┌───────────────────┐
│ SFU 转发服务 │ ←── 连麦方 B、C 通过 WebRTC 接入
│ (低延迟双向转发) │
└────────┬──────────┘
│
↓
┌───────────────────┐
│ MCU 合流转码 │ ←── 将 A/B/C 三路流合成为一路
│ (服务端混流) │
└────────┬──────────┘
│ RTMP 推流
↓
┌───────────────────┐
│ CDN 源站 │
└────────┬──────────┘
│ HTTP-FLV / HLS
↓
┌───────────────────┐
│ 所有观众 │
└───────────────────┘合流转码拓扑图
输入流 A ─→┐
├──→ 音频混音器 ─→ 音频编码器 ─→┐
输入流 B ─→┤ ├──→ MP4 muxer ─→ RTMP 推流
│ │
输入流 C ─→└──→ 视频拼接器 ─→ 视频编码器 ─→┘
│
↓
布局引擎 (Layout Engine)
├── 平铺模式 (Tile): 等分画面
├── 画中画模式 (PiP): 主播大窗 + 连麦方小窗
└── 自定义模式: 指定坐标和尺寸延时优化
平台连麦延迟目标
直播连麦业务的延迟目标分三个层级:
| 场景 | 目标延迟 | 协议 |
|---|---|---|
| 连麦双方互通 | < 300ms | WebRTC (UDP) |
| 连麦转推 CDN | 1-3s | RTMP / HTTP-FLV |
| 普通观众拉流 | 3-5s | HLS / HTTP-FLV |
影响延迟的关键因素
端到端延迟拆解:
采集 → 编码 → 网络传输 → 服务端处理 → 网络分发 → 解码 → 渲染
各环节优化策略:
采集:
- 使用系统原生采集接口,减少 buffer 环节
- 根据网络状态动态调整采集分辨率
编码:
- 使用硬件编码器(x264/MediaCodec/VideoToolbox)
- 开启编码器低延迟模式(-tune zerolatency)
- GOP 大小控制在 1-2s(帧间预测范围越小,延迟越低)
网络传输:
- 使用 UDP 协议替代 TCP(RTMP 基于 TCP)
- 实现 FEC(前向纠错)减少重传
- NACK + 关键帧请求结合
服务端处理:
- 避免转码(直接转发)
- 若需转码,使用低延迟编码参数
- 边缘节点就近处理
网络分发:
- 使用 WebRTC 推流到 CDN 边缘节点
- 多 CDN 动态切换
- 预热关键内容到边缘节点UDP 相比 RTMP 的优势
| 对比维度 | RTMP(TCP 协议) | WebRTC(UDP 协议) |
|---|---|---|
| 传输协议 | TCP | UDP + SRTP/SCTP |
| 连接建立 | 三次握手 | 非连接导向 |
| 拥塞控制 | TCP 内置(可能引入延迟) | GCC (Google Congestion Control) |
| 丢包重传 | 自动重传(可能阻塞后续数据) | NACK 选择性重传 |
| 头部开销 | 较大 | 较小 |
| 典型延迟 | 1-5s | 100-300ms |
| 适用场景 | 直播推流 | 实时通信 |
TCP 的可靠传输机制(ACK 确认、超时重传、拥塞窗口)在弱网环境下会引入"队头阻塞"问题——一个数据包丢失会导致后续所有数据包排队等待重传。UDP 没有这个限制,配合 FEC 和 NACK 可以在保证一定可靠性的同时大幅降低延迟。
WebRTC 连麦
WebRTC 信令交互
WebRTC 使用 SDP Offer/Answer 模型进行媒体协商,通过 ICE 框架进行网络穿透。信令本身不属于 WebRTC 标准,由应用层实现(可通过 WebSocket / HTTP / MQTT 等传输)。
主播 A 信令服务 连麦方 B
│ │ │
│ 1. Create Offer (SDP) │ │
│─────────────────────────────→│ │
│ │ 2. Forward Offer (SDP) │
│ │─────────────────────────────→│
│ │ │ 3. Create Answer (SDP)
│ │ │←────────────────────
│ │ 4. Forward Answer (SDP) │
│←─────────────────────────────│ │
│ │ │
│ 5. ICE Candidate │ │
│─────────────────────────────→│ │
│ │ 6. Forward ICE Candidate │
│ │─────────────────────────────→│
│ │ │
│←─────────────────────────────│ 7. ICE Candidate │
│ 8. Forward ICE Candidate │ │
│ │ │
│◄═══════════════════════════════════════════════════════════►│
│ 9. P2P 媒体通道建立 (SRTP/SCTP) │
│ │信令交互步骤:
- 主播创建 Offer:主播 A 创建 RTCPeerConnection,调用
createOffer()生成 SDP Offer,包含支持的音视频编解码器、网络信息等。 - 信令服务转发:服务端将 SDP Offer 转发给目标连麦方 B。
- 连麦方创建 Answer:连麦方 B 收到 Offer 后,调用
setRemoteDescription()设置远端 SDP,再调用createAnswer()生成 SDP Answer。 - 信令服务回传:服务端将 SDP Answer 转发回主播 A,主播 A 调用
setRemoteDescription()设置远端 SDP。 - ICE 候选者交换:双方各自收集 ICE Candidate(本机 IP、STUN 反射地址、TURN 中继地址),通过信令服务交换。
- P2P 连接建立:ICE 连通性检查通过后,建立 SRTP(加密媒体)和 SCTP(数据通道)连接。
信令消息格式
{
"type": "offer",
"roomId": "room_12345",
"from": "user_1001",
"to": "user_2002",
"sessionId": "sess_abc123",
"sdp": {
"type": "offer",
"sdp": "v=0\r\no=- 123456 2 IN IP4 0.0.0.0\r\n..."
}
}
{
"type": "ice_candidate",
"roomId": "room_12345",
"from": "user_1001",
"to": "user_2002",
"candidate": {
"candidate": "candidate:1 1 UDP 2122252543 192.168.1.1 54321 typ host",
"sdpMid": "0",
"sdpMLineIndex": 0
}
}信令服务端示例(Node.js + WebSocket)
// 信令服务器(简化实现)
const WebSocket = require('ws');
const wss = new WebSocket.Server({ port: 8080 });
const rooms = new Map();
wss.on('connection', (ws) => {
ws.on('message', (message) => {
const data = JSON.parse(message);
switch (data.type) {
case 'join':
// 加入房间
const room = rooms.get(data.roomId) || new Map();
room.set(data.userId, ws);
rooms.set(data.roomId, room);
ws.roomId = data.roomId;
ws.userId = data.userId;
// 通知房间其他成员
broadcast(room, {
type: 'user_joined',
userId: data.userId,
users: Array.from(room.keys())
}, ws);
break;
case 'offer':
case 'answer':
case 'ice_candidate':
// 转发到目标用户
const targetRoom = rooms.get(data.roomId);
if (targetRoom) {
const targetWs = targetRoom.get(data.to);
if (targetWs && targetWs.readyState === WebSocket.OPEN) {
targetWs.send(JSON.stringify(data));
}
}
break;
case 'leave':
leaveRoom(ws);
break;
}
});
ws.on('close', () => leaveRoom(ws));
});媒体协商
SDP 交换
SDP(Session Description Protocol)描述媒体会话的参数,包括媒体类型、编解码器、传输地址等。
v=0
o=- 1234567890 2 IN IP4 0.0.0.0
s=-
t=0 0
a=group:BUNDLE 0 1
a=msid-semantic: WMS
m=audio 9 UDP/TLS/RTP/SAVPF 111 103 104 9 0 8 106
c=IN IP4 0.0.0.0
a=rtcp:9 IN IP4 0.0.0.0
a=mid:0
a=msid:stream1 audio1
a=rtpmap:111 opus/48000/2
a=rtpmap:103 ISAC/16000
a=rtpmap:104 ISAC/32000
a=fmtp:111 minptime=10;useinbandfec=1
m=video 9 UDP/TLS/RTP/SAVPF 96 97 98 99 100
c=IN IP4 0.0.0.0
a=mid:1
a=msid:stream1 video1
a=rtpmap:96 VP8/90000
a=rtpmap:97 VP9/90000
a=rtpmap:98 H264/90000
a=fmtp:98 profile-level-id=42e01f;level-asymmetry-allowed=1
a=rtpmap:100 H265/90000关键字段说明:
| 字段 | 含义 |
|---|---|
m=audio/video | 媒体行,声明音频或视频流 |
a=rtpmap | 编解码器映射,格式为 payload_type codec/clockrate/channels |
a=fmtp | 编解码器特定参数 |
a=group:BUNDLE | 将多个媒体流复用在同一传输通道 |
a=ice-ufrag/ice-pwd | ICE 验证凭据 |
a=fingerprint | DTLS 指纹,用于 SRTP 加密密钥协商 |
编解码器协商
浏览器端 WebRTC 支持的编解码器优先级(按兼容性和性能排序):
| 编码类型 | 浏览器支持 | 码率范围 | 特点 |
|---|---|---|---|
| Opus(音频) | Chrome/Firefox/Safari | 6-510 kbps | 开源、低延迟、支持 FEC |
| VP8(视频) | 全平台 | 100-2000 kbps | WebRTC 基础编码器 |
| H.264(视频) | 全平台 | 100-2000 kbps | 硬件编码支持广泛 |
| VP9(视频) | Chrome/Firefox | 100-1500 kbps | 同等画质下码率更低 |
| AV1(视频) | Chrome 90+ | 50-1000 kbps | 最新开源编码器,计算开销大 |
编码器优先级调整示例:
// 优先使用 H.264(硬件编码器,低功耗)
const pc = new RTCPeerConnection(configuration);
const transceiver = pc.addTransceiver('video', { direction: 'sendrecv' });
// 设置编码器优先级:H.264 > VP8 > VP9
const codecPreferences = [];
const capabilities = RTCRtpReceiver.getCapabilities('video');
for (const codec of capabilities.codecs) {
if (codec.mimeType === 'video/H264') {
// H.264 设置为最高优先级
codecPreferences.unshift(codec);
} else if (codec.mimeType === 'video/VP8') {
codecPreferences.push(codec);
}
}
transceiver.setCodecPreferences(codecPreferences);服务端 SDP 协商策略
在合流转推场景中,服务端需要与 WebRTC 客户端进行 SDP 协商,将客户端发来的媒体流转接到合流转码模块。服务端通常使用 mediasoup、Janus 或 LiveKit 等 SFU 框架。
// 服务端接收客户端 Offer 并创建 Answer(mediasoup 示例)
async function handleOffer(roomId, userId, sdpOffer) {
const router = rooms.get(roomId).router;
// 创建 WebRTC 传输
const transport = await router.createWebRtcTransport({
listenIps: [{ ip: '0.0.0.0', announcedIp: 'public.ip.address' }],
enableUdp: true,
enableTcp: true,
preferUdp: true
});
// 设置客户端 SDP 并生成服务端 SDP Answer
const connectOptions = {
dtlsParameters: {
fingerprints: sdpOffer.fingerprints,
role: 'auto'
}
};
await transport.connect({ dtlsParameters: connectOptions.dtlsParameters });
// 创建 Producer(接收客户端媒体)
const producer = await transport.produce({
kind: 'video',
rtpParameters: sdpOffer.rtpParameters
});
return {
transport,
producer,
sdpAnswer: transport.sdpAnswer
};
}NAT 穿透
企业内网和运营商网络中,NAT(Network Address Translation)阻止了公网 IP 的直接访问。WebRTC 通过 ICE 框架实现 NAT 穿透。
STUN / TURN / ICE
┌──────────────────┐
│ STUN Server │
│ (NAT 地址探测) │
└────────┬─────────┘
│ STUN 请求/响应
│
NAT 设备 ──────────────┼──────────────────
│ │
↓ ↓
┌──────────┐ ┌──────────────┐
│ 主播端 │◄──────►│ 连麦方 B │
│ (内网IP) │ P2P │ (内网/公网) │
└────┬─────┘ └──────────────┘
│ ↑
│ ┌───────────────┘
↓ │
┌─────────┴──────┐
│ TURN Server │
│ (中继转发) │
└────────────────┘| 组件 | 全称 | 作用 | 部署位置 |
|---|---|---|---|
| STUN | Session Traversal Utilities for NAT | 探测客户端的公网 IP 和端口,返回映射地址 | 公网服务器 |
| TURN | Traversal Using Relays around NAT | 当 P2P 无法打通时,作为中继转发媒体数据 | 公网服务器(高带宽) |
| ICE | Interactive Connectivity Establishment | 收集候选地址并排序,通过连通性检查选择最佳路径 | 客户端库 |
ICE 候选者优先级排序:
- Host(本机 IP 地址)—— 延迟最低
- srflx(STUN 反射地址)—— 可获得公网 IP
- relay(TURN 中继地址)—— 兜底方案,延迟最高
ICE 连通性检查流程
主播 A ICE 流程 连麦方 B
│ │
│ 收集候选地址: │
│ host: 192.168.1.1:54321 │
│ srflx: 203.0.113.1:12345 │
│ relay: turn.example.com:3478 │
│ │
│ 交换候选地址 │
│──────────────────────────信令─────────────────────────────→│
│ │
│ 连通性检查 (STUN binding request): │
│ ──→ host -> srflx (直接) │
│ ──→ srflx -> srflx (P2P) │
│ ──→ relay -> relay (TURN 中继) │
│ │
│ 选择最佳路径: host <-> srflx (P2P 成功) │
│═══════════════════════════════════════════════════════════► │
│ 媒体传输 (SRTP/SCTP) │ICE 配置代码示例
const configuration = {
iceServers: [
{ urls: 'stun:stun.l.google.com:19302' },
{ urls: 'stun:stun1.l.google.com:19302' },
{
urls: 'turn:turn.example.com:3478',
username: 'webrtc_user',
credential: 'your_turn_credential'
},
{
urls: 'turn:turn.example.com:3478?transport=tcp',
username: 'webrtc_user',
credential: 'your_turn_credential'
}
],
iceCandidatePoolSize: 10
};
const pc = new RTCPeerConnection(configuration);TURN 服务自建 vs 云服务
| 对比项 | 自建 TURN(coturn) | 云 TURN(阿里云/腾讯云) |
|---|---|---|
| 部署成本 | 需独立服务器 | 按量付费 |
| 带宽成本 | 需自付服务器带宽 | 按使用量计费 |
| 节点覆盖 | 受限 | 全球多节点 |
| 维护成本 | 需自行维护 | 免运维 |
| 延迟 | 取决于服务器位置 | 智能调度 |
WebRTC ↔ RTMP 协议转换方案
直播场景中,主播端通常使用 RTMP 推流到 CDN,但连麦双方使用 WebRTC 实现低延迟互通。这就需要在服务端进行协议转换。
架构方案
┌──────────────────────────────────┐
│ 协议转换网关 │
│ │
WebRTC ──→ SRTP ──→│ RTP 解包 │
│ ↓ │
│ 音频: Opus → AAC │ ←── 转码(可选)
│ 视频: H264 → H264 (passthrough) │ ←── 不解码直接 remux
│ ↓ │
│ FLV muxer → RTMP 推流 ─────────→│──→ CDN
└──────────────────────────────────┘方案选项对比
| 方案 | 实现方式 | 优点 | 缺点 |
|---|---|---|---|
| FFmpeg 转推 | 通过 pipe 或 UDP 将 RTP 转 RTMP | 成熟稳定,社区支持好 | 延迟较高(需完整编解码) |
| SRS 服务器 | SRS 4.0+ 原生支持 WebRTC→RTMP 转换 | 延迟低,配置简单 | 功能定制有限 |
| Janus Gateway | 通过 VideoRoom 插件实现 | 功能丰富,可扩展 | 架构复杂 |
| 自定义网关 | 基于 librtmp + WebRTC 库自研 | 灵活可控 | 开发成本高 |
FFmpeg 转推示例
# 将 WebRTC 输出的 RTP 流转为 RTMP 推流
ffmpeg -f lavfi -i anullsrc -c:a aac -ar 44100 -ac 2 \
-f rtsp -rtsp_transport tcp -i rtsp://localhost:8554/live/stream \
-c:v copy -c:a copy \
-f flv rtmp://cdn-push.example.com/live/stream_mergedSRS 配置示例
SRS(Simple Realtime Server)是目前国内最流行的直播服务器之一,原生支持 WebRTC 推流并转 RTMP。
# srs.conf - WebRTC → RTMP 转换配置
listen 1935;
listen 8080;
http_api {
enabled on;
listen 1985;
}
rtc_server {
enabled on;
listen 8000;
# WebRTC over UDP
protocol udp;
}
vhost __defaultVhost__ {
rtc {
enabled on;
# 启用 WebRTC 到 RTMP 的自动转换
rtc_to_rtmp on;
}
http_remux {
enabled on;
mount [vhost]/[app]/[stream].flv;
}
hls {
enabled on;
hls_path ./objs/nginx/html/hls;
}
}客户端推流(WebRTC → SRS):
async function publishToSRS(streamUrl) {
const pc = new RTCPeerConnection({
iceServers: [{ urls: 'stun:stun.l.google.com:19302' }]
});
// 添加本地音视频轨道
const localStream = await navigator.mediaDevices.getUserMedia({
video: true, audio: true
});
localStream.getTracks().forEach(track => pc.addTrack(track, localStream));
// 创建 SDP Offer
const offer = await pc.createOffer();
await pc.setLocalDescription(offer);
// 通过 WHIP 协议推流到 SRS
const response = await fetch(streamUrl, {
method: 'POST',
headers: { 'Content-Type': 'application/sdp' },
body: offer.sdp
});
const answerSdp = await response.text();
await pc.setRemoteDescription({
type: 'answer',
sdp: answerSdp
});
}云直播服务
阿里云直播
阿里云视频直播(ApsaraVideo Live)提供一站式的直播解决方案,包括推流、转码、录制、截图、审核等功能。
推流域名 / 播流域名配置
阿里云直播使用推流域名和播流域名分离的架构,两者需要进行 CNAME 解析到阿里云直播平台。
域名配置:
推流域名: push.example.com
CNAME → push.example.com.w.alikunlun.com
用途: 主播端 RTMP 推流
播流域名: pull.example.com
CNAME → pull.example.com.w.alikunlun.com
用途: 观众端拉流(支持 HTTP-FLV / HLS / UDP)控制台或 API 创建域名:
# 阿里云 CLI 添加推流域名
aliyun live AddLiveDomain \
--DomainName push.example.com \
--LiveDomainType liveVideo \
--Region cn-shanghai \
--Scope domestic
# 添加播流域名
aliyun live AddLiveDomain \
--DomainName pull.example.com \
--LiveDomainType liveEdge \
--Region cn-shanghai \
--Scope domestic鉴权 URL
为防止盗播和恶意推流,阿里云直播支持 URL 鉴权(AuthKey 模式)。
鉴权 URL 格式:
RTMP 推流:
rtmp://push.example.com/live/{AppName}?auth_key={timestamp}-{rand}-{uid}-{md5hash}
HTTP-FLV 拉流:
http://pull.example.com/live/{AppName}.flv?auth_key={timestamp}-{rand}-{uid}-{md5hash}
HLS 拉流:
http://pull.example.com/live/{AppName}.m3u8?auth_key={timestamp}-{rand}-{uid}-{md5hash}鉴权 MD5 计算方式:
md5hash = md5(
"/{AppName}--{timestamp}-{rand}-{uid}-{private_key}"
)服务端生成鉴权 URL 示例(Java):
public String generateAuthUrl(String domain, String appName, String streamName,
String protocol, String privateKey, long expireSeconds) {
String uri = "/" + appName + "/" + streamName;
long timestamp = System.currentTimeMillis() / 1000 + expireSeconds;
String rand = "0";
String uid = "0";
// 拼接待签名字符串
String signStr = uri + "-" + timestamp + "-" + rand + "-" + uid + "-" + privateKey;
String md5 = DigestUtils.md5DigestAsHex(signStr.getBytes(StandardCharsets.UTF_8));
// 构造带有鉴权参数的 URL
String authParams = String.format("auth_key=%d-%s-%s-%s", timestamp, rand, uid, md5);
if ("rtmp".equals(protocol)) {
return String.format("rtmp://%s%s?%s", domain, uri, authParams);
} else if ("flv".equals(protocol)) {
return String.format("http://%s%s.flv?%s", domain, uri, authParams);
} else if ("hls".equals(protocol)) {
return String.format("http://%s%s.m3u8?%s", domain, uri, authParams);
}
throw new IllegalArgumentException("Unsupported protocol: " + protocol);
}录制、转码、截图 API
阿里云直播提供统一的 API 对直播流进行实时处理:
// 创建录制配置 - 录制为 HLS 存储到 OSS
CreateLiveStreamRecordIndexFilesRequest request = new CreateLiveStreamRecordIndexFilesRequest();
request.setDomainName("push.example.com");
request.setAppName("live");
request.setStreamName("stream_001");
request.setOssEndpoint("oss-cn-shanghai.aliyuncs.com");
request.setOssBucket("live-record-bucket");
request.setOssObject("record/{AppName}/{StreamName}/{Date}/{Hour}/{Minute}.m3u8");
// 创建实时转码配置
AddLiveStreamTranscodeRequest transcodeRequest = new AddLiveStreamTranscodeRequest();
transcodeRequest.setDomain("push.example.com");
transcodeRequest.setApp("live");
transcodeRequest.setTemplate("sd"); // 标清模板
// 创建截图配置
AddLiveStreamSnapshotConfigRequest snapshotRequest = new AddLiveStreamSnapshotConfigRequest();
snapshotRequest.setDomainName("push.example.com");
snapshotRequest.setAppName("live");
snapshotRequest.setTimeInterval(5); // 每 5 秒截一张
snapshotRequest.setOssEndpoint("oss-cn-shanghai.aliyuncs.com");
snapshotRequest.setOssBucket("live-snapshot-bucket");腾讯云直播
腾讯云提供 LVB(Live Video Broadcasting,云直播)和 CSS(Cloud Streaming Services)两个直播产品线。CSS 是 LVB 的升级版,功能更丰富。
LVB / CSS 服务
腾讯云直播域名配置:
推流域名: push.example.com
CNAME → push.example.com.livepush.myqcloud.com
播流域名: pull.example.com
CNAME → pull.example.com.liveplay.myqcloud.com推流地址生成:
RTMP 推流地址:
rtmp://push.example.com/live/{stream_id}?txSecret={md5}&txTime={hex_timestamp}
拉流地址:
HTTP-FLV: http://pull.example.com/live/{stream_id}.flv?txSecret={md5}&txTime={hex_timestamp}
HLS: http://pull.example.com/live/{stream_id}.m3u8?txSecret={md5}&txTime={hex_timestamp}
UDP: http://pull.example.com/live/{stream_id}.udp?txSecret={md5}&txTime={hex_timestamp}腾讯云推流鉴权实现:
public class TencentLiveAuth {
private static final String PLAY_KEY = "your_play_auth_key";
private static final String PUSH_KEY = "your_push_auth_key";
public static String generatePushUrl(String domain, String streamId, long expireHours) {
String txTime = Long.toHexString(System.currentTimeMillis() / 1000 + expireHours * 3600);
String txSecret = DigestUtils.md5DigestAsHex(
(PUSH_KEY + streamId + txTime).getBytes(StandardCharsets.UTF_8)
);
return String.format("rtmp://%s/live/%s?txSecret=%s&txTime=%s",
domain, streamId, txSecret, txTime.toUpperCase());
}
public static String generatePlayUrl(String domain, String streamId,
String protocol, long expireHours) {
String txTime = Long.toHexString(System.currentTimeMillis() / 1000 + expireHours * 3600);
String txSecret = DigestUtils.md5DigestAsHex(
(PLAY_KEY + streamId + txTime).getBytes(StandardCharsets.UTF_8)
);
String suffix = "";
if ("flv".equals(protocol)) suffix = ".flv";
else if ("m3u8".equals(protocol)) suffix = ".m3u8";
return String.format("%s://%s/live/%s%s?txSecret=%s&txTime=%s",
"flv".equals(protocol) ? "http" : "http",
domain, streamId, suffix, txSecret, txTime.toUpperCase());
}
}快直播 WebRTC 低延迟
腾讯云 快直播(LEB,Live Event Broadcasting)基于 WebRTC 协议实现毫秒级低延迟直播,适用于连麦、互动直播等场景。它与标准 WebRTC 连麦不同:快直播是单向低延迟分发(主播→观众),而非双向通信。
快直播架构:
主播 RTMP 推流 ──→ 腾讯云直播 ──→ WebRTC 转码 ──→ 观众端 (WebRTC 拉流)
处理中心 ↓
延迟 < 1s快直播拉流地址:
WebRTC 拉流地址(腾讯云):
webrtc://pull.example.com/live/{stream_id}?txSecret={md5}&txTime={hex_timestamp}前端播放快直播(WebRTC):
import { WebRTCPlayer } from 'trtc-web-player';
const player = new WebRTCPlayer({
url: 'webrtc://pull.example.com/live/stream_001',
autoplay: true,
muted: false,
controls: true
});
player.on('error', (error) => {
console.error('播放失败:', error);
// 降级到 HTTP-FLV
fallbackToFlvPlayer('stream_001');
});七牛云 / PaaS 直播
七牛云直播(Pili)提供 API 驱动的一站式直播云服务,核心功能封装为 RESTful API 和 SDK。
上行 / 下行 API 对比
| 功能 | 阿里云直播 API | 腾讯云直播 API | 七牛云 Pili API |
|---|---|---|---|
| 推流地址创建 | CreateLiveStream | CreateLiveStream | hub.createStream() |
| 推流地址管理 | DescribeLiveStreams | DescribeLiveStreams | hub.listStreams() |
| 禁播/断流 | ForbidLiveStream | SetLiveStreamForbidden | stream.disable() |
| 拉流地址 | 自动生成(按域名规则) | 自动生成(按域名规则) | publishUrl / rtmpPlayUrl / hlsPlayUrl / flvPlayUrl |
| 实时转码 | AddLiveStreamTranscode | CreateTranscodeTemplate | pipeline.transcoding |
| 直播录制 | CreateLiveRecordIndexFiles | CreateLiveRecord | recording.start() |
| 直播截图 | CreateLiveSnapshot | CreateLiveSnapshot | snapshot.start() |
管理 API 对比
| 功能 | 阿里云 | 腾讯云 | 七牛云 |
|---|---|---|---|
| SDK 语言 | Java/Python/Go/PHP/Node.js/C# | Java/Python/PHP/Node.js/Go | Go/Java/Python/PHP/Node.js/Ruby |
| 认证方式 | AccessKey + SecretKey | SecretId + SecretKey | AK + SK |
| API 网关 | live.aliyuncs.com | live.tencentcloudapi.com | pili.qiniu.com |
| 回调机制 | HTTP 回调 + MNS 消息 | HTTP 回调 + CMQ 消息 | HTTP 回调 |
| WebSocket 信令 | 阿里云 IM(连麦信令) | TRTC(连麦信令) | 需要自建信令服务 |
| 全球加速 | 海外节点较少 | 海外节点丰富(配合腾讯云 CDN) | 海外节点有限 |
回调 API 对比
| 回调类型 | 阿里云 | 腾讯云 | 七牛云 |
|---|---|---|---|
| 推流回调 | 推流开始/结束通知 | 推流开始/结束通知 | stream.connected / stream.disconnected |
| 录制回调 | 录制文件生成完成 | 录制文件生成完成 | recording.finished |
| 截图回调 | 截图文件生成完成 | 截图文件生成完成 | snapshot.generated |
| 转码回调 | 转码进度/完成 | 转码完成 | transcoding.finished |
| 审核回调 | 内容审核结果 | 内容审核结果 | 无原生审核回调 |
七牛云 Pili API 示例
// 七牛云 Pili SDK 示例
import com.qiniu.pili.*;
public class QiniuLiveService {
private final Hub hub;
public QiniuLiveService(String accessKey, String secretKey, String hubName) {
Client client = new Client(accessKey, secretKey);
this.hub = client.newHub(hubName);
}
// 创建直播流
public Stream createStream(String streamKey) throws Exception {
Hub.CreateStreamBuilder builder = hub.createStreamBuilder();
builder.setTitle(streamKey);
builder.setPublishKey("your_publish_key");
builder.setPublishSecurity("dynamic"); // 动态推流鉴权
return builder.build();
}
// 获取推流地址和播放地址
public Map<String, String> getStreamUrls(String streamId) {
Stream stream = hub.getStream(streamId);
Map<String, String> urls = new HashMap<>();
urls.put("rtmpPublishUrl", stream.rtmpPublishUrl());
urls.put("rtmpPlayUrl", stream.rtmpLiveUrls().get("ORIGIN"));
urls.put("hlsPlayUrl", stream.hlsLiveUrls().get("ORIGIN"));
urls.put("flvPlayUrl", stream.flvLiveUrls().get("ORIGIN"));
return urls;
}
// 开始录制
public void startRecording(String streamId, String bucket, long segmentDuration) {
// 七牛云通过转码 pipeline 触发录制
// 实际录制配置在服务端预设,推流后自动触发
}
}直播转码与录制
转码配置
直播转码将原始推流(通常为 1080p/高码率)实时转码为多个分辨率和码率档位,以适应不同网络条件和终端设备。
标清 / 高清 / 超清 模板
| 模板 | 分辨率 | 视频码率 | 音频码率 | 帧率 | 适用场景 |
|---|---|---|---|---|---|
| 超清 (LD) | 544 x 960 | 800-1000 kbps | 64 kbps | 24 fps | 手机 4G 网络 |
| 高清 (SD) | 720 x 1280 | 1500-2000 kbps | 96 kbps | 25 fps | Wi-Fi / 稳定 4G |
| 超清 (HD) | 1080 x 1920 | 3000-5000 kbps | 128 kbps | 30 fps | 宽带有线网络 |
| 蓝光 (2K) | 1440 x 2560 | 6000-10000 kbps | 192 kbps | 30 fps | 大屏电视/光纤 |
| 极清 (4K) | 2160 x 3840 | 15000-30000 kbps | 256 kbps | 30 fps | 赛事/演唱会直播 |
阿里云转码模板配置
{
"TemplateName": "hd_transcode",
"Video": {
"Codec": "H.264",
"Profile": "high",
"Width": 1280,
"Height": 720,
"Fps": 25,
"Bitrate": 2000,
"BitrateOpt": "recommend",
"Gop": "2s",
"RemoveDuplicatedFrame": true
},
"Audio": {
"Codec": "AAC",
"Profile": "aac_low",
"Samplerate": 44100,
"Bitrate": 96,
"Channels": 2
},
"TransConfig": {
"IsCheckAudioBitrate": false,
"IsCheckVideoBitrate": false
}
}自定义水印
直播水印支持图片叠加,可配置位置、大小和透明度。
// 阿里云添加直播水印
AddLiveStreamWatermarkRequest request = new AddLiveStreamWatermarkRequest();
request.setDomain("push.example.com");
request.setApp("live");
request.setStream("stream_001");
// 水印图片 OSS URL
request.setPictureUrl("https://your-bucket.oss-cn-hangzhou.aliyuncs.com/watermark.png");
// 水印位置:左上角 (X: 10px, Y: 10px)
request.setXPosition(0.02); // 横向位置(百分比)
request.setYPosition(0.02); // 纵向位置(百分比)
request.setWidth(0.15); // 水印宽度(百分比)
request.setHeight(0.0); // 0 表示按比例缩放
// 腾讯云水印模板
CreateWatermarkTemplateRequest request = new CreateWatermarkTemplateRequest();
request.setType("image");
request.setXPosition(10);
request.setYPosition(10);
request.setWidth(150);
request.setHeight(60);
request.setCoordinate("Absolute");DRM 加密
直播 DRM(Digital Rights Management)防止直播内容被非法录制和分发。
| 加密方案 | 原理 | 适用场景 |
|---|---|---|
| HLS AES-128 | 对 HLS 分片使用 AES-128 加密,密钥通过 HTTPS 分发 | HLS 直播 |
| HLS SAMPLE-AES | 仅对视频关键帧加密,兼容 FairPlay | iOS HLS 直播 |
| Widevine (CENC) | Google 的通用加密方案,支持自适应码率 | Android / Web |
| FairPlay | Apple 的 DRM 方案 | iOS/tvOS |
| 自定义加密 | 推流端加密 + 播放端解密 | 私有协议 |
HLS AES-128 加密配置:
# FFmpeg 加密推流
ffmpeg -i rtmp://push.example.com/live/stream_001 \
-c:v copy -c:a copy \
-hls_key_info_file key_info.txt \
-f hls -method PUT \
http://pull.example.com/live/stream_001.m3u8
# key_info.txt 内容格式
# key URI (播放器获取密钥的地址)
# key file path (本地密钥文件路径)
# IV (初始化向量,可选)
https://license.example.com/getkey?stream=stream_001
/path/to/encrypt.key
0123456789abcdef0123456789abcdef录制配置
直播录制将实时直播流存储为 HLS 或 MP4 文件,用于点播回放、内容审核或存档。
HLS / MP4 录制
| 录制格式 | 优点 | 缺点 | 适用场景 |
|---|---|---|---|
| HLS (m3u8 + ts) | 支持分段存储、边录边播、容错性好 | 文件碎片多 | 默认录制格式,长直播 |
| MP4 | 单文件、播放兼容性好 | 录制中断后文件不可用 | 短直播、重要场次 |
| FLV | 直播领域常用 | 播放器支持有限 | 点播回放 |
录制配置 API
// 阿里云创建录制配置
private void createRecordConfig(String domain, String appName, String streamName) {
CreateLiveStreamRecordIndexFilesRequest request = new CreateLiveStreamRecordIndexFilesRequest();
request.setDomainName(domain);
request.setAppName(appName);
request.setStreamName(streamName);
request.setOssEndpoint("oss-cn-shanghai.aliyuncs.com");
request.setOssBucket("live-record");
request.setOssObject("record/{AppName}/{StreamName}/{Date}/{Hour}/{Minute}_{Second}");
// 录制周期(单位:秒),超过此时长自动创建新文件
request.setRecordDuration(3600); // 1小时一个文件
// HLS 录制配置
request.setFormat("hls");
// 单个 ts 分片时长(秒)
request.setSliceDuration(30);
// 保留 ts 文件数量
request.setRetainInterval(30);
liveClient.createLiveStreamRecordIndexFiles(request);
}录制回调
录制完成后,云服务商会通过 HTTP POST 方式通知业务服务器。
阿里云录制回调内容示例:
{
"domain": "push.example.com",
"app": "live",
"stream": "stream_001",
"uri": "record/live/stream_001/2026/07/12/14_30_01.m3u8",
"duration": 3599.5,
"start_time": 1657614601,
"stop_time": 1657618201,
"file_size": 524288000,
"oss_endpoint": "oss-cn-shanghai.aliyuncs.com",
"oss_bucket": "live-record",
"oss_object": "record/live/stream_001/2026/07/12/14_30_01.m3u8",
"sign": "md5(domain+app+stream+uri+private_key)"
}回调签名验证:
public boolean verifyCallback(Map<String, String> params, String privateKey) {
String domain = params.get("domain");
String app = params.get("app");
String stream = params.get("stream");
String uri = params.get("uri");
String sign = params.get("sign");
String rawStr = domain + app + stream + uri + privateKey;
String expectedSign = DigestUtils.md5DigestAsHex(
rawStr.getBytes(StandardCharsets.UTF_8)
);
return expectedSign.equals(sign);
}存储到 OSS / COS
录制文件最终存储在云厂商的对象存储服务中:
存储目录结构(推荐):
live-record/
├── {AppName}/
│ ├── {StreamName}/
│ │ ├── 2026/
│ │ │ ├── 07/
│ │ │ │ ├── 12/
│ │ │ │ │ ├── 14_30_01.m3u8
│ │ │ │ │ ├── 14_30_01_000.ts
│ │ │ │ │ ├── 14_30_01_001.ts
│ │ │ │ │ └── ...
│ │ │ │ └── ...
│ │ │ └── ...
│ │ └── screenshots/
│ │ ├── 2026/07/12/14_30_05.jpg
│ │ └── ...
│ └── ...
└── transcoded/
└── {StreamName}/
└── sd/
└── hd/
└── ...录制文件生命周期管理:
# OSS 生命周期规则
生命周期:
- ID: live-record-archive
规则:
前缀: record/
过期时间:
- 30天: 转为低频存储 (Infrequent Access)
- 90天: 转为归档存储 (Archive)
- 365天: 自动删除截图与审核
直播截图
直播截图功能周期性抓取直播画面,用于封面展示、回放索引或内容审核。
// 阿里云配置直播截图
private void configureSnapshot(String domain, String appName) {
AddLiveStreamSnapshotConfigRequest request = new AddLiveStreamSnapshotConfigRequest();
request.setDomainName(domain);
request.setAppName(appName);
// 截图间隔(秒),最小 5 秒
request.setTimeInterval(10);
// 截图存储到 OSS
request.setOssEndpoint("oss-cn-shanghai.aliyuncs.com");
request.setOssBucket("live-snapshot");
// 截图文件名模板:{AppName}/{StreamName}/{Date}/{Hour}/{Minute}_{Second}.jpg
request.setOssObject("snapshot/{AppName}/{StreamName}/{Date}/{Hour}/{Minute}_{Second}.jpg");
liveClient.addLiveStreamSnapshotConfig(request);
}断帧检测
断帧(黑帧/花屏/绿屏)检测通过分析截图内容判断画面质量是否异常。
// 断帧检测伪代码
public class FrameQualityDetector {
public static DetectionResult analyzeFrame(BufferedImage frame) {
int width = frame.getWidth();
int height = frame.getHeight();
// 1. 检测黑帧:计算平均亮度
double avgBrightness = calculateAverageBrightness(frame);
if (avgBrightness < 10.0) {
return DetectionResult.BLACK_FRAME;
}
// 2. 检测绿屏:检测绿色通道占比
double greenRatio = calculateGreenChannelRatio(frame);
if (greenRatio > 0.9) {
return DetectionResult.GREEN_SCREEN;
}
// 3. 检测花屏:计算相邻像素差异
double blockNoise = calculateBlockNoise(frame);
if (blockNoise > 0.5) {
return DetectionResult.BLOCK_NOISE;
}
// 4. 检测冻帧:与上一帧比较
if (isFrameFrozen(frame, previousFrame)) {
return DetectionResult.FROZEN_FRAME;
}
return DetectionResult.NORMAL;
}
private static double calculateAverageBrightness(BufferedImage image) {
long totalBrightness = 0;
int pixels = image.getWidth() * image.getHeight();
for (int y = 0; y < image.getHeight(); y += 10) {
for (int x = 0; x < image.getWidth(); x += 10) {
int rgb = image.getRGB(x, y);
int r = (rgb >> 16) & 0xFF;
int g = (rgb >> 8) & 0xFF;
int b = rgb & 0xFF;
totalBrightness += (r + g + b) / 3;
}
}
return (totalBrightness * 100.0) / (pixels / 100);
}
}内容审核 API 对接
内容审核可以对直播截图进行涉黄、涉政、暴恐、广告等违规内容识别。
// 阿里云内容审核对接示例
public class LiveContentModeration {
private static final String ACCESS_KEY_ID = "your_access_key";
private static final String ACCESS_KEY_SECRET = "your_secret";
// 审核直播截图
public ModerationResult moderateImage(String imageUrl) {
// 调用阿里云内容安全服务(Green)
Client client = createGreenClient();
ScanImageRequest request = new ScanImageRequest();
// 设置审核场景
request.setScenes(Arrays.asList(
"porn", // 涉黄
"terrorism", // 暴恐
"ad", // 广告
"political" // 敏感人物/政治违规
));
// 添加待审核图片
ScanImageRequest.Task task = new ScanImageRequest.Task();
task.setImageURL(imageUrl);
task.setDataId(UUID.randomUUID().toString());
request.setTasks(Arrays.asList(task));
try {
ScanImageResponse response = client.scanImage(request);
return parseModerationResult(response);
} catch (Exception e) {
throw new RuntimeException("内容审核调用失败", e);
}
}
private ModerationResult parseModerationResult(ScanImageResponse response) {
if (response.getCode() != 200) {
return ModerationResult.error(response.getCode(), response.getMsg());
}
ScanImageResponse.Data data = response.getData().get(0);
List<ScanImageResponse.Result> results = data.getResults();
for (ScanImageResponse.Result result : results) {
if ("pass".equals(result.getSuggestion())) continue;
// 违规或疑似违规
String label = result.getLabel();
double confidence = result.getRate();
String suggestion = result.getSuggestion(); // review / block
return ModerationResult.violation(label, suggestion, confidence);
}
return ModerationResult.pass();
}
// 腾讯云内容审核
public ModerationResult tencentModerate(String imageUrl) {
// 调用腾讯云内容安全(CMS)
Credential cred = new Credential(SECRET_ID, SECRET_KEY);
ImageModerationClient client = new ImageModerationClient(cred, "ap-shanghai");
ImageModerationRequest request = new ImageModerationRequest();
request.setImageUrl(imageUrl);
request.setBizType("live_review");
ImageModerationResponse response = client.ImageModeration(request);
return new ModerationResult()
.setPass(response.getSuggestion().equals("PASS"))
.setLabel(response.getLabel())
.setScore(response.getScore());
}
}审核流程自动化:
推流中 ──→ 周期性截图 ──→ OSS/COS 存储 ──→ 内容审核 API ──→ 审核结果回调
│
┌────────────────┤
↓ ↓
审核通过 审核违规(疑似)
│ │
↓ ├── 人工复核 → 确认违规
继续推流 │ ↓
│ 自动断流/禁播
│
└── 标记待审 → 转人工队列互动场景设计
连麦 PK
连麦 PK 是直播互动中最常见的场景之一,两个主播进行实时对决。
PK 流程
主播 A 服务端 主播 B
│ │ │
│ 1. 发起 PK 请求 │ │
│ ───────────────────────────────→ │ │
│ │ 2. 推送 PK 邀请 │
│ │ ────────────────────────────→ │
│ │ │
│ │ 3. 接受 / 拒绝 PK │
│ │ ←──────────────────────────── │
│ 4. PK 开始通知 │ │
│ ←───────────────────────────────│ │
│ │ 5. 发起合流转码 │
│ │ (MCU 合成 A+B 画面到一路流) │
│ │ │
│ 6. 倒计时开始 (30s/60s/120s) │ │
│ ←───────────────────────────────│─────────────────────────────→ │
│ │ │
│ 7. 观众送礼/投票 │
│ ˑˑˑˑˑˑˑˑˑˑˑˑˑˑˑˑˑˑˑˑˑˑˑˑˑˑˑˑˑˑˑˑˑˑˑˑˑˑˑˑˑˑˑˑˑˑˑˑˑˑˑˑˑˑˑˑˑˑˑˑˑˑˑ │
│ │ │
│ 8. PK 结束 │ │
│ ←───────────────────────────────│─────────────────────────────→ │
│ │ 9. 停止合流转码 │
│ │ │
│ 10. 展示 PK 结果 │ │
│ ←───────────────────────────────│─────────────────────────────→ │
│ │ │
│ 11. 关闭连麦通道 │ │
│ ←───────────────────────────────│─────────────────────────────→ │PK 状态机
+----------+
| 空闲态 |
+----+-----+
|
发起 PK (发起者)
|
v
+----------+
| 等待中 | ←── 等待对方接受
+----+-----+
|
+----------+----------+
| |
接受/自动接受 拒绝/超时
| |
v v
+----------+ +----------+
| PK 进行 | | 已拒绝 |
+----+-----+ +----------+
|
倒计时结束/手动结束
|
v
+----------+
| PK 结束 | ──→ 展示结果
+----+-----+
|
关闭连麦
|
v
+----------+
| 空闲态 |
+----------+PK 合流布局
PK 倒计时画面(5-3-2-1):
┌──────────────────┐ ┌──────────────────┐
│ │ │ │
│ 主播 A 画面 │ │ 主播 B 画面 │
│ │ │ │
│ 分数: 12500 │ │ 分数: 9800 │
└──────────────────┘ └──────────────────┘
┌──────────────────────────────────────────┐
│ 倒计时: 00:32 │
└──────────────────────────────────────────┘PK 接口设计
@RestController
@RequestMapping("/api/pk")
public class PKController {
// 发起 PK
@PostMapping("/start")
public ResponseEntity<PKResult> startPK(@RequestBody PKStartRequest request) {
// 参数校验
String fromRoomId = request.getFromRoomId();
String toRoomId = request.getToRoomId();
int duration = request.getDuration(); // PK 时长(秒)
// 检查对方直播间状态
LiveRoom targetRoom = liveRoomService.getRoom(toRoomId);
if (targetRoom.getStatus() != RoomStatus.LIVE) {
return ResponseEntity.badRequest().body(PKResult.error("对方不在直播中"));
}
if (targetRoom.isInPK()) {
return ResponseEntity.badRequest().body(PKResult.error("对方正在进行 PK"));
}
// 发送 PK 邀请给主播 B
pkService.sendInvite(fromRoomId, toRoomId, duration);
return ResponseEntity.ok(PKResult.success("PK 邀请已发送"));
}
// 接受 PK
@PostMapping("/accept")
public ResponseEntity<PKResult> acceptPK(@RequestBody PKAcceptRequest request) {
String roomId = request.getRoomId();
String inviteId = request.getInviteId();
// 验证邀请有效性
PKInvite invite = pkService.getInvite(inviteId);
if (invite == null || invite.isExpired()) {
return ResponseEntity.badRequest().body(PKResult.error("邀请已过期"));
}
// 创建 PK 会话
PKSession session = pkService.createSession(invite, request.getDuration());
// 发起合流转码
mcuService.startPKMix(session.getStreamA(), session.getStreamB(),
MixLayout.SPLIT_HALF);
return ResponseEntity.ok(PKResult.success(session));
}
// PK 结束
@PostMapping("/end")
public ResponseEntity<PKResult> endPK(@RequestBody PKSessionRequest request) {
String sessionId = request.getSessionId();
// 计算 PK 结果
PKResultData result = pkService.calculateResult(sessionId);
// 停止合流转码
mcuService.stopPKMix(sessionId);
// 广播结果
pkService.broadcastResult(sessionId, result);
return ResponseEntity.ok(PKResult.success(result));
}
// 查询 PK 状态
@GetMapping("/status/{sessionId}")
public ResponseEntity<PKSession> getPKStatus(@PathVariable String sessionId) {
PKSession session = pkService.getSession(sessionId);
return ResponseEntity.ok(session);
}
}
// PK 相关数据模型
@Data
public class PKStartRequest {
private String fromRoomId;
private String toRoomId;
private int duration; // PK 时长(秒)
}
@Data
public class PKAcceptRequest {
private String roomId;
private String inviteId;
private int duration;
}
@Data
public class PKSession {
private String sessionId;
private String streamA;
private String streamB;
private String roomA;
private String roomB;
private int duration;
private long startTime;
private PKStatus status;
}
public enum PKStatus {
WAITING, // 等待开始
COUNTDOWN, // 倒计时
IN_PROGRESS, // PK 进行中
ENDED // 已结束
}多主播直播间
多主播直播间允许多个主播同时在线互动,观众可以看到所有主播的画面。
架构设计
┌────────────────────────┐
│ MCU 合流转码服务 │
│ │
│ ┌── 主播 A ──┐ │
│ │ 主播 B │ ← Tile Layout
│ │ 主播 C │ │
│ │ 主播 D │ │
│ └────────────┘ │
└───────────┬────────────┘
│ RTMP 推流到 CDN
↓
┌────────────────────────┐
│ 观众(CDN 拉流) │
└────────────────────────┘
┌────────────────────────┐
│ SFU 转发服务 │
│ (主播间低延迟互通) │
│ │
│ 主播 A ◄────► 主播 B │
│ ▲ ▲ │
│ │ │ │
│ 主播 C ◄────► 主播 D │
└────────────────────────┘多主播合流布局
平铺布局 (2x2): 平铺布局 (1+3):
┌──────────┬──────────┐ ┌──────────────────┬────┐
│ │ │ │ │ B │
│ 主播 A │ 主播 B │ │ 主播 A ├────┤
│ │ │ │ (主持人/主窗) │ C │
├──────────┼──────────┤ │ ├────┤
│ │ │ │ │ D │
│ 主播 C │ 主播 D │ └──────────────────┴────┘
│ │ │
└──────────┴──────────┘多主播房间管理
@Service
public class MultiHostRoomService {
private final Map<String, MultiHostRoom> rooms = new ConcurrentHashMap<>();
// 创建多主播房间
public MultiHostRoom createRoom(String creatorId, String roomName, int maxHosts) {
MultiHostRoom room = new MultiHostRoom();
room.setRoomId(UUID.randomUUID().toString());
room.setRoomName(roomName);
room.setMaxHosts(maxHosts);
room.setHosts(new CopyOnWriteArrayList<>());
room.setStatus(RoomStatus.WAITING);
// 将创建者设为主播
room.getHosts().add(new RoomHost(creatorId, HostRole.ANCHOR, Instant.now()));
// 创建合流任务
room.setMixTaskId(mcuService.createMixTask(maxHosts));
rooms.put(room.getRoomId(), room);
return room;
}
// 邀请主播加入
public boolean inviteHost(String roomId, String inviterId, String inviteeId) {
MultiHostRoom room = rooms.get(roomId);
if (room == null) return false;
if (room.getHosts().size() >= room.getMaxHosts()) return false;
// 发送邀请通知
notificationService.sendInvite(inviteeId, roomId, inviterId);
return true;
}
// 主播加入房间
public RoomHost joinAsHost(String roomId, String userId) {
MultiHostRoom room = rooms.get(roomId);
RoomHost host = new RoomHost(userId, HostRole.GUEST_HOST, Instant.now());
room.getHosts().add(host);
// 将主播的推流加入 MCU 合流
mcuService.addToMix(room.getMixTaskId(), userId, getNextLayoutPosition(room));
// 通知房间内其他主播
broadcastHostChange(room, "host_joined", host);
return host;
}
// 主播离开房间
public void leaveHost(String roomId, String userId) {
MultiHostRoom room = rooms.get(roomId);
room.getHosts().removeIf(h -> h.getUserId().equals(userId));
// 从 MCU 合流中移除
mcuService.removeFromMix(room.getMixTaskId(), userId);
// 通知房间内其他主播
broadcastHostChange(room, "host_left", userId);
// 如果没有主播了,关闭房间
if (room.getHosts().isEmpty()) {
closeRoom(roomId);
}
}
// 关闭房间
public void closeRoom(String roomId) {
MultiHostRoom room = rooms.remove(roomId);
if (room != null) {
mcuService.stopMix(room.getMixTaskId());
notificationService.broadcastToRoom(roomId, "room_closed", null);
}
}
// 获取下一个布局位置
private int getNextLayoutPosition(MultiHostRoom room) {
return room.getHosts().size(); // 0-based 位置索引
}
// 广播主播变更
private void broadcastHostChange(MultiHostRoom room, String event, Object data) {
notificationService.broadcastToRoom(room.getRoomId(), event, data);
}
}
// 多主播房间数据模型
@Data
public class MultiHostRoom {
private String roomId;
private String roomName;
private int maxHosts;
private List<RoomHost> hosts;
private String mixTaskId;
private RoomStatus status;
private Instant createdAt;
}
@Data
@AllArgsConstructor
public class RoomHost {
private String userId;
private HostRole role; // ANCHOR(房主)或 GUEST_HOST(嘉宾主播)
private Instant joinedAt;
}
public enum HostRole {
ANCHOR, // 房主
GUEST_HOST // 嘉宾主播
}
public enum RoomStatus {
WAITING, // 等待开始
LIVE, // 直播中
ENDED // 已结束
}
### 观众连线申请
观众连线(观众上麦)允许观众申请与主播进行实时音视频互动,常见于直播连麦问答、语音互动、才艺展示等场景。
#### 连线流程
```text
观众 服务端 主播
│ │ │
│ 1. 发起连麦申请 │ │
│ ─────────────────────────────→ │ │
│ │ 2. 推送连麦申请通知 │
│ │ ────────────────────────────→ │
│ │ │
│ │ 3. 审核申请 │
│ │ ├── 同意 │
│ │ └── 拒绝 │
│ │ │
│ 4. 连麦结果通知 │ │
│ ←─────────────────────────────│ │
│ │ │
│ 5. 建立 WebRTC 连接 │ │
│═══════════════════════════════════════════════════════════════►│
│ 连麦通道(双向音视频 / 仅音频) │
│ │
│ 6. 观众画面加入 MCU 合流 │ │
│ │──────────────────────────────→│
│ │ │
│ 7. 主播 / 观众主动下麦 │ │
│ ←─────────────────────────────│──────────────────────────────→ │
│ │ 8. 停止合流、断开 WebRTC │连线状态机
+----------+
| 空闲态 |
+----+-----+
|
申请连麦 (观众)
|
v
+----------+
| 申请中 | ←── 等待主播审核
+----+-----+
|
+----------+----------+
| |
主播同意 主播拒绝/超时
| |
v v
+----------+ +----------+
| 连麦中 | | 已拒绝 |
+----+-----+ +----------+
|
主动下麦/主播踢下麦
|
v
+----------+
| 空闲态 |
+----------+连麦申请接口设计
@RestController
@RequestMapping("/api/live/connect")
public class LiveConnectController {
@Autowired
private LiveConnectService connectService;
// 观众申请连麦
@PostMapping("/apply")
public ResponseEntity<ConnectResult> applyForConnect(
@RequestBody ConnectApplyRequest request) {
String roomId = request.getRoomId();
String userId = request.getUserId();
boolean audioOnly = request.isAudioOnly(); // 仅音频连麦
// 校验房间状态
LiveRoom room = liveRoomService.getRoom(roomId);
if (room == null || room.getStatus() != RoomStatus.LIVE) {
return ResponseEntity.badRequest()
.body(ConnectResult.error("直播间不在直播状态"));
}
// 检查连麦人数限制
if (connectService.getConnectedCount(roomId) >= room.getMaxConnections()) {
return ResponseEntity.badRequest()
.body(ConnectResult.error("连麦人数已达上限"));
}
// 创建连麦申请
ConnectApplication application = connectService.createApplication(
roomId, userId, audioOnly);
// 通知主播审核
notificationService.notifyAnchor(room.getAnchorId(),
"connect_apply", application);
return ResponseEntity.ok(ConnectResult.success(application));
}
// 主播审核连麦申请
@PostMapping("/review")
public ResponseEntity<ConnectResult> reviewApplication(
@RequestBody ReviewRequest request) {
String applicationId = request.getApplicationId();
boolean approved = request.isApproved();
ConnectApplication application = connectService.getApplication(applicationId);
if (application == null || application.getStatus() != ApplyStatus.PENDING) {
return ResponseEntity.badRequest()
.body(ConnectResult.error("申请不存在或已处理"));
}
if (approved) {
// 同意连麦
connectService.approveApplication(application);
// 建立 WebRTC 连接(服务端发起或通过信令告知双方)
connectService.establishWebRTCConnection(application);
// 将观众画面加入合流
if (!application.isAudioOnly()) {
mcuService.addToMix(
application.getRoomId(),
application.getUserId(),
MixLayout.PIP_SMALL // 画中画模式,小窗显示观众
);
}
// 通知观众准备连麦
notificationService.notifyUser(application.getUserId(),
"connect_approved", application);
} else {
// 拒绝连麦
connectService.rejectApplication(application);
notificationService.notifyUser(application.getUserId(),
"connect_rejected", application);
}
return ResponseEntity.ok(ConnectResult.success());
}
// 主动下麦
@PostMapping("/disconnect")
public ResponseEntity<ConnectResult> disconnect(
@RequestBody DisconnectRequest request) {
String roomId = request.getRoomId();
String userId = request.getUserId();
// 从合流中移除
mcuService.removeFromMix(roomId, userId);
// 断开 WebRTC 连接
connectService.disconnect(roomId, userId);
// 通知双方
LiveRoom room = liveRoomService.getRoom(roomId);
notificationService.notifyUser(room.getAnchorId(),
"audience_disconnected", userId);
notificationService.notifyUser(userId,
"you_are_disconnected", null);
return ResponseEntity.ok(ConnectResult.success());
}
// 查询连麦列表
@GetMapping("/list/{roomId}")
public ResponseEntity<List<ConnectedUser>> getConnectedUsers(
@PathVariable String roomId) {
List<ConnectedUser> users = connectService.getConnectedUsers(roomId);
return ResponseEntity.ok(users);
}
}连麦申请数据模型
@Data
public class ConnectApplyRequest {
private String roomId;
private String userId;
private boolean audioOnly; // true: 仅音频连麦, false: 音视频连麦
}
@Data
public class ConnectApplication {
private String applicationId;
private String roomId;
private String userId;
private String userName;
private boolean audioOnly;
private ApplyStatus status;
private Instant createdAt;
private Instant reviewedAt;
}
public enum ApplyStatus {
PENDING, // 待审核
APPROVED, // 已同意
REJECTED, // 已拒绝
EXPIRED // 已过期
}
@Data
public class ConnectedUser {
private String userId;
private String userName;
private String avatar;
private boolean audioOnly;
private Instant connectedAt;
private String rtcStatus; // connected / reconnecting / disconnected
}连麦超时与异常处理
@Service
public class LiveConnectService {
// 连麦申请超时时间(秒)
private static final long APPLY_TIMEOUT_SECONDS = 30;
// 连麦申请自动清理间隔(秒)
private static final long CLEANUP_INTERVAL_SECONDS = 60;
@Scheduled(fixedDelay = CLEANUP_INTERVAL_SECONDS * 1000)
public void cleanupExpiredApplications() {
Instant threshold = Instant.now()
.minus(APPLY_TIMEOUT_SECONDS, ChronoUnit.SECONDS);
applications.values().removeIf(app -> {
if (app.getStatus() == ApplyStatus.PENDING
&& app.getCreatedAt().isBefore(threshold)) {
// 标记为过期
app.setStatus(ApplyStatus.EXPIRED);
// 通知用户申请超时
notificationService.notifyUser(app.getUserId(),
"connect_expired", app);
return true;
}
return false;
});
}
// 检测连麦方连接状态
@Scheduled(fixedDelay = 10000)
public void checkConnectionHealth() {
connectedUsers.forEach((sessionId, session) -> {
if (!rtcService.checkConnection(sessionId)) {
log.warn("连麦连接断开: sessionId={}", sessionId);
// 尝试重连
boolean reconnected = rtcService.tryReconnect(sessionId);
if (!reconnected) {
// 重连失败,强制下麦
forceDisconnect(sessionId, "连接超时");
}
}
});
}
// 强制下麦
private void forceDisconnect(String sessionId, String reason) {
ConnectedSession session = connectedUsers.get(sessionId);
if (session == null) return;
mcuService.removeFromMix(session.getRoomId(), session.getUserId());
rtcService.closeConnection(sessionId);
connectedUsers.remove(sessionId);
// 通知双方
notificationService.notifyAnchor(session.getRoomId(),
"audience_forced_disconnect",
Map.of("userId", session.getUserId(), "reason", reason));
notificationService.notifyUser(session.getUserId(),
"forced_disconnect", Map.of("reason", reason));
}
}连麦布局与体验优化
| 场景 | 布局模式 | 说明 |
|---|---|---|
| 单人连麦 | 画中画 (PiP) | 主播全屏,观众小窗(右下角) |
| 多人连麦 | 平铺 (Tile) | 所有画面等分排列 |
| 仅音频连麦 | 头像 + 声浪 | 显示观众头像和音频波动画 |
| 问答连麦 | 主播全屏 + 观众半屏 | 教育/咨询场景 |
画中画布局 (1 个观众连麦):
┌──────────────────────────────────────┐
│ │
│ │
│ 主播全屏 │
│ │
│ ┌──────────┐│
│ │ 观众画面 ││
│ │ (小窗) ││
│ └──────────┘│
└──────────────────────────────────────┘
平铺布局 (3 个观众连麦):
┌────────────┬────────────┬────────────┐
│ │ │ │
│ 主播 │ 观众 A │ 观众 B │
│ │ │ │
├────────────┴────────────┴────────────┤
│ 观众 C │
└──────────────────────────────────────┘连麦安全性设计
连麦安全策略:
1. 权限控制:
- 仅认证用户可申请连麦
- 主播可设置连麦门槛(关注时长、等级、付费)
- 主播可拉黑/禁言特定用户
2. 风控策略:
- 单用户连麦频率限制(N 次/小时)
- 自动审核用户画像(新号、异常行为检测)
- 连麦内容实时截图审核
3. 降级方案:
- 观众端 WebRTC 连接失败 → 降级为仅音频
- 服务端合流异常 → 断开连麦,恢复单主播画面
- TURN 带宽不足 → 限制连麦人数