第 4 章:发帖、分发与 Feed Fan-out

一、发帖不是一次写入,而是一组副作用

用户点击发布时,最先需要保证的是“这条状态作为事实合法地存在”。之后才是把它传播给不同收件人。Mastodon 把这两个阶段分开,核心入口是 PostStatusService,后续传播由多个 worker/service 承担。

同步阶段:可见性、媒体、提及、引用、反垃圾、事务
异步阶段:Feed、通知、标签流、链接抓取、邮件、ActivityPub

这个边界既保护请求延迟,也防止联邦网络的失败回滚本地事实。

二、PostStatusService#call 的五步

1. 捕获输入上下文

Service 保存 account、text、reply、quote、media、visibility、scheduled_at、idempotency 等选项。这里的 @options 不是随意的 Hash,它代表 UI、REST API、定时发布和内部调用之间的共同用例接口。

2. 媒体校验

媒体必须属于当前账户、尚未绑定到另一个 status、处理完成,并符合数量与音视频组合限制。把校验放在保存前,可以避免出现“状态已经对外可见但媒体仍未完成处理”的状态。

3. 预处理属性

敏感标记、默认隐私、被 silenced 账户的可见性、quoted status 的限制和 scheduled time 都会在这里归一化。注意默认值往往来自用户设置和账户状态,而不是只来自请求参数。

4. 创建关系并事务保存

@status = @account.statuses.new(status_attributes)
process_mentions_service.call(@status)
safeguard_mentions!(@status)
attach_tagged_objects!(@status)
attach_quote!(@status)
 
ApplicationRecord.transaction do
  @status.save!
end

事务保护的是“状态与其必要关系”的一致性。保存成功以后,链接抓取、时间线、通知和远端投递不应继续占用当前 HTTP 请求。

5. 提交后编排副作用

postprocess_status! 会注册 hashtag、安排链接抓取、触发本地 distribution、投递 ActivityPub,并在有 poll/quote/email 时安排专用 worker。它不是一个大而全的同步循环,而是把工作转成队列消息。

三、Fan-out-on-write:把计算成本放在写入时

首页时间线的两种策略:

策略读取时做什么写入时做什么适合
fan-out-on-read每次打开首页查询关注者的最新状态只写状态低写、高变化或关系很小
fan-out-on-write读取预先生成的 feed写入每个收件人的 feed高读、实时首页

Mastodon 对本地 home/list/tag feed 大量采用写入时分发,并通过批量 worker、Redis 和缓存 payload 降低成本。

四、FanOutOnWriteService 的收件人分类

flowchart TB
  S[新 Status]
  S --> L[本地收件人]
  S --> P[公共收件人]
  S --> T[公共 Redis streams]
  L --> Self[作者自己的 home]
  L --> Followers[本地 followers]
  L --> Lists[lists]
  L --> Mentions[mentions / conversations]
  P --> HashtagFollowers[hashtag followers]
  T --> Public[public / local / remote]
  T --> Hashtag[hashtag channels]
  T --> Media[public media]

本地收件人

fan_out_to_local_recipients! 首先处理作者本人,然后根据状态可见性选择:

  • public/unlisted/private:所有本地分发范围内的 followers 与 lists;
  • limited:被提及账户的 followers;
  • direct:会话参与者。

此外,它还会安排引用、提及和编辑更新通知。

公共收件人

公共状态可以进入本地 tag followers 的 feed,也可以广播到公共流。reblog、silenced 作者和 reply 的特殊条件会影响是否 broadcastable。

Redis 公共频道

broadcast_to_public_streams! 会发布 timeline:publictimeline:public:localtimeline:public:remote,有媒体时再发布媒体频道。Streaming server 订阅这些频道后,不需要为每个连接重新查询全部数据库。

五、为什么先 warm payload cache

同一个状态可能被成百上千个 follower 接收。如果每次 FeedInsertWorker 都重新渲染 status,会重复执行 serializer、权限字段组合和媒体 URL 计算。warm_payload_cache! 预先生成可复用的 rendered payload:

Status
  └─ InlineRenderer.render(status, viewer=nil, :status)
       └─ Rails.cache[fan-out/status-id]
            ├─ feed worker 使用
            ├─ streaming anonymous payload 使用
            └─ remote ActivityPub serializer 另行处理

注意缓存 payload 不是所有 viewer 的最终 REST JSON。不同用户的关系字段、可见性和权限可能不同,所以它主要服务于匿名/公共投影或由接收方再次裁剪的路径。

六、FeedManager:Redis 投影的核心

app/lib/feed_manager.rb 同时包含 key 设计、插入/删除、时间线查询和 streaming 触发逻辑。阅读它时先画出 key:

home:<account_id>
list:<list_id>
hashtag:<tag_id>
timeline:public
timeline:public:local
timeline:public:remote

然后追踪三个动作:

  1. push_to_home:把状态 ID 写入用户 home feed;
  2. remove_from_home:删除、屏蔽、取消关注时移除;
  3. merge_into_home / maintenance:关系变化后重建或合并。

实际实现还会处理 reblog、reply、过滤器、忽略关系、缓存状态和是否存在 streaming subscriber。

七、批量 worker 的意义

FeedInsertWorker.push_bulk 这种批量入队方式,把“每个 follower 一条任务”的调度成本降下来。worker 里再按批次读取账户和状态,减少数据库往返。

一个状态
  └─ 找到 50,000 个本地 followers
      └─ push_bulk 生成批量 job
          └─ 多个 Sidekiq 线程分别写 Redis

这条路径仍然有压力:名人账户会产生热点,Redis 写入和 Sidekiq 队列会瞬间放大。因此代码中会用 followers_for_local_distribution、批量查询、缓存 payload、限制 feed 长度和专用维护命令控制成本。

八、幂等与重复投递

异步系统必须接受“同一 job 可能执行多次”。常见防线:

  • Feed 写入使用可覆盖的 status ID,而不是盲目追加;
  • 通知按关系和唯一约束避免重复;
  • 远端 delivery 以 activity/inbox/来源组合追踪;
  • 状态发布支持 idempotency key;
  • 删除和清理操作设计成重复执行也安全。

不要把“队列只执行一次”当作前提。Sidekiq retry、进程崩溃、网络超时和人工重跑都会打破这个假设。

九、用一个具体状态追踪分发

假设 Alice 发布一个 public status,包含 #ruby,提及 Bob,同时 Alice 有 3 个本地 followers:

PostStatusService
  ├─ Status.save!
  ├─ Mention(Bob).save
  └─ enqueue DistributionWorker + ActivityPub::DistributionWorker

DistributionWorker
  └─ FanOutOnWriteService
      ├─ home:Alice ← status
      ├─ home:Follower1..3 ← status
      ├─ tags follower feeds ← status
      ├─ LocalNotificationWorker(Bob)
      ├─ Redis timeline:public
      └─ Redis timeline:hashtag:ruby

ActivityPub::DistributionWorker
  └─ 为 Alice 的 followers/inbox 生成 Create activity
      └─ ActivityPub::DeliveryWorker(host inbox)

如果此时 Redis 暂时不可用,PostgreSQL 中的 status 仍存在;恢复后可以通过维护任务重建 feed,而不需要重新让用户发帖。

十、源码入口清单

十一、小结

发帖链路的设计中心是“事实与传播分离”:数据库事务保证合法性,Sidekiq 保证传播可重试,Redis 保证投影快速,Streaming 保证在线客户端及时看到变化。下一章把传播的一部分拉到实例边界之外,追踪 ActivityPub 的入站和出站。