02 消息模型与存储¶
状态:可实施 更新日期:2026-09-01 上游契约:
docs/PLAN.md§6.2~§6.4、§7.0~§7.2、§7.4、§7.15、§18、§25.1 实施决策:ADR-0006(Rust + Redpanda)、ADR-0007(一期简化形态)、ADR-0012(Fanout 无 Kafka 事务)
0. 文档边界¶
本文定义:ScyllaDB 建表 DDL、压缩策略、二级查询、MessageStore 读写路径、
媒体元数据编码、ID 在存储层的字节表示。
本文不得重新定义:
§7.0 的分区键 / 聚簇键 / 单分区上界
retention_class 的取值集合
§6 的序列格式与语义
不在本文范围:MailboxStore(UserMailboxEntry / UserSessionProjection /
UserBadgeState / DispatchProgress)属 docs/03;会话投影读路径属 docs/04。
本文只覆盖 MessageStore 侧。
1. ID 在存储层的字节表示¶
契约层的 ID 是 u128 / u64(§6),CQL 没有无符号类型,必须写死映射规则。
| 契约类型 | CQL 类型 | 规则 |
|---|---|---|
message_id / last_activity_id(u128) |
blob(16 字节) |
大端,字节序比较 == 数值比较(§6.1) |
conversation_seq(u64) |
bigint |
见下方有符号约束 |
mailbox_seq(u64 复合) |
bigint |
同上 |
event_id(u64 = blake3[0:8]) |
bigint |
仅作等值匹配,不参与范围比较,符号无影响 |
tenant_id / conversation_id / user_id |
uuid |
|
dek_id |
frozen<tuple<text,text,int>> |
(key_scope, key_id, key_version),§21.3 |
有符号约束(必须写进实现并加断言):
CQL bigint 是有符号 64 位。要让 CQL 的排序与 u64 无符号排序一致,
必须保证 bit63 恒为 0:
conversation_seq 由 ConversationWriter 单调分配,实际值远小于 2^63
mailbox_seq 高 16 位是 shard_epoch,§6.5 已约束 shard_epoch <= 0x7FFF
因此 bit63 恒零
实现断言:写入前检查 (value & 0x8000_0000_0000_0000) == 0,违反即 panic。
这条一旦被破坏,聚簇键排序会在跨过 2^63 时整体翻转,历史分页立刻错乱。
message_id 用 blob 而不是拆成两个 bigint,是因为 §6.9.2 的跨会话时间轴
需要它整体可比;拆列后 CQL 无法对"两列合成的 128 位值"做单一范围比较。
2. 建表 DDL¶
2.1 MessageRecord¶
CREATE TABLE message_record (
tenant_id uuid,
conversation_id uuid,
seq_bucket bigint, -- conversation_seq / message_seq_bucket_width
conversation_seq bigint,
message_id blob, -- u128 大端 16 字节
last_activity_id blob, -- u128 大端 16 字节
sender_id uuid,
sender_device_id uuid,
client_message_id text,
message_type tinyint,
custom_type text,
schema_version int,
payload_or_ciphertext blob,
media_metadata blob, -- Protobuf 编码,见 §5
reply_to_conversation_seq bigint,
mention_targets blob, -- Protobuf 编码,见 §5.2
state tinyint, -- NORMAL|RECALLED|EDITED|DELETED
recall_event_conversation_seq bigint, -- §13.4.1 固化顺序的判定位
edited_at_activity_id blob,
created_at bigint,
retention_class tinyint,
dek_id frozen<tuple<text,text,int>>,
PRIMARY KEY ((tenant_id, conversation_id, seq_bucket), conversation_seq)
) WITH CLUSTERING ORDER BY (conversation_seq DESC)
AND compaction = {'class':'TimeWindowCompactionStrategy',
'compaction_window_unit':'DAYS',
'compaction_window_size':7}
AND gc_grace_seconds = 86400;
主键即幂等:(tenant_id, conversation_id, seq_bucket, conversation_seq) 唯一定位一行,
CQL INSERT 天然是同键覆盖(upsert)。这不是巧合而是要求——至少一次投递下,
重放的写入必须落到同一槽位而不是追加成第二行。对照:OpenIM 的 MongoDB 消息模型
(DocID = conversationID:docIndex,块内按 seq % 100 定位槽位)用另一种存储
达成同一性质,反向验证了"会话分块 + 序号定位 + 同键幂等"这一组合;
一期 Redis 实现的 zset "同 score 同 member 覆盖"也遵守同一契约(§18.1.2b)。
没有表级 default_time_to_live——这是与 user_mailbox_entry(§18.1.3)的关键差别:
message_record 的保留期由 retention_class 逐行决定(§7.1):
default 按租户策略,写入时按租户配置逐行设 USING TTL
ephemeral_24h 写入时 USING TTL 86400
compliance_hold **不设 TTL**,只能由 §21 的合规流程解除后删除
tenant_custom 按租户配置逐行设 TTL
因此表级 TTL 必须为 0,否则 compliance_hold 的行会被表级 TTL 静默删除,
直接违反 §21.5 的法务保留要求。
> **当前实现状态(2026-09-26)**:本表由 `HistoryArchiver` 异步批量写入(`ADR-0018`),
> writer 的提交热路径只写 Redis 热层(`history_hot_window`,7 天),不再同步写本表;归档为
> 普通 INSERT(先 `message_record` 后 `history_index`,无 LWT,同值幂等)。`ephemeral_24h`
> 按剩余期限 `USING TTL` 写入;`default` 按 `ADR-0019` 保留 30 天(可配),同样按剩余期限写
> `USING TTL`。`tenant_custom` 的租户策略来源未实现(失败闭合为不到期);一期不做对象存储冷归档
> 与 `PULL_HISTORY` 归档回读。
TWCS 窗口取 7 天(而不是邮箱表的 1 天):正文保留期以年计,1 天窗口会产生 数百个 SSTable;7 天窗口在"整文件过期丢弃"与"文件数可控"之间平衡。
2.2 MessageIndex¶
CREATE TABLE message_index (
tenant_id uuid,
message_id blob,
conversation_id uuid,
conversation_seq bigint,
seq_bucket bigint,
PRIMARY KEY ((tenant_id, message_id))
) WITH compaction = {'class':'LeveledCompactionStrategy'};
LCS 而非 TWCS:这张表是点查表、无时间局部性、行极小,LCS 的读放大最低。
它不在任何热路径上(§7.2):撤回、编辑、回复、举报在协议层强制携带
(conversation_id, conversation_seq),正常路径直接定位分区。本表只服务
Moderation/Admin 与故障排查。若线上出现该表 QPS 与消息量同数量级,说明有实现
绕过了协议层的坐标要求,属缺陷。
2.3 ClientDedup¶
CREATE TABLE client_dedup (
tenant_id uuid,
sender_id uuid,
client_message_id text,
message_id blob,
conversation_seq bigint,
state tinyint, -- RESERVED | COMMITTED(§8.3)
created_at bigint,
PRIMARY KEY ((tenant_id, sender_id, client_message_id))
) WITH default_time_to_live = 7200 -- client_dedup_ttl_seconds(ADR-0008)
AND compaction = {'class':'TimeWindowCompactionStrategy',
'compaction_window_unit':'HOURS',
'compaction_window_size':6}
AND gc_grace_seconds = 3600;
唯一使用 LWT 的表。占位用 INSERT ... IF NOT EXISTS,跨 ConversationWriter
实例的并发重试在此收敛(§8.3)。LWT 的 Paxos 开销只落在这一张小表的单行上,
不进入 message_record 与 conversation_head 的写路径。
gc_grace_seconds 取 1 小时而非默认 10 天:TTL 表的墓碑必须尽快回收,
且本表允许在极端修复场景下丢失(丢失只会让 24 小时内的重试产生重复消息,
由客户端 message_id 去重兜底,不影响正确性)。
2.4 ConversationHead¶
CREATE TABLE conversation_head (
tenant_id uuid,
conversation_id uuid,
latest_conversation_seq bigint,
earliest_available_conversation_seq bigint,
last_message_id blob,
last_activity_id blob,
last_sender_id uuid,
preview_or_placeholder blob, -- 用会话 DEK 加密(§21.3.1)
member_count int,
updated_at bigint,
fencing_epoch bigint,
head_version bigint,
PRIMARY KEY ((tenant_id, conversation_id))
) WITH compaction = {'class':'LeveledCompactionStrategy'};
LCS:单行反复覆盖写,LCS 能把同一分区的多个版本尽快合并,读放大最低。 STCS 会让一个热点会话的行散落在多个层级的 SSTable 里。
写入方式(§7.4 已定,此处给出 CQL 形态):
-- 正常路径:blind write,无 LWT
UPDATE conversation_head SET latest_conversation_seq=?, last_message_id=?, ...
WHERE tenant_id=? AND conversation_id=?;
-- 故障切换窗口内:条件更新
UPDATE conversation_head SET ...
WHERE tenant_id=? AND conversation_id=?
IF fencing_epoch < ? OR (fencing_epoch = ? AND head_version < ?);
latest_conversation_seq 只允许单调前进。实现必须在应用层比较后再写,
任何会造成回退的写入直接丢弃并计入告警——blind write 本身不提供该保证。
2.5 表情回应(§7.15)¶
CREATE TABLE message_reaction_summary (
tenant_id uuid,
conversation_id uuid,
seq_bucket bigint,
conversation_seq bigint,
counts map<text,int>,
summary_version bigint,
updated_at bigint,
PRIMARY KEY ((tenant_id, conversation_id, seq_bucket), conversation_seq)
) WITH CLUSTERING ORDER BY (conversation_seq DESC)
AND compaction = {'class':'LeveledCompactionStrategy'};
CREATE TABLE message_reaction (
tenant_id uuid,
conversation_id uuid,
conversation_seq bigint,
user_id uuid,
reaction_key text,
created_at bigint,
PRIMARY KEY ((tenant_id, conversation_id, conversation_seq), user_id)
) WITH compaction = {'class':'LeveledCompactionStrategy'};
message_reaction_summary 与 message_record 分区键完全相同(§7.15)。
这是有意的:§4.3 的分组批量读可以在同一次分区访问中把消息与其回应聚合一起取回,
回应功能在读路径上零额外往返。
counts 用 map<text,int>:回应种类是短枚举字符串、数量有限(几十个),
map 的整体覆盖写正好匹配 REACTION_UPDATE 的绝对值语义(§13.6.2)。
3. 压缩策略汇总¶
| 表 | 策略 | 理由 |
|---|---|---|
message_record |
TWCS 7 天窗口 | 按时间写入、按 TTL 整文件过期;避免逐行墓碑 |
message_index |
LCS | 点查、行小、无时间局部性 |
client_dedup |
TWCS 1 小时窗口 | 2 小时 TTL(ADR-0008),短窗口让整文件尽快过期 |
conversation_head |
LCS | 单行反复覆盖,需尽快合并版本 |
message_reaction_summary |
LCS | 同上,counts 反复覆盖 |
message_reaction |
LCS | 点查与分区扫描,无 TTL |
user_mailbox_entry |
TWCS 1 天窗口 | 见 §18.1.3,与本文表策略不同 |
统一约束:所有表禁止显式 DELETE(§18.1.3 的删除约束推广到 MessageStore)。
删除一律走两条路径之一:
到期删除 USING TTL 逐行设置,整 SSTable 过期后丢弃
合规删除 加密擦除(销毁 dek_id 指向的密钥,§21.3),行本体留待 TTL 自然回收
理由:TWCS 不做跨窗口 compaction,范围墓碑既不能提前释放空间,还会在读路径上
被反复扫描并击穿 §26.1 的 tombstones_scanned == 0 断言。
4. MessageStore 读路径¶
4.1 seq_bucket 分页算法¶
conversation_seq 允许空洞(§6.3:预留窗口崩溃、治理删除、可见性裁剪),
因此"扫到空桶"是正常现象,翻页必须有终止条件。
输入:conversation_id、direction、anchor_conversation_seq、limit
常量:W = message_seq_bucket_width(4096)
边界:H = ConversationHead.latest_conversation_seq
E = ConversationHead.earliest_available_conversation_seq
【最新一页】anchor 缺省时取 H
【older 方向】
b := anchor / W
empty_run := 0
loop:
rows := SELECT * FROM message_record
WHERE tenant=? AND conversation_id=? AND seq_bucket=b
AND conversation_seq < anchor
ORDER BY conversation_seq DESC LIMIT (limit - 已收集)
若 rows 为空 -> empty_run++;否则 empty_run := 0
已收集 += rows
若 已收集 >= limit -> 返回,has_more = true
若 b * W <= E -> 返回,has_more = false(触达保留边界)
若 empty_run > history_max_empty_bucket_scan(8, 附录 B.3)
-> 返回 has_more = false 并告警
b--; anchor := (b+1) * W
【newer 方向】同构,但 b++ 且需要 ORDER BY conversation_seq ASC
newer 方向是反向扫描(聚簇序是 DESC),在 Scylla 上有额外代价。
因此协议层的默认与热路径一律用 older:进入会话取最新一页 → 向上翻页,
newer 只用于"从某个锚点向下补齐"(如 mention_first_conversation_seq 跳转后回补)。
为什么必须有 empty_run 上限:一次预留窗口崩溃可能让 conversation_seq
跳过数千个值(§6.3 的预留窗口是 4096),恰好等于一个整桶。连续多次崩溃就会
留下连续空桶。没有上限时,翻页会退化为对空分区的连续扫描。
4.2 单分区上界的自我校验¶
seq_bucket 保证单分区 ≤ 4096 行(§7.0)。实现应在读路径埋一个断言:
单次分区读返回行数 > 4096 即为分桶实现错误(例如 W 被改过或 seq_bucket
计算用了浮点除法)。§27.2.4 已把 message_seq_bucket_width 列为建表后不可变项。
4.3 分组批量读(§18.3.1 的强制规则)¶
这是离线回填性能的决定性实现,也是 docs/01 之外与 §18.3.1 直接对应的落点。
输入:一批待 join 正文的条目(来自 MAILBOX_BATCH / PUSH_EVENTS 的组装)
1. 内联条目(ADR-0005,条目自带正文)直接跳过,不进入 join
2. 其余按 message_id 查本节点正文 LRU,命中的剔除
3. 未命中的按 (conversation_id, seq_bucket) 分组
4. 每组一次范围读,组间并发:
SELECT ... WHERE tenant=? AND conversation_id=? AND seq_bucket=?
AND conversation_seq IN (?, ?, ...)
同组内的 conversation_seq 连续时改用 range:
AND conversation_seq >= ? AND conversation_seq <= ?
5. 回填 LRU;同一 message_id 每批只 join 一次、只编码一次
**禁止逐 message_id 点查 message_record**(要点查只能走 message_index,
而它不在热路径上,见 §2.2)
代价对比(§18.3.1 已给量级):
1 个大群的 500 条积压 -> 1 次范围读 单条 join 成本 0.002
20 个单聊各 5 条 -> 20 次范围读 单条 join 成本 0.2
朴素实现(逐 message_id 点查)会把第一行也变成 500 次点查
IN 子句的元素数必须设上限(建议 100,超出拆多次请求)——Scylla 对超大 IN
会在协调者上放大内存与延迟。
4.4 回应聚合的同批取回¶
因为 message_reaction_summary 与 message_record 分区键相同,第 4 步可以
在同一分区的同一次访问中追加一条对 summary 表的范围读。实现上是两条 CQL 语句
但命中同一副本集与同一分区,协调者侧无额外跳数。
self_reaction_keys(请求者自己点过什么)不在 summary 表中,需要时按
(conversation_id, conversation_seq, user_id) 点查 message_reaction;
仅在用户查看单条消息详情时发起,不进批量路径。
5. 写路径¶
5.1 提交顺序(§8.1、§8.3)¶
1. ClientDedup 占位 INSERT ... IF NOT EXISTS,state=RESERVED
占位失败 -> 说明是重试,读出既有 message_id 并回放 SEND_ACK
2. 分配序号 同一临界区内分配 message_id / conversation_seq / last_activity_id
3. 写 MessageRecord USING TTL(按 retention_class,见 §2.1)
4. 追加 Outbox Redpanda 幂等 producer + acks=all(ADR-0006/0012)
5. ClientDedup 置位 state=COMMITTED
6. 回 SEND_ACK 仅表示 COMMITTED(§11.2),不表示任何接收者已收到
7. 异步更新 ConversationHead 失败重试并告警,**不阻断 fanout**(§7.4)
第 5 步与第 6 步之间崩溃:客户端重试命中 state=COMMITTED 的占位,直接回放
SEND_ACK,不产生第二条消息(§8.3 的崩溃点恢复表)。
5.2 mention_targets 与 media_metadata 的编码¶
两者都用 Protobuf 编码进 blob,与 docs/01 共用同一份 .proto:
mention_targets {
repeated bytes user_ids = 1; // 16 字节 uuid
bool at_all = 2;
}
media_metadata {
string object_id = 1;
string mime = 2;
uint64 bytes = 3;
uint32 width = 4;
uint32 height = 5;
uint32 duration_ms= 6;
bytes thumbnail = 7; // <= media_thumbnail_max_bytes(32 KiB)
string blurhash = 8;
bytes checksum = 9;
}
用 Protobuf 而不是 CQL 的 frozen<udt>:UDT 的 schema 演进需要
ALTER TYPE 并影响全表读路径,而 Protobuf 的字段增删由 §27.1.2 的规则覆盖,
且客户端与服务端共用同一定义,不需要在两处维护同构 schema。
thumbnail 内联进 media_metadata(不单独建表):它 ≤ 32 KiB,
且读消息时必然需要,单独建表会给每条媒体消息多一次读。原文件走对象存储,
message_record 只存 object_id(§3 禁令 7)。
5.3 撤回与编辑¶
固化顺序见 §13.4.1,本文只给存储侧动作:
撤回:1. 先写 CONTROL 事件到提交日志(分配自己的 conversation_seq)
2. 再 UPDATE message_record SET state=RECALLED,
recall_event_conversation_seq=<事件坐标>
WHERE tenant=? AND conversation_id=? AND seq_bucket=? AND conversation_seq=?
顺序不可颠倒:recall_event_conversation_seq 非空即表示"事件已入日志",
崩溃恢复据此判定是否需要补发事件(§13.4.4 的幂等条件)。
先改 state 后写日志会在崩溃窗口内永久丢失撤回事件。
编辑:同构,写 edited_at_activity_id,payload_or_ciphertext 整体替换。
正文不保留历史版本(产品未要求;若要求需另建版本表并重算容量)。
撤回不删除行(§3 的删除约束):正文清空由加密擦除或 TTL 承担,
state=RECALLED 只是渲染标志。
6. 加密字段¶
payload_or_ciphertext 与 ConversationHead.preview_or_placeholder 都用
会话 DEK 加密(§21.3.1)。AAD 绑定必须包含定位坐标,防止密文被搬到别处:
AAD = tenant_id || conversation_id || conversation_seq || dek_id.key_version
读路径解密点与 DEK 缓存见 §21.3.4;dek_cache_ttl 见附录 B.5.3。
服务端对 E2EE 密文(flags.ENCRYPTED)不做任何解密——那一层是端到端密钥,
与此处的静态加密 DEK 是两回事,不要混淆:
静态加密 DEK 服务端持有,用于合规擦除与静态数据保护,服务端可解密
E2EE 会话密钥 服务端不持有,§22.2 的 sender-key 模型,服务端只搬运
两者可叠加:E2EE 密文再经 DEK 加密后落盘
7. 容量核算(§25.1)¶
一期简化形态(ADR-0007,千人群、R_avg ≈ 22)下 MessageStore 侧的量级:
消息正文(式见 §25.1)
驻留 = M_day × b_row × lsm_space_amp × RF × 保留天数
b_row 含 Protobuf 编码后的 media_metadata 与列名开销
起步档代入(M_day = 2000 万,b_row = 600 B,space_amp 1.2,RF=3,365 天):
2e7 × 600 × 1.2 × 3 × 365 ≈ 15.8 TB
单分区上界自检:4096 行 × b_row ≈ 2.4 MB,远低于 Scylla 大分区告警阈值
message_index 与 client_dedup 的容量相对正文可忽略(前者行 ≈ 60 B 且与消息
同数量级,后者 2 小时 TTL,ADR-0008)。回应两表的容量取决于产品使用率,
压测前不进容量结论(§25.6)。
8. 验收¶
S-1【发布阻断】分组批量读
构造两种回填分布(1 个大群 500 条 / 20 个单聊各 5 条),清空正文 LRU:
(a) 对 message_record 的读请求次数 <= 2
(b) 读请求次数 <= 21
两种分布下"逐 message_id 点查 message_record"次数 == 0
message_index 的 QPS == 0(正常路径不得触碰)
—— 与 §26 用例 26.1.5 同源,本条断言落到 CQL 层
S-2【发布阻断】有符号边界
构造 conversation_seq 与 mailbox_seq 跨越 2^62 的写入:
写入前断言触发(bit63 检查),不产生排序翻转的数据
S-3 空洞翻页终止
构造 conversation_seq 从 1 跳到 100000 的会话(连续 24 个空桶):
PULL_HISTORY 扫描的空桶数 <= history_max_empty_bucket_scan(8)
返回 has_more=false 并计入告警;不出现无限扫描
S-4 单分区上界
单个 seq_bucket 的返回行数恒 <= 4096;超出即断言失败
S-5 compliance_hold 不被 TTL 删除
写入 retention_class=compliance_hold 的消息,等待超过租户默认 TTL:
行仍可读;表级 default_time_to_live == 0 的 schema 断言
S-6 墓碑
正常运行 7 天后对 message_record 做全表采样读:
tombstones_scanned == 0(禁止显式 DELETE 的验证)
S-7 撤回顺序
在"写日志后、改 state 前"注入崩溃:
恢复后 recall_event_conversation_seq 为空 -> 补发事件;
撤回事件最终可见;丢失撤回事件次数 == 0
9. 待办¶
[x] .proto 落地:media_metadata / mention_targets 与 docs/01 共用同一文件(`common.proto`)
[ ] 租户 TTL 配置到 retention_class 的映射表(§21.1 的分层策略输入)
[ ] earliest_available_conversation_seq 的推进器实现(保留期裁剪与治理删除时更新)
[ ] IN 子句元素数上限(建议 100)的实测校准
[ ] b_row 实测回填(附录 B.7),当前 600 B 为估算值