第 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:public、timeline:public:local、timeline: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
然后追踪三个动作:
push_to_home:把状态 ID 写入用户 home feed;remove_from_home:删除、屏蔽、取消关注时移除;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,而不需要重新让用户发帖。
十、源码入口清单
app/services/post_status_service.rbapp/services/fan_out_on_write_service.rbapp/lib/feed_manager.rbapp/workers/distribution_worker.rbapp/workers/feed_insert_worker.rbapp/workers/push_update_worker.rb
十一、小结
发帖链路的设计中心是“事实与传播分离”:数据库事务保证合法性,Sidekiq 保证传播可重试,Redis 保证投影快速,Streaming 保证在线客户端及时看到变化。下一章把传播的一部分拉到实例边界之外,追踪 ActivityPub 的入站和出站。