跳转至

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。

启动与恢复

  1. 读取全部 checkpoint 和元数据,随后读取所有 dispatch partition 的 broker low/high watermarks,形成本次启动高水位快照 H0[p]。
  2. 对每个 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} 内容和活跃时间单调, 没有重复或缺项。