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 条):
- 多设备同步:发送者手机发的消息,电脑通过 PULL_MAILBOX 拉到同一条 entry
- 统一去重:entry 携带
client_message_id(仅发送者条目),多设备按此去重 - 会话列表更新:当前 qsession 统一消费同一 dispatch 的收发两侧 recipients, 不需要“发送侧”单独逻辑;阶段二才由 mailbox-tail 投影压缩器承担同一职责
- 统一同步路径:一次 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 |
四道防线正交,各解决不同层面的问题,缺一不可。