跳转至

Q-IM 消息处理流程

状态:配套文档,所有机制以 docs/PLAN.md 契约核心为唯一权威来源。 更新日期:2026-09-01

当前实现覆盖:本文以现行日志驱动链路为主;10 万人群、MailboxNode mailbox-tail 投影压缩器、UserBadgeState 与 NotificationService 属目标形态。当前一期群上限为 1000,qsession 独立消费 dispatch,角标/离线推送完整链路尚未实现。

本文以具体场景追踪消息从发送到接收的完整链路,说明各服务组件在每个阶段做什么、数据存在哪里、lane 水位如何推进。


1. 总体数据流

客户端 SEND_MESSAGE
  │
  ▼
ConnectionNode(准入校验、帧编解码)
  │
  ▼
ConversationWriter(分配序号、写正文到 Redis 热层、追加 Outbox)──> SEND_ACK 返回客户端
  │                                                              ↑ COMMITTED 层级
  ├──> HistoryArchiver(独立消费 Outbox、批量归档 ScyllaDB、发布归档水位)  ↑ ARCHIVED(异步,不在 ACK 路径)
  ▼
FanoutCoordinator(消费 Outbox、固化成员版本、写分发日志)
  │
  ▼
MailboxNode(消费分发日志、物化邮箱引用、推进 W[lane])  ↑ MAILBOXED 层级
  │
  ├──> 在线推送:PushBatch -> ConnectionNode -> PUSH_EVENTS -> 客户端  ↑ PUSHED 层级
  └──> 离线推送:PushTask -> NotificationService -> APNs/FCM          ↑ NOTIFIED 层级

三个独立阶段,各自有独立的失败域与重试语义(§4.1 原则 1): 1. 消息提交(ConversationWriter) 2. 收件人索引(FanoutCoordinator + MailboxNode) 3. Socket 投递(ConnectionNode)


2. 群聊消息完整流程

目标容量场景:10 万人群,张三发送"你好"

初始状态

群成员:10 万人
MailboxShard 7 上有 3000 个群成员,分散在 64 个 lane
  Lane 5:  47 人(含张三)
  Lane 17: 0 人
  其他 lane: 各 ~47 人

当前水位:W[0..63] 全部 = 9998
consumed_offset = 9998

步骤 1-5:提交阶段(ConversationWriter)

1. 客户端 -> ConnectionNode
   SEND_MESSAGE{client_message_id=abc, conversation_id=群3, payload="你好"}

2. ConnectionNode
   L1 发送者维度准入(per_sender_in_conversation_rate: 1 msg / 3s)
   按 conversation_id 路由到该会话 Home Region 的 ConversationWriter

3. ConversationWriter
   a. 身份与成员关系校验(membership_state 必须为 ACTIVE)
   b. 内容尺寸校验
   c. L1 复核 + L2 会话维度准入 + L3 租户 fanout 配额
   d. ClientDedup 幂等占位(IF NOT EXISTS)
   e. 同一临界区内分配:
      message_id = 0x1234...(HLC,u128)
      conversation_seq = 42(单会话严格递增)
      last_activity_id = 0x5678...(与 message_id 同源同序)

4. ConversationWriter 写存储
   -> Redis 热层: MessageRecord + 历史索引(与序号分配同一 Lua;正文"你好"热层 1 份)
   -> Redpanda: 追加 Outbox 日志(1 条记录,含完整 canonical MessageRecord)
   -> ScyllaDB: 由 HistoryArchiver 异步批量归档(永久 1 份);ConversationHead 异步更新(失败不阻断 fanout)

5. ConversationWriter -> SEND_ACK -> ConnectionNode -> 客户端
   SEND_ACK{client_message_id=abc, message_id=0x1234, conversation_seq=42}

   ↑ 此时到达 COMMITTED 层级
   正文与 Outbox 已可靠提交,但任何接收者的邮箱引用尚未物化

步骤 6-8:分发阶段(FanoutCoordinator)

