第 6 章:Redis、Streaming 与实时同步
一、实时更新的两段式设计
Mastodon 把“事件产生”和“连接管理”拆开:
Rails / Sidekiq
└─ redis.publish(channel, payload)
└─ Node Streaming 订阅
└─ 找到在线 sessions
└─ filter + websocket.send
└─ React streaming client
Rails 不需要知道每个浏览器 socket 的生命周期,Streaming 也不需要理解完整的发帖业务。Redis 是两者之间的低延迟事件总线。
二、频道命名是协议
streaming/index.js 声明了 system、user、user:notification、list、direct、public、public:media、public:local、public:remote、hashtag 等频道。命名不是内部随意字符串,因为它同时被:
- Rails
redis.publish; - Node subscription router;
- API 文档中的 streaming parameter;
- 前端连接 URL;
- 测试 fixture;
共同消费。修改频道名必须视作协议变更。
三、连接建立 sequence
sequenceDiagram
participant B as Browser
participant N as Node Streaming
participant R as Redis
participant P as PostgreSQL
B->>N: WebSocket /api/v1/streaming/user?access_token=...
N->>P: resolve token/account/permissions
N->>R: subscribe user channel
N-->>B: upgrade accepted
R-->>N: event payload
N->>N: filter by session / language / permissions
N-->>B: JSON event
WebSocket upgrade 前必须先鉴权;否则匿名连接可能订阅用户流。Node 服务会解析 token、scope、account id、语言、permissions 和缓存过滤器,并把这些结果挂到 request/session 上。
四、为什么 Streaming 还要查 PostgreSQL
Redis 消息只包含事件和渲染 payload,但连接是否应该收到它,可能依赖当前数据库事实:
- 用户是否仍关注作者;
- list 成员关系是否改变;
- 用户是否 block/mute 了作者;
- 用户权限是否允许查看某个 feed;
- 语言/过滤器设置是否改变。
因此 Streaming 服务使用 PostgreSQL pool 做必要的关系查询。它不是完全无状态的 Redis relay,而是一个带轻量权限过滤的连接服务。
五、Redis Pub/Sub 与 Feed 的区别
Feed
Feed 是可回放的有序投影:用户断线后重新打开首页,可以根据 Redis/database 读取过去的状态。
Pub/Sub
Pub/Sub 是瞬时事件:消息只在订阅期间到达;连接断开时不自动补发。因此实时事件一般不能作为唯一事实来源。
数据库 / feed:可回放
Redis Pub/Sub:低延迟通知
客户端:收到事件后更新内存 store
重新连接:重新拉取时间线,修复丢失窗口
这就是为什么 streaming 客户端需要重连和重新获取,而不是假设每条 websocket message 永不丢失。
六、从 FanOutOnWriteService 到 socket
以公共状态为例:
redis.publish('timeline:public', anonymous_payload)
redis.publish('timeline:public:local', anonymous_payload) if @status.local?
Node 侧订阅对应频道后,onRedisMessage 解析 JSON,再按照频道/连接类型找到 sessions。对于 user、notification、list、direct,频道通常带 account/list id;对于 public/hashtag,连接级过滤更重要。
七、前端事件类型与 reducer
Streaming 事件通常有 event、payload、stream 等字段。React 客户端收到后会转成内部 action,例如:
update → accounts + statuses + timeline insert
status.update → replace existing status
notification → notifications append + unread count
delete → remove status from all relevant timelines
filters_changed → refresh or update filter state
事件处理不应该只更新一个数组:状态卡片依赖 account、media、poll、relationship 和 status entities。Redux 的 normalized store 让同一实体被多个页面共享。
八、心跳、关闭与反压
长连接必须处理:
- ping/pong 保活;
- 客户端主动关闭;
- Redis 连接重连;
- Node 进程优雅退出;
- 单个慢客户端的发送缓冲;
- 大量连接下的内存与 metrics。
streaming/index.js 中的 WebSocketSession、isAlive、close/error 处理和 metrics,是读“实时服务如何避免僵尸连接”的重点。client.send 失败不能拖垮 Redis 消费循环,异常连接要被及时清理。
九、频道权限与安全
/api/v1/streaming/user、notification、direct 显然不是匿名频道。安全检查应同时覆盖:
- access token 是否有效;
- token scope 是否包含必要能力;
- account 是否仍可用;
- path 中的 account/list/hashtag 是否属于请求者可见范围;
- Redis payload 是否被误发到公共频道;
- 连接关闭后订阅是否真正取消。
实时协议的漏洞常常不是 REST endpoint 本身,而是“连接已授权但频道切换后没有重新检查”。因此读 subscription routing 和 authorization helper 时要一起看。
十、事件丢失和重连策略
合理的客户端策略是:
connected
├─ 收到 update → 立即写入 store
├─ 收到 delete → 删除实体和时间线引用
├─ socket error/close → 指数退避重连
└─ 重连成功 → 拉取最新 timeline,补齐断线窗口
事件流只负责“快”,API 拉取负责“准”。二者配合,才能在 Pub/Sub 瞬时消息的语义下提供可靠体验。
十一、性能观察点
Rails/Redis 侧
- 一个 status 被 publish 到多少公共频道;
InlineRenderer和缓存命中率;- Pub/Sub payload 大小;
- fan-out 写入和 streaming publish 是否在同一 worker 中串行。
Streaming 侧
- connected clients;
- subscribed channels;
- Redis messages received;
- per-connection send count;
- 数据库查询延迟和 pool 使用率;
- 单连接 outgoing buffer 是否增长。
十二、源码入口清单
streaming/index.jsstreaming/database.jsstreaming/redis.jsstreaming/metrics.jsapp/services/fan_out_on_write_service.rbapp/workers/push_update_worker.rbapp/javascript/mastodon/streaming
十三、小结
实时体验不是 WebSocket 单点功能,而是 Rails publish、Redis channel、Node subscription、数据库过滤和 React normalized state 的合作。下一章转向浏览器内部,追踪 React/Redux 如何把 REST 与 streaming 的两种输入合并成 UI。