Durable Session Projection 增量设计¶
目标与边界¶
将 SessionProjection 的输入从 MailboxNode 到 qsession 的易失 TCP Projection 推送,改为
qsession 自己消费已提交的 dispatch log。这样 qsession 停机、出站队列满或 Redis 短暂故障
只能造成投影滞后,不能造成会话集合永久缺项,也不阻塞 mailbox 的消息可达性。
本增量只定义 qsession 的后续实现;不改变 canonical MessageRecord、个人邮箱、MailboxNode
水位、fanout 事务或客户端协议。qsession 不提交 Kafka consumer-group offset,也不复用
MailboxNode 的 DispatchProgress / W / W_floor。
复用边界¶
直接复用 qim-commit-log::RdkafkaDispatchLogConsumer 与
RdkafkaDispatchLogConsumerConfig:它已提供静态 partition assign、显式每 partition 起点、
auto.offset.reset=error 和 low/high watermark readiness 校验。qsession 用独立 client/group
identity 创建自己的 reader;group id 仅是 librdkafka client 身份,正确性位点始终是下述
Redis checkpoint。
无需为 qsession 创建第二套 Kafka seek/readiness 实现。后续若需要提取公共接口,只能从现有 dispatch consumer 的“给定 partition 起点并验证保留范围”能力抽象,不能把 qsession 的 checkpoint 语义塞入 MailboxNode API。
Checkpoint 存储契约¶
Redis 使用一个带版本的元数据 hash 和按 partition 独立的 checkpoint hash;当前 standalone Redis 前提下由同一连接管理,任何不一致均 fail closed。
session_projection:meta
format_version = 1
dispatch_topic = 配置的 topic 精确字符串
partition_count = broker / 部署预期的精确分区数
cluster_id = broker `fetch_cluster_id` 的精确字符串
topic_id = broker DescribeTopics 返回的 immutable UUID
session_projection:checkpoint:{dispatch_topic}:{partition}
format_version = 1
dispatch_topic = 同 meta 的精确字符串
partition_count = 同 meta 的精确数值
cluster_id = 同 meta 的精确字符串
topic_id = 同 meta 的 immutable UUID
partition = 当前 partition
next_offset = 已完整投影的下一条 Kafka offset
next_offset 是下一条待处理记录,不是最后已处理 offset。它必须只在该 dispatch 的所有
recipient ZADD convs:{user} GT created_at conversation_id 已成功后推进到
record.position.offset + 1。
元数据、每个 checkpoint 的格式版本、topic、partition count、cluster id、topic id 和 partition 自身必须完全匹配。 任一字段缺失、无法解析、版本过新、cluster 切换、topic 改名、partition count 改变或同名 topic 重建均拒绝 启动;禁止根据当前配置重写旧 checkpoint,也禁止把不匹配状态当作 offset 0。
启动与恢复¶
- 读取全部 checkpoint 和元数据,随后读取所有 dispatch partition 的 broker low/high
watermarks,形成本次启动高水位快照
H0[p]。 - 对每个 partition 计算显式起点:
有 checkpoint: low[p] <= next_offset[p] <= high[p],否则拒绝启动
无 checkpoint: 仅 low[p] == 0 时可从 0 bootstrap;low[p] > 0 必须拒绝启动
“无 checkpoint”包括 checkpoint 表被删除。它在 low>0 时不能证明投影过历史记录,
所以不能用 earliest、当前 high 或配置默认值恢复。
3. 用上述起点创建现有 RdkafkaDispatchLogConsumer,并执行其 readiness。消费者只读
read_committed dispatch。
4. 逐 partition 保持顺序,物化到 H0[p]。当每个 partition 的 checkpoint 均已到达
H0[p],才开放 PULL_SESSION_LIST / session readiness。运行期间新增的 offset 不改变
H0;它们由同一持续消费循环继续物化。
5. 一条 dispatch 的 recipient 写入失败、checkpoint 写入失败或消费者读错误均使 qsession
变为 not-ready/dirty 并停止对外服务;不能跳到下一 offset。进程重启从最后 checkpoint
重放。
由于 ZADD ... GT 对相同 (user, conversation, created_at) 幂等,网络失败导致的“部分
recipient 已写、checkpoint 未推进”只会重放,不能漏会话。实现可先校验所有 pipeline 结果,
再写 checkpoint;不得在任何失败或未确认的批次后推进 checkpoint。
与现有易失投影的迁移¶
切换版本先部署新 checkpoint reader 并完成从 low=0 的首次投影;在它完成全分区 catch-up
前,qsession readiness 不得为绿。验证会话集合和 canonical history 一致后,删除
MailboxNode 的 session_tx 易失投影出口。迁移期间不得让两条路径竞争推进同一 checkpoint;
若临时双写,只能允许旧 TCP 路径做无 checkpoint 的冗余 ZADD GT。
topic 重建、partition count 扩容/缩容、checkpoint schema 升级不是在线兼容事件。操作流程必须 显式创建新的 format version / checkpoint namespace,并在完整 dispatch 历史可读时重建;否则 停止服务并人工恢复。不得自动清空旧 checkpoint 伪装成首次启动。
Broker 能力前提¶
本设计要求 broker 通过标准 Kafka DescribeTopics 返回 KIP-516 的非零 immutable topic UUID,
并允许客户端读取 fetch_cluster_id、topic metadata、watermark 和 retention.ms。任一 identity
字段缺失、全零或在启动的 inspect/readiness 前后变化,qsession 均 fail closed。不得用人工环境变量
generation、topic 名或 Redpanda 私有 Admin API 替代此证明。Redpanda v24.3.6 的 metadata 不提供
topic id,不能承载 durable projection;部署的最低能力应为实现该标准字段的 broker(当前 Redpanda
基线为 v26.2.2)。
必测门禁¶
- qsession 停机或 Redis pipeline 失败后重启:未完成 dispatch 重放,所有 recipients 的会话 均出现,且 checkpoint 不越过失败记录。
- checkpoint 已存在且位于
[low, high]:从next_offset恢复;等于 high 的已追平分区可 直接通过。 - checkpoint 缺失/被删除且
low=0:允许 offset 0 bootstrap;缺失且low>0:启动失败。 - checkpoint 小于 low、大于 high、topic identity 不同、partition count 不同、format version 不同或字段损坏:全部启动失败,不执行自动 reset。
- 支持 KIP-516 的真实 broker:同名单分区 topic 写入旧 checkpoint 后删除重建,新 topic 的 high
大于旧 checkpoint 且旧 checkpoint 仍落在
[low, high]时,immutable topic id 改变并拒绝启动。 - 启动后记录
H0,在 catch-up 过程中追加新 dispatch:所有 partition 达到H0后 readiness 变绿,新增记录继续被消费但不使启动目标无限后移。 - dispatch 重放 100 次及“部分 recipient 已写后断开”:
convs:{user}内容和活跃时间单调, 没有重复或缺项。