6. FanoutCoordinator 消费 Outbox
   读到 {message_id=0x1234, conversation_id=群3, conversation_seq=42, ...}

7. FanoutCoordinator 固化 membership_version
   取 GroupMembership 当前已提交版本作为惰性快照
   写入 GroupDispatch.membership_version
   固化点:conversation_seq 分配之后、追加分片分发日志之前

8. FanoutCoordinator 写分发日志
   按 membership_version 读取分片成员 Bitmap
   得到目标 MailboxShard 列表(S 个,10 万人群 ≈ 256 个分片)

   对每个目标 MailboxShard 追加 1 条 GroupDispatch 到该分片的分发日志:
   dispatch_id = blake3("qim.group-dispatch.v1\0" ||
                        message_id_be16 || target_mailbox_shard_be4)[0:16]

   GroupDispatch 内容:
     tenant_id, dispatch_id, message_id, conversation_id,
     conversation_seq, last_activity_id, sender_id,
     fencing_epoch(目标字段;当前编解码尚缺,发布阻断), membership_version,
     target_mailbox_shard,
     event_template, committed_at

   注意:GroupDispatch 不含正文,正文留在 MessageStore

   同一目标分片只生成 1 个任务,无论该分片上有多少成员

步骤 9-12:物化阶段(MailboxNode,以 Shard 7 为例)

9. MailboxNode 消费分发日志
   从 Redpanda 分区消费到 offset=9999 的 GroupDispatch
   consumed_offset = 9999

   按 dispatch_id 去重(dispatch_progress_retention = 7 天窗口内)

10. MailboxNode 加载成员 Bitmap 并按 lane 分组
    加载该 membership_version 的本分片不可变 RoaringBitmap(LRU 缓存)
    本分片 3000 成员,按 lane_id 分组:
      Lane 0:  48 人 -> 创建子任务 0
      Lane 5:  47 人 -> 创建子任务 5(含张三)
      Lane 17: 0 人 -> 不创建子任务 ← 关键
      ...
      Lane 63: 51 人 -> 创建子任务 63

    成员边界过滤(目标契约;当前尚未落地 joined/left 序号区间):
      跳过 left_at_conversation_seq <= 42 的成员(退群成员)
      跳过 joined_at_conversation_seq >= 42 的成员(新入群成员)

    每条条目携带 visibility_floor_conversation_seq

    一期实现(ADR-0021):成员变更与序号分配同一原子区,版本 V 的快照恰是分配
    本消息序号时的成员,按精确 membership_version 展开即等价于上述过滤;历史与会话
    列表另按成员区间裁剪。visibility_floor_conversation_seq 尚未写入邮箱条目。

11. MailboxNode 各子任务并行写入 UserMailboxEntry
    每个子任务按 mailbox_writebatch_max_entries(2000)分块
    条目与 DispatchProgress{lane_id, chunk_done_bitmap} 按顺序约束提交:
      先写全部 UserMailboxEntry(LOCAL_QUORUM)
      再写 DispatchProgress 的 chunk 位

    每条 UserMailboxEntry 内容(~110 字节,不含正文):
      tenant_id, user_id, mailbox_seq=9999,
      event_ordinal=0 (MESSAGE), event_id=确定性哈希,
      message_id=0x1234, conversation_id=群3, conversation_seq=42,
      sender_id=张三, flags={counts_unread=true, affects_session_order=true},
      created_at=dispatch.committed_at

12. MailboxNode 推进持久 64-lane 水位
    当前日志模型按同一 dispatch partition 顺序处理 record。offset=9999 的全部
    lane/chunk/UserMailboxEntry durable,且每个 lane 的 expected_previous=9998
    通过 CAS 连续校验前,任何 lane 都不能把 W 越过 9998。

    约 50ms 后该 record 的全部子任务完成:
      advance_log_watermarks(observed=9999, expected_previous=9998)
      W[0..63] -> 9999

    空 lane 不代表可独立越过当前 record;旧直写模型的 pending 集合与
    W=min(pending)-1 公式不适用于现行实现。

    ↑ 此时到达 MAILBOXED 层级
    接收者的邮箱引用已可靠物化并被 W[lane] 覆盖

