系统架构总览
架构图
┌─────────────────────── 客户端 ───────────────────────┐
│ OBS 推流(RTMP/WHIP) Web 前端(Vue3 + dash.js)│
└───────────────┬──────────────────────────┬────────────┘
│ WebRTC │ HTTP /api/v1
▼ ▼
┌──────────────────┐ ┌─────────────────────┐
│ live777 SFU │ │ Traefik 统一入口 │
│ (Rust 直播分发) │ │ 按路径分流 │
└──────────────────┘ └──────┬──────────────┘
│ /api/v1/auth|user │ 其余 /api/v1
▼ ▼
┌────────────────────────── 应用服务(可独立扩容)─────────────────────────┐
│ auth (8081 HTTP + 50052 gRPC) api (8080 HTTP) │
│ worker (Kafka 消费者 · 编排) transcoder (50051 gRPC · FFmpeg) │
└───────┬──────────────┬───────────────┬──────────────────┬──────────────┘
│ │ │ │
▼ ▼ ▼ ▼
┌──────────────────────────────────────────────────────────────────┐
│ PostgreSQL Redis MinIO Kafka etcd │
│ 元数据 缓存/计数 对象存储 任务队列 注册中心/领导选举 │
└──────────────────────────────────────────────────────────────────┘设计思想
1. 角色拆分:进程级微服务化
Vistack 采用「单二进制 + 多角色」模式:cmd/vistack 编译出一个二进制,通过 VISTACK_ROLE(或启动参数)分发为不同进程:
- api:HTTP 入口,处理上传、预签名、视频、弹幕、评论、互动等请求,投递 Kafka 消息;
- worker:Kafka 消费者,编排转码、删除文件、弹幕/评论异步落库,运行重试派发器与 Watchdog;
- transcoder:gRPC 转码服务,只连 MinIO(下载原片 / 上传 DASH),不连 DB / Redis / Kafka,完全无状态;
- auth:独立认证服务,持 RSA 私钥签发 RS256 JWT,对外提供 HTTP(注册/登录/资料/JWKS)与 gRPC(用户查询),并注册到 etcd。
每个角色可独立部署、独立扩容,互不阻塞。详见 角色与启动流程。
2. 异步解耦:Kafka 任务队列
上传、转码、删除、弹幕、评论等耗时/异步操作全部通过 Kafka 解耦:
transcode:转码任务(worker 消费,编排远程转码);delete_file:文件删除(worker 消费,异步清理对象与元数据);danmaku/comment:弹幕、评论异步落库;- 失败任务进入 Redis ZSet 延迟队列,指数退避重试,Watchdog 兜底超时任务。
3. 远程计算:gRPC 转码
FFmpeg 隔离在独立的 transcoder 容器中,worker 通过 gRPC ProcessVideo 调用,输入/输出均走 MinIO(S3),天然无状态、可任意副本接任务。transcoder 向 etcd 注册,worker 经 etcd 动态发现 + gRPC round_robin 负载均衡。
4. 对象存储分发
原片、DASH 分片(manifest.mpd + init-*.m4s / chunk-*.m4s)、封面全部存放在 MinIO;上传走预签名直传(Multipart + 秒传去重),播放走预签名 URL / STS 临时凭证,带宽卸载到对象存储,API 服务器不承担大流量。
5. 高并发应用层能力
在业务之上,内置了完整的高并发组件:
| 组件 | 解决的核心问题 |
|---|---|
| Redis 缓存层 | 缓存穿透 / 击穿 / 雪崩三件套 |
| 分布式限流 | 单机令牌桶 + 分布式滑动窗口 |
| 点赞/收藏/播放量 | 计数与落库解耦、热门榜单 |
| 弹幕系统 | 高吞吐写入 + 三级缓存读取 |
| 评论系统 | 异步批量落库 + 敏感词审核 |
一次上传到播放的完整链路
1. 前端分片上传原片 → MinIO(预签名直传,Multipart + 秒传去重)
2. 上传完成 → api(CompleteVideoUpload) → 写 video_sources → 投递 Kafka[transcode]
3. worker 消费 transcode 任务
├─ DB 状态先行校验 + Redis SetNX 租约(幂等)
├─ 置 video_transcodes.status = processing
└─ etcd 发现 transcoder → gRPC ProcessVideo
4. transcoder:MinIO 下载原片 → ffprobe 探测 → 抽封面 → DASH 多档位转码 → MinIO 上传
5. worker 收到结果 → 事务写库(manifest + profiles + 封面 + 视频 published)
6. 前端拉取视频详情 → dash.js 播放 MinIO 上的 DASH 流
7. 失败兜底:指数退避重试(Redis ZSet)→ Watchdog 超时重投 → 7 次后放弃
8. 删除视频 → 投递 Kafka[delete_file] → worker 异步清理对象与元数据可靠性设计
| 场景 | 机制 |
|---|---|
| 消息重复消费 | DB 状态机 + Redis SetNX 租约(幂等) |
| 转码失败 | Redis ZSet 延迟队列指数退避(1min → 2min → … → 8h + jitter) |
| 转码卡死 | Watchdog 每分钟扫描 processing 超 15 分钟的任务,无租约则重投 |
| 消息丢失(pending 卡住) | Watchdog 兜底扫描 pending 超 10 分钟的任务 |
| 多 worker 重复派发 | etcd 领导选举,仅 leader 运行 dispatcher + watchdog |
| 优雅停机 | 信号感知 + 排空在途任务(worker 30s drain / api 30s shutdown) |
| 并发 ID | Snowflake,node_id 自动从 hostname/POD_IP 派生 |
