第 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 事件通常有 eventpayloadstream 等字段。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 中的 WebSocketSessionisAlive、close/error 处理和 metrics,是读“实时服务如何避免僵尸连接”的重点。client.send 失败不能拖垮 Redis 消费循环,异常连接要被及时清理。

九、频道权限与安全

/api/v1/streaming/usernotificationdirect 显然不是匿名频道。安全检查应同时覆盖:

  1. access token 是否有效;
  2. token scope 是否包含必要能力;
  3. account 是否仍可用;
  4. path 中的 account/list/hashtag 是否属于请求者可见范围;
  5. Redis payload 是否被误发到公共频道;
  6. 连接关闭后订阅是否真正取消。

实时协议的漏洞常常不是 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 是否增长。

十二、源码入口清单

十三、小结

实时体验不是 WebSocket 单点功能,而是 Rails publish、Redis channel、Node subscription、数据库过滤和 React normalized state 的合作。下一章转向浏览器内部,追踪 React/Redux 如何把 REST 与 streaming 的两种输入合并成 UI。