步骤 13-14:在线推送阶段

13. MailboxNode 合并 PushBatch
    物化完成并推进 W[lane] 后:
    a. 成员 Bitmap ∩ 在线成员 Bitmap(快速过滤)
    b. 经 PresenceDirectory 本地缓存展开 per-device 明细
       (推送路径 0 次同步远程调用)
    c. 按 ConnectionShard 合并成 PushBatch

    PushBatch 内容:
      1 份 common_encoded_body(正文只编码 1 次)
      + 多个轻量个性化 recipients[](含 connection_id, session_epoch,
        mailbox_seq, event_ordinal, event_id, flags, mention_type 等)

    默认拓扑下每个 MailboxShard <= 4 个 ConnectionShard
    -> 1 条 GroupDispatch 产生 <= 4 个 PushBatch

14. ConnectionNode 展开 PushBatch -> 写入客户端 Socket
    for r in batch.recipients:
      校验 connection_id 存在且 session_epoch 匹配
      不匹配 -> 丢弃该收件人,回 PRESENCE_STALE
      匹配 -> 生成 PUSH_EVENTS 帧,引用 common_encoded_body(零拷贝)
      写入对应 stream(stream 1 实时流)

    ↑ 此时到达 PUSHED 层级

离线推送(并行)

MailboxNode 在物化后判断离线设备:
  对每个无有效 PresenceEntry 的设备产生 PushTask
  PushTask 主键 = (tenant_id, user_id, device_id)

  -> NotificationService
  -> APNs / FCM / 厂商通道

  ↑ 此时到达 NOTIFIED 层级(仅表示已交给外部通道)

3. 单聊消息流程

单聊与群聊复用同一条 fanout 路径(§8.1),不设特例短路。

与群聊的差异

维度 群聊(10万人) 单聊(2人)
GroupDispatch 数量 S ~256 ≤ 2(发送者与接收者各一片,同片合并为 1)
每分片邮箱引用数 ~3000 1-2
正文存储 1 份 MessageRecord 1 份 MessageRecord
邮箱引用存储 10 万条 UserMailboxEntry 2 条 UserMailboxEntry
子任务耗时 ~50ms/lane ~1ms/lane

发送者也写邮箱条目的原因

单聊产生 2 条 UserMailboxEntry(发送者 1 条 + 接收者 1 条):

  1. 多设备同步:发送者手机发的消息,电脑通过 PULL_MAILBOX 拉到同一条 entry
  2. 统一去重:entry 携带 client_message_id(仅发送者条目),多设备按此去重
  3. 会话列表更新:当前 qsession 统一消费同一 dispatch 的收发两侧 recipients, 不需要“发送侧”单独逻辑;阶段二才由 mailbox-tail 投影压缩器承担同一职责
  4. 统一同步路径:一次 PULL_MAILBOX 恢复全部(收到的 + 发出的)

发送者条目的 flags.counts_unread = false(自己发的消息不计未读)。


4. lane 水位推进机制(现行日志模型)

核心条件

对 dispatch record(offset = s):
  仅当该 record 的全部 lane/chunk/entries 已 durable,
  且 expected_previous == 当前持久连续覆盖点时,
  才以 CAS 将 W[0..63] 推进到 observed_log_offset = s。

任何失败、未完成 chunk 或 expected_previous 缺口:W 保持旧值。
  • observed_log_offset:当前顺序处理的 dispatch partition offset。
  • expected_previous:持久 CAS 的连续性前提,禁止从 9998 直接跳到 10000。
  • W[lane]:该 lane 的持久可见水位;用户拉取邮箱时只能看到 mailbox_seq <= W[lane_id(user)] 的条目。
  • W = min(在途) - 1 与分配时登记/完成时销账只属于旧 TCP 直写序号模型。

