用Node.js做直播系统是很多前端和全栈开发者都会遇到的需求场景。相比传统的C++流媒体方案,Node.js虽然不适合直接处理音视频编解码这类CPU密集型工作,但凭借出色的事件驱动和IO处理能力,它非常适合承担流管理、信令控制、弹幕消息转发这类工作。这篇文章把直播系统拆成推拉流和弹幕互动两个核心模块,逐个讲清楚实现思路和关键代码。

直播系统的整体架构设计
一个最小可用的直播系统由三部分组成:主播端推流、流媒体服务器中转分发、观众端拉流播放。流媒体服务器推荐使用Nginx配合nginx-rtmp-module,或者SRS(Simple RTMP Server),Node.js负责业务逻辑层,包括鉴权、房间管理、推流地址生成、弹幕消息分发等。
这里要厘清一个概念:Node.js本身并不解码视频流。音视频数据从主播的OBS等推流软件通过RTMP协议发往流媒体服务器,服务器把数据切片转封装,观众端再用HLS或HTTP-FLV协议拉取。Node.js在整个链路里扮演的是调度员角色,它决定谁能推流、哪个流是合法的、观众该去哪里拉流,这个定位理解错了后面架构就容易走偏。
推拉流的完整流程是这样的:主播先向Node.js业务服务请求推流地址,Node.js校验身份后生成带签名的RTMP地址返回;主播用OBS向该地址推流;观众打开页面时,Node.js返回对应的播放地址;同时观众端的浏览器通过WebSocket连上弹幕服务,形成双向消息通道。整条链路中Node.js只处理信令和文本消息,不碰音视频数据,这样才能保证服务稳定性。
推拉流的实现细节与鉴权控制
推流鉴权是必须做的,否则任何人拿到服务器地址都能推流,直播内容不可控。常见做法是Node.js生成一个带过期时间戳和签名的推流密钥,nginx-rtmp-module通过on_publish回调把这个密钥发回Node.js校验,校验通过才允许推流建立。
下面是Node.js侧生成推流地址和校验的核心代码:
const crypto = require('crypto');
// 生成带签名的推流地址
function generatePushUrl(streamKey, expireSeconds) {
const expire = Math.floor(Date.now() / 1000) + expireSeconds;
const secret = 'your-rtmp-secret';
const sign = crypto
.createHash('md5')
.update(`${streamKey}-${expire}-${secret}`)
.digest('hex');
return `rtmp://live.ippipp.com/live/${streamKey}?expire=${expire}&sign=${sign}`;
}
// nginx-rtmp-module 的 on_publish 回调处理
const express = require('express');
const app = express();
app.use(express.urlencoded({ extended: false }));
app.post('/rtmp/auth', (req, res) => {
const { name, expire, sign } = req.body;
const secret = 'your-rtmp-secret';
const expect = crypto
.createHash('md5')
.update(`${name}-${expire}-${secret}`)
.digest('hex');
if (String(expire) > String(Math.floor(Date.now() / 1000)) && sign === expect) {
return res.status(200).send('ok');
}
res.status(403).send('forbidden');
});
app.listen(3000);拉流侧的协议选择需要权衡延迟和兼容性。HLS兼容性最好,手机浏览器都能直接播,但延迟通常在10秒以上;HTTP-FLV延迟可以压到2到3秒,配合flv.js在浏览器播放效果很好,是目前国内直播平台的主流方案。nginx-rtmp-module开启这两个协议只需要简单配置,HLS会自动切片成ts文件并生成m3u8索引,HTTP-FLV则通过http-live模块的stat地址暴露。
观众端播放用flv.js的代码大致如下,如果浏览器不支持MSE,再降级到HLS播放:
<video id="player" controls autoplay></video>
<script src="flv.min.js"></script>
<script>
if (flvjs.isSupported()) {
const player = flvjs.createPlayer({
type: 'flv',
url: 'http://live.ippipp.com/live/room001.flv'
});
player.attachMediaElement(document.getElementById('player'));
player.load();
player.play();
}
</script>基于WebSocket的弹幕系统实现
弹幕是直播互动的灵魂,实现难点不在于把消息发出去,而在于高并发下怎么保证消息不丢、不乱、不把服务器打爆。Node.js的ws库是构建弹幕服务的首选,单进程就能维持数万条WebSocket连接,配合cluster模式可以榨满多核。
先看一个基础版本的弹幕服务,包含房间隔离和广播逻辑:
const WebSocket = require('ws');
const wss = new WebSocket.Server({ port: 8080 });
// 房间模型:roomId -> Set of clients
const rooms = new Map();
function joinRoom(roomId, ws) {
if (!rooms.has(roomId)) rooms.set(roomId, new Set());
rooms.get(roomId).add(ws);
ws.roomId = roomId;
}
function broadcast(roomId, message) {
const room = rooms.get(roomId);
if (!room) return;
const data = JSON.stringify(message);
for (const client of room) {
if (client.readyState === WebSocket.OPEN) {
client.send(data);
}
}
}
wss.on('connection', (ws, req) => {
const roomId = new URL(req.url, 'http://localhost').searchParams.get('room');
joinRoom(roomId, ws);
ws.on('message', (raw) => {
try {
const msg = JSON.parse(raw);
if (msg.type === 'danmaku') {
// 服务端补充时间戳,用于历史弹幕对齐
broadcast(roomId, {
type: 'danmaku',
user: msg.user,
text: msg.text.slice(0, 50), // 限制长度,防止刷屏
ts: Date.now()
});
}
} catch (e) {
// 非法消息直接丢弃
}
});
ws.on('close', () => {
const room = rooms.get(ws.roomId);
if (room) {
room.delete(ws);
if (room.size === 0) rooms.delete(ws.roomId);
}
});
});这个基础版本能跑,但有几个生产环境必须补上的点。第一是消息合并,当一个房间同时有几万人发弹幕时,逐条广播会造成大量小包发送,正确做法是按200毫秒左右的窗口把消息攒成数组批量推送,网络包数量能降一个数量级。第二是敏感词过滤,弹幕文本必须在服务端过一遍过滤逻辑再广播,不能信任客户端。第三是历史弹幕回放,把弹幕按时间戳存进Redis的有序集合,观众进入直播间时拉取最近一段时间的弹幕回放,体验会好很多。
多进程部署时还有一个坑要注意:cluster模式下每个进程维护自己的房间集合,观众A连在进程1,观众B连在进程2,弹幕就串不起来。解决方案是用Redis的发布订阅做进程间消息总线,收到弹幕后publish到频道,每个进程subscribe后各自广播给自己持有的连接。
const Redis = require('ioredis');
const sub = new Redis();
const pub = new Redis();
// 进程收到弹幕后先发到Redis频道
ws.on('message', (raw) => {
const msg = JSON.parse(raw);
pub.publish(`danmaku:${ws.roomId}`, JSON.stringify({
type: 'danmaku', user: msg.user, text: msg.text, ts: Date.now()
}));
});
// 每个进程订阅频道,收到消息后广播给本进程持有的连接
sub.psubscribe('danmaku:*');
sub.on('pmessage', (pattern, channel, message) => {
const roomId = channel.split(':')[1];
broadcast(roomId, JSON.parse(message));
});性能优化与横向扩容
弹幕系统真正的压力来自广播的扇出效应。一条弹幕在一个十万人的房间会产生十万次send调用,这中间序列化是重复计算的,正确做法是把消息序列化一次,把Buffer分发给所有连接,ws库的send方法直接传Buffer即可,这个优化在高并发场景能省下大量CPU。
单机连接数有上限,通常单进程稳定维持三五万连接,再多就该扩容了。扩容时WebSocket网关和Redis消息总线配合,每台网关机器只服务一部分观众,弹幕通过Redis在网关之间流转。如果规模继续上量,可以把Redis换成Kafka,按房间ID做分区,保证同一房间的消息有序消费。
流媒体侧的扩容是另一条线。Nginx-rtmp单机带宽有限,通常用边缘节点分发HTTP-FLV和HLS,源站把流推给多个边缘,观众按地域调度到最近的节点,这套架构在SRS里有内置的集群方案,接入成本比自建低不少。Node.js在这一层继续负责调度决策,根据观众的地域和节点负载返回最优播放地址即可。
最后补充几个容易踩的坑:OBS推流地址里的流密钥要和播放地址的流名严格一致,否则会出现推流成功但拉不到流的情况;flv.js播放结束或流中断时要监听错误事件做重连,直播流是持续的,一次拉流失败不代表流没了;弹幕接口务必加频率限制,比如同一用户每秒最多发两条,否则刷屏脚本几分钟就能把整个房间搞瘫。把这些细节处理好,一套支撑万人同时在线的互动直播系统就基本成型了。