时间线示例

场景:Shard 7 上,offset=9999 与 offset=10000 两条 record 顺序到达

时刻               当前 record    durable/CAS 状态                    W[0..63]
──────────────────────────────────────────────────────────────────────────────
初始               —              连续覆盖到 9998                    9998
T=0 开始 9999      9999           部分 chunk 未完成                  9998
T=5 完成 9999      9999           全部 durable;CAS(9998 -> 9999)    9999
T=5 开始 10000     10000          部分 chunk 未完成                  9999
T=5.1 完成 10000   10000          全部 durable;CAS(9999 -> 10000)   10000
  • W 表示连续覆盖证明,不是“见过的最大 offset”。
  • 64-lane 是现行持久格式与用户可见水位接口;稳态 packed-v2 可压成一份全同记录。
  • 当前 ordered record 路径不会让空 lane 提前越过尚未完成的 record;若未来恢复 lane 独立推进,必须另立 ADR、给出跨重启连续性证明并更新本节。

record 内并行,record 间有序

MailboxNode 可在一个 record 内并行执行 lane/chunk 写入,但持久可见水位按 record 有序:

poll record(s) -> 按 lane/chunk 并行写 UserMailboxEntry / DispatchProgress
               -> 等该 record 全部 durable
               -> expected_previous CAS 推进 W
               -> 再确认下一连续 record 可推进

后续 record 可以已被 broker 预取,但不能越过未完成前项发布可见水位。


5. 离线同步流程(§9.3)

登录同步时序

客户端 -> AUTH{access_token, device_id, mailbox_cursor}

服务端游标判定(顺序固定):
  1. 游标签名无效 or seq 越界 -> ERROR{CURSOR_INVALID}
  2. cursor.last_applied < effective_trim(user) -> ERROR{CURSOR_EXPIRED, rebuild_required=true}
  3. cursor.shard_epoch 落后但边界可解析 -> ERROR{CURSOR_REBASED, new_cursor, replay_from_seq}
  4. 以上都不触发 -> AUTH_OK

服务端 -> AUTH_OK{
    session_epoch, lane_id, lane_watermark, trim_watermark,
    sync_to_seq = AUTH 时刻的 lane_watermark 快照,  ← 权威等式
    has_offline, pending_entry_count_hint, pending_bytes_hint,
    total_unread, total_mention, muted_unread,
    projection_complete, preferred_endpoint,
    sync_delay_hint_ms, next_ping_interval_ms
}

客户端 -> PULL_MAILBOX{after_seq=cursor.last_applied, up_to_seq=sync_to_seq, ...}
服务端 -> MAILBOX_BATCH{entries[], covered_through_seq, lane_watermark, has_more}
  可流水线:在途请求数 <= pull_mailbox_window = 4
  批次切分只能在 mailbox_seq 边界(事件组不可切分)

客户端 -> SYNC_COMPLETE{sync_to_seq}
服务端校验:
  1. sync_to_seq == 本次 AUTH_OK 分配值
  2. 各子区间 covered_through_seq 已连续覆盖到 sync_to_seq
  任一不满足 -> ERROR{SYNC_INCOMPLETE}

服务端 -> ONLINE_READY(解除登录屏障)

sync_to_seq 是地基:把设备事件一刀切为"客户端拉"与"服务端推"两半。

空洞跳过的前置条件(§9.3.3 规则 3)

仅当 after_seq >= effective_trim(user) 时,
(after_seq, covered_through_seq] 区间内的序列空洞才可安全跳过。

否则必须返回 ERROR{CURSOR_EXPIRED},
禁止退化为"返回空批次 + has_more=false"的静默成功。

实时推送不得推进游标(§6.8 不变量)

设备游标只能由 MAILBOX_BATCH 连续推进。
PUSH_EVENTS 送达的事件不得推进 cursor.last_applied_mailbox_seq。

违反此不变量会导致:
  实时推送越过尚未拉取的区间 -> 游标被推到更高位置
  -> 未拉取区间的消息按规则 3 被判为"可跳过"
  -> 永久丢失,且不可检测

6. REBUILD 流程(§9.6)

触发条件:CURSOR_EXPIRED(游标早于 mailbox_trim_watermark)

0. 重新握手:保留旧令牌身份,以 mailbox_cursor=空 重新 AUTH
   -> AUTH_OK{sync_to_seq=lane_watermark, has_offline=false}
   -> SYNC_COMPLETE -> ONLINE_READY

1. 保留本地已有消息,不清空本地库

2. PULL_SESSION_LIST 取得权威会话集合
   来源分层:
     UserConversationState = 会话集合持久权威
     ConversationHead = 每个会话公共最新状态
     UserSessionProjection = 可重建缓存

3. 逐会话补齐:
   anchor = 本地该会话最大连续 conversation_seq
   PULL_HISTORY{conversation_id, direction=newer, anchor, limit}
   直到 has_more=false 或达到 latest_conversation_seq
   按 message_id 去重

4. 未读重算:
   latest - base <= unread_precise_limit(200) -> 精确重算
   超过 -> unread_count=200, unread_exact=false

5. 窗口外提及数不保证精确,UI 展示为"99+"

6. 服务端换发签名游标令牌:
   cursor.last_applied := 本次 REBUILD AUTH_OK 的 sync_to_seq

7. 补齐期间新到事件正常走实时队列

撤回/编辑不需要邮箱事件流即可收敛:PULL_HISTORY 返回的 MessageRecord.state 已是终态。


7. 会话列表更新流程(目标契约与当前实现对照)

三层模型(§12.2)

层 结构 权威性 更新触发
会话公共最新状态 ConversationHead 会话维度权威 每条消息一次,与成员数无关
用户会话集合与主动状态 UserConversationState 持久权威 只因用户操作或成员关系变化
用户会话列表视图 UserSessionProjection 可重建的物化视图 当前 qsession 消费 dispatch;阶段二可由个人邮箱增量合并

当前一期:qsession durable dispatch 投影(§17.2)

qsession 独立消费 dispatch 日志;同一 record 的全部 recipients 以 Redis `ZADD GT`
幂等写入成功后,才 CAS 推进该 partition 的 durable checkpoint。读取会话列表时以
`convs:{user}` 为候选集合,并从 canonical MessageStore 读取会话摘要与定义式未读。

MailboxNode 的 TCP Projection 只作尽力低延迟补充,不拥有 qsession checkpoint。
UserConversationState(一期为会话集合子集,ADR-0021):投影 Lua 在会话首次进入集合时写入
待登记队列,后台批量落 ScyllaDB;checkpoint 越窗或投影全失时递增投影纪元,按用户惰性重建。

阶段二目标:mailbox-tail 投影压缩器(§12.3)

压缩器 = MailboxNode 内的常驻协程,每 lane 一个

输入:本节点本地邮箱写入流(与物化 UserMailboxEntry 是同一次 WriteBatch 产生的事件序列)
消费上界 = W[lane]:只按 mailbox_seq 顺序消费 <= W[lane] 的连续前缀

聚合:内存中按 (user_id, conversation_id) 折叠
      同一用户同一会话在一个窗口内只保留一份聚合结果

flush:每 projection_compaction_window(5s)或 512 条事件
       按 user_id 分组,每用户一个 WriteBatch,
       同批写入 UserSessionProjection + UserBadgeState
       projection_mailbox_seq 推进到 min(本批最大 mailbox_seq, W[lane])

明令禁止按用户主键全表扫描。

上述 mailbox-tail、projection_mailbox_seq 与同批 UserBadgeState 尚未按本形态实现, 不得与当前 qsession 路径混写成一条已上线链路。

SESSION_DELTA(§12.6)

幂等绝对值帧,版本源 = projection_mailbox_seq

客户端规则:
  if delta.projection_mailbox_seq > local[conv_id].projection_mailbox_seq:
      用绝对值直接覆盖 unread_count / mention_count / latest / preview 等
  else:
      整帧丢弃(迟到帧或重复帧,不做补偿)

权威声明:
  客户端未读的权威输入是 UserMailboxEntry 与 read_conversation_seq。
  SESSION_DELTA 只是低延迟缓存,任何冲突以邮箱重算为准。

8. 聊天室消息流程(§14)

聊天室与普通群不共用持久邮箱语义:

RoomWriter
  -> 校验发送权限与 room_msg_rate(20 msg/s/房间)
  -> 分配 message_id(用于去重与举报)与 room_seq
  -> 写短期 RoomRecord(保留 room_log_retention_minutes = 30 分钟)
  -> 按"存在该房间连接"的 ConnectionShard 合并广播
ConnectionNode
  -> 按 room_outbound_frame_rate(10 frame/s/连接)合并
  -> 写入 ROOM_BATCH(stream 3 房间流)
维度 普通群 聊天室
消息序列 conversation_seq(持久) room_seq(短期)
邮箱引用 UserMailboxEntry(N 条) 无(不产生)
会话列表投影 有 无
离线推送 有 无
离线可达 精确可达(30 天内) 仅短窗口回放(30 分钟)
丢消息 不允许 允许(§14.3 明确声明)

9. 故障恢复流程

MailboxNode 接管(§10.4.2)

MailboxNode 的状态性随 MailboxStore 实现阶段而变(§17.1):

阶段 MailboxStore MailboxNode 权威数据位置 接管方式
一期 Redis 无状态计算节点 Redis(节点外) 漂移接管(形态 B)
阶段一 ScyllaDB 无状态计算节点 ScyllaDB(节点外) 漂移接管(形态 B)
阶段二 自研 LSM 有状态,热备可选 节点本地 LSM 热备(形态 A)或漂移(形态 B)

一期/阶段一(无主备,漂移接管):

T=0:    MailboxNode 宕机
        权威数据全在 Redis/ScyllaDB(RF=3),不丢失

T=15s:  shard_lease_ttl 到期,ShardRegistry 授予新节点租约
        递增 shard_epoch,写入 EpochBoundary

T=15s:  新节点恢复:
  1. 从 Redis/ScyllaDB 读取 DispatchProgress(chunk_done_bitmap)
  2. 重算 W[lane](从已完成的 dispatch 推导,纯内存计算)
  3. 从 W[lane] 对应 offset + 1 重放分发日志
     - 按 dispatch_id 去重(已完成的跳过)
     - 未完成的 -> 创建子任务,写入 UserMailboxEntry
     - 确定性 event_id 保证同值覆盖(§6.7)
  4. 逐 lane 满足追平判据后恢复服务

T=~20s: 追平完成,逐 lane 恢复
        游标落后的客户端收到 CURSOR_REBASED

实际 RTO ≈ 15s(租约)+ ~5s(重算+重放)= ~20s
(shard_rto_target ≤ 10min 为保守上界,覆盖极端积压场景)

不需要从 S3 加载检查点:权威数据在 Redis/ScyllaDB 中,MailboxNode 只是计算节点。

阶段二(有状态,热备可选):

形态 A 热备接管(可选 RTO 优化):
  备 MailboxNode 一直在消费同一分发日志并物化相同索引
  主节点宕机 -> 备节点满足追平判据后接管
  RTO = mailbox_takeover_rto_target(30s)

形态 B 漂移接管(无热备或主备同时丢失):
  新节点从 S3 加载检查点(全量基线 + 增量段)
  从 checkpoint.log_offset + 1 重放分发日志
  RTO = shard_rto_target(10min)

两种形态的判定规则完全相同(§10.4.2):

(1) 接管起始位点 = min_j(W[j]) 对应的 log_offset + 1
    不取消费者组已提交位点(否则"先提交位点后物化"的崩溃会使该 lane 永久停滞)

(2) 校验 重算 W[j] >= 持久化的 W_floor[j],不满足即拒绝服务 + P1

(3) 逐 lane 判定:追平后才对该 lane 服务,未追平返回 SHARD_MOVED
    禁止用未追平的水位回答 PULL_MAILBOX(会让 §9.3 规则 3 看到假空洞)

接管必然递增 shard_epoch,旧游标按 EpochBoundary 换发(CURSOR_REBASED)

分片分裂流程(§5.5)

触发:单 MailboxShard entry/s 持续超过 per_shard_entry_budget 的 70%

1. ShardRegistry 发起,冻结映射表版本
2. 递增 shard_epoch,为新旧分片各分配新 epoch
3. 写 ShardSplitBoundary
4. 双写窗口开启(迁出桶的 dispatch 同时写旧分片与新分片)
5. 新分片追平(一期/阶段一从外部存储重算 W[lane] + 日志重放;阶段二从 S3 检查点 + 日志重放)
6. 切读(映射表版本 +1)
7. 旧分片停写
8. 游标换发(下发 CURSOR_REBASED)

lane_id 跨分裂稳定(独立哈希导出,不含 mailbox_shard_count)

10. 端到端延迟分解

阶段                              耗时         累计      SLO
──────────────────────────────────────────────────────────────
1. 客户端 -> ConnectionNode       ~5ms        5ms
2. ConnectionNode -> ConvWriter   ~5ms        10ms
3. ConvWriter 写 MessageRecord    ~1ms        11ms      (Redis 融合 Lua,热层;ADR-0018)
4. ConvWriter 追加 Outbox         ~3-5ms      16ms      (Redpanda produce)
5. 返回 SEND_ACK                  ~1ms        17ms      ← P99 ≤ 150ms(10k msg/s 实测 P50 8 / P99 17 ms)
──────────────────────────────────────────────────────────────
6. FC 消费 Outbox                 ~1-5ms      35ms      (Redpanda poll)
7. FC 写分发日志                  ~3-5ms      40ms      (Redpanda produce)
8. MailboxNode 消费分发日志       ~1-5ms      45ms      (Redpanda poll)
9. MailboxNode 物化邮箱引用       ~5-50ms     95ms      (取决于成员数)
10. MailboxNode 发 PushBatch      ~5ms        100ms     (内部 RPC)
11. ConnectionNode 写 Socket      ~1ms        101ms
──────────────────────────────────────────────────────────────
端到端(发送到对方收到)                               ← P99 ≤ 300ms

以上分段仍以设计估算为主;当前链路有两项实测可校准:10k msg/s × 180 s 下 SEND_ACK P50 8 / P99 17 ms,端到端 P50 73 / P99 108 ms(开发机,非目标硬件)。ACK 之后约 65 ms(P50) 主要是 fanout 批周期(每批约 30 ms,卡在 librdkafka 客户端逐分区请求而非 broker, ADR-0016/0017)+ mailbox 物化 + 推送。步骤 6 + 8 的拉取延迟仍无独立分量实测, 不得按 300ms 预算反推出固定占比。MailboxNode 可预取后续 record,但当前可见水位必须等待 前一 record 全部 durable,慢大群会阻塞同 shard 后续 record 的可见推进。


11. 四道阻塞防线

防线 机制 解决问题 章节
lane 水位隔离(目标) W[lane_count] 向量水位 目标是大群只阻塞同 lane;当前整条 record 完成后统一推进 64 lane,尚未实现 §6.5.1、§9.2
公平调度 大群子任务低优先级,单聊高优先级 资源竞争(物理层):大群不吃满 CPU/IO §10.3
租户配额 tenant_fanout_quota 按租户限流 一个租户的大群风暴不影响其他租户 §10.5
lane 停滞超时 lane_stall_failover = 120s 极端情况下拖住有上界,上界到了换节点 §9.2.4

四道防线正交,各解决不同层面的问题,缺一不可。