我们的内部 Agent 能力平台在上线前夜发现了一个尴尬的事实:所有对话历史都只在 Redis 里,TTL 24 小时,没开持久化,没挂 volume。 Redis 一重启,所有对话记录灰飞烟灭。
这篇文章记录我们如何把对话消息从"Redis 里的易失 JSON"迁移到"ClickHouse 主存储 + Redis 热缓存"的架构,以及过程中的几个关键决策:为什么不是 MySQL、为什么是 ClickHouse、表结构怎么设计、为什么 OLAP 场景下字段冗余胜过 JOIN、以及 ClickHouse 不能用 Flyway 之后我们怎么做 schema 管理。
1. 问题:消息数据的四个坏消息
改造前的现状:
会话元数据(session 表) → MySQL(✅ 持久化)
会话消息(messages[]) → Redis(⚠️ 无持久化,TTL 24h)
消息以整个会话为单位存成一个 JSON blob(SessionMemory),key 为 memory:{sessionId}。四个问题:
| 问题 | 影响 |
|---|---|
| Redis 未开 AOF/RDB、未挂 volume | 重启即所有对话历史永久丢失 |
| TTL 24h | 超过一天的会话消息自动被清除 |
| 改为永不过期? | Redis 内存持续膨胀,用内存当主存是烧钱 |
| 消息不在分析链路 | 无法做 Agent 效果评估、调用量统计、token 分析等任何 OLAP 查询 |
1.1 选型:为什么不是别的方案
| 方案 | 否决原因 |
|---|---|
| MySQL 单表 | 消息量大且只追加不修改,MySQL 行存持续膨胀,早晚触发分库分表;且后续分析查询会拖垮业务库 |
| 对象存储(MinIO/S3) | 团队不愿为消息再维护一套文件存储链路;按 session 读取整条历史也不方便 |
| MongoDB | 引入全新组件,团队要额外维护一套集群 |
| ClickHouse | ✅ 公司已有部署,零新增组件;架构规划里本就定位为分析库;append-only + 列存,与消息数据形态天然吻合 |
核心判断就一句话:对话消息是"写多、只追加、不改、按会话范围读、按维度聚合分析"的数据——这是 OLAP 列存的教科书场景,不是 OLTP 的。
1.2 设计目标
- 消息永久持久化,不受 Redis TTL / 重启影响;
- 零新组件,复用现有 ClickHouse;
- 读取性能不低于原 Redis 方案(热数据仍有缓存);
- 写入不阻塞推理主流程;
- 为分析预留——支持 Agent 效果评估、用量统计等 OLAP 查询。
2. 架构:ClickHouse 主存储 + Redis 热缓存
写入链路
─────────────────
chat() 流式完成
│
├──① INSERT chat_message (ClickHouse,批量)
│
└──② SET memory:{sid} (Redis 热缓存,TTL 1h)
读取链路
─────────────────
loadMessages(sessionId)
│
├──① GET memory:{sid}
│ └─ 命中 → 直接返回(< 1ms)
│
└──② miss → ClickHouse SELECT
WHERE session_id = ? ORDER BY sequence
│
└─ 回填 Redis → 返回
职责划分:
- Redis:活跃会话的热缓存,TTL 1 小时。命中率高,响应亚毫秒。不再是主存储,可丢。
- ClickHouse:所有消息的权威持久化副本。
这是一个标准的 cache-aside 模式,只是"数据库"那一侧换成了 ClickHouse。冷启动(缓存 miss)时多一次 ClickHouse 范围扫描,按下面的表结构设计,这个查询是毫秒级的。
3. 表结构设计:每个决策都有出处
CREATE TABLE chat_message (
message_id UInt64,
session_id String, -- 会话 ID
role LowCardinality(FixedString(9)), -- 'user' | 'assistant' | 'system'
content String, -- 消息正文
sequence UInt16, -- 会话内序号,保证顺序
created_at DateTime,
space_id String, -- 冗余:跨会话分析
instance_id String, -- 冗余:按 Agent 聚合
user_id Int64, -- 冗余:按内部用户聚合
end_user_id String, -- 冗余:终端主账号
end_sub_user_id String -- 冗余:终端子账号
) ENGINE = MergeTree()
PARTITION BY toYYYYMM(created_at)
ORDER BY (session_id, sequence)
TTL created_at + INTERVAL 180 DAY DELETE;
| 决策 | 理由 |
|---|---|
ORDER BY (session_id, sequence) |
读取整场会话 = 一次范围扫描,无需二次排序。这是最高频查询,排序键为它服务 |
PARTITION BY toYYYYMM |
按月分区。TTL 到期自动整分区删除,清理成本为零(不需要 DELETE mutation) |
LowCardinality(FixedString(9)) on role |
只有 3 个枚举值,字典编码后压缩近似 tinyint |
冗余 space_id / instance_id / user_id / end_user_id |
见下一节——ClickHouse 的黄金法则:宽表冗余,绝不 JOIN |
TTL 180 DAY |
控制存储成本,可按合规要求调整 |
sequence UInt16 |
会话内排序的稳定依据;同一 session 内单调递增 |
3.1 思辨:为什么宁可冗余也不 JOIN
表设计评审时被挑战过:end_user_id 明明可以通过 session_id → 会话元数据表(MySQL) 查到,为什么还要在 chat_message 里冗余一列?
第一层:常规聊天展示确实不需要冗余。 加载消息按 session_id 查,终端用户信息可以从 MySQL 的会话表拿。只看这个场景,冗余是多余的。
第二层:OLAP 分析强依赖冗余。 ClickHouse 的定位是分析库,核心价值在于不做 JOIN 的聚合查询:
-- 场景 A:哪个终端客户的会话量最大?
SELECT end_user_id, count(DISTINCT session_id)
FROM chat_message GROUP BY end_user_id;
-- 场景 B:某客户最常被问到的问题 TOP 10
SELECT end_user_id, substring(content, 1, 80), count()
FROM chat_message
WHERE role = 'user'
GROUP BY end_user_id, content
ORDER BY count() DESC LIMIT 10;
-- 场景 C:各终端子账号的回复质量侧面指标
SELECT end_sub_user_id, avg(length(content))
FROM chat_message
WHERE role = 'assistant'
GROUP BY end_sub_user_id;
如果 end_user_id 只存在 MySQL 里,上面每个查询都要跨两个异构数据库做 JOIN——应用层先查 MySQL 拿映射、再在 ClickHouse 里过滤,列存优势荡然无存,代码复杂度翻倍。
结论:ClickHouse 的最佳实践是宽表冗余——宁可多存字段,也不跨库 JOIN。 顺带一个一致性论证:user_id 同样是冗余的(也能通过会话表查到),没理由保留 user_id 却删掉 end_user_id。要么都冗余,要么都别冗余——而"都不冗余"在分析场景下不可行。
4. 读写流程与写入策略
4.1 写入:流式完成后批量插入
推理是 SSE 流式的,消息落库挂在流完成的回调(doOnComplete)里:
1. 收集器聚合完整的 assistant 回复
2. 构建两条消息:
- {session_id, "user", userMessage, seq=N, ts}
- {session_id, "assistant", assistantContent, seq=N+1, ts}
3. ChatMessageGateway.insertBatch(...) → ClickHouse
4. 同步更新 Redis 热缓存(保证读一致性)
4.2 当前:同步 fire-and-forget;优化方向:异步批量
当前实现是同步写入但不阻塞:doOnComplete 里直接调用 insertBatch(),失败只记 error 日志 + 发布工作日志事件,不影响已完成的响应。用户已经拿到完整回复,落库失败只是数据问题,不是可用性问题。
后续量上来之后,可以演进为内存队列 + 定时批量 flush,进一步降低 ClickHouse 连接开销:
chat() 完成
│
├── 同步: 更新 Redis 缓存
└── 异步: 消息追加到内存队列
│
▼
ChatMessageWriter (@Scheduled 1s)
│
├── 攒够 100 条 → flush
└── 超过 5s → flush
│
▼
ClickHouse PreparedStatement.addBatch() 批量写入
ClickHouse 的官方建议本来就是大批量、低频次写入(每次几千上万行),避免高频小批次产生大量 part 触发合并压力。对话消息场景单轮只有 2 行,异步攒批是早晚要做的优化;但 MVP 阶段同步 fire-and-forget 足够,先把持久化正确性拿到手。
4.3 读取与缓存失效
loadMessages(sessionId):
1. GET memory:{sid}
2. 命中 → 解析 JSON → 返回
3. miss → ChatMessageGateway.selectBySession(sid)
→ 构建 SessionMemory → SET Redis(回填)→ 返回
缓存失效:新建消息 / 删除会话 / 更新会话 → DEL memory:{sid}
5. Schema 管理:ClickHouse 用不了 Flyway 怎么办
一个容易踩的坑:ClickHouse 不支持事务,而 Flyway 的 flyway_schema_history 记录 + 迁移执行依赖事务保证原子性——把 Flyway 直接怼到 ClickHouse 上无法可靠工作。
我们采用了更简单的方案:@Component 实现 ApplicationRunner,应用启动时执行幂等建表脚本:
ClickHouseSchemaInitializer implements ApplicationRunner
│
│ 注入 JDBC 连接配置(独立的连接,非 Spring DataSource Bean)
│
└── run() →
CREATE TABLE IF NOT EXISTS chat_message (...)
ALTER TABLE chat_message ADD COLUMN IF NOT EXISTS xxx -- 兼容存量表
↑
幂等:表存在则跳过,重复执行安全
与 Flyway 的对比:
| 维度 | MySQL(Flyway) | ClickHouse(ApplicationRunner) |
|---|---|---|
| 执行时机 | Bean 初始化阶段 | ApplicationRunner,所有 Bean 就绪后 |
| 幂等保证 | flyway_schema_history 表 |
SQL 自带 IF NOT EXISTS |
| 并发安全 | 数据库事务锁 | 无事务,但 IF NOT EXISTS 本身是原子操作 |
| 变更追溯 | history 表 + checksum | Git 中保留版本化 SQL 文件,Initializer 按序执行 |
后续表结构变更(加列、改 TTL)遵循同样方式:新建版本号命名的 SQL 文件放进 db/migration-clickhouse/ 目录,由 Initializer 按顺序执行。牺牲了一点工具链的严谨(没有 checksum 校验),换来的是零额外依赖和完全可控的启动流程——对小团队是划算的交换。
6. 分层落地:ClickHouse 只是一个 Gateway 的实现细节
按项目的 COLA 约束,Domain 定义 SPI,Infrastructure 实现:
domain/memory/
├── ChatMessage.java # 实体(纯 Java,零存储依赖)
└── ChatMessageGateway.java # SPI 接口
infrastructure/clickhouse/
├── ClickHouseConfig.java # JDBC 连接配置
├── ClickHouseSchemaInitializer.java # 启动建表
└── ChatMessageGatewayImpl.java # SPI 实现
SPI 接口只有两个方法:
package com.example.platform.domain.memory;
public interface ChatMessageGateway {
/** 批量写入消息 */
void insertBatch(List<ChatMessage> messages);
/** 按 sessionId 查询全部消息,按 sequence 升序 */
List<ChatMessage> selectBySession(String sessionId);
}
上层的调用关系值得注意:
ReasoningAppService只依赖MemoryGateway,不感知 ClickHouse 的存在;MemoryGatewayRedisImpl.loadSession()在 Redis miss 时通过ChatMessageGateway.selectBySession()回退到 ClickHouse——缓存与主存的协作被封装在 infra 层内部。
收益:如果哪天 ClickHouse 换成了别的(比如数据量真的大到需要分布式方案,或者公司基础设施变了),改动范围被严格限制在 infra 层的一个目录里,接口签名纹丝不动。存储选型是局部决策,不是架构决策——这正是模块化单体 + SPI 分层想买到的东西。
7. 迁移与成本
历史数据:不迁移。 改造前 Redis 里的消息本来就带 24h TTL,会自然过期;当时还是内部 MVP,没有生产数据;上线后新消息直接写 ClickHouse。不给自己加一场没有收益的迁移工程。
存储成本估算(中等规模):
| 参数 | 值 |
|---|---|
| 日均活跃会话 | 200 |
| 每会话消息 | 15 轮(30 条) |
| 每消息平均 | 500 字符(~1KB) |
| 日增量 | 200 × 30 × 1KB ≈ 6 MB |
| 180 天存量(TTL 上限) | ~1 GB |
| 实际磁盘占用(列存 3~8× 压缩) | 200~400 MB |
一年攒不下 1GB 磁盘。这就是"用对存储"的红利:同样的数据放 MySQL,一年后你要面对一张几亿行的 InnoDB 表和一堆慢查询;放 ClickHouse,它甚至还没热身。
8. 全景:四种存储各司其职
改造后平台的存储矩阵:
| 数据 | 存储 | 用途 |
|---|---|---|
| 业务数据 | MySQL | 工作空间、Agent 配置、知识库元数据、工具、画像 |
| 会话元数据 | MySQL | session 表 |
| 会话消息 | ClickHouse | chat_message 表(本次新增) |
| 消息热缓存 | Redis | 活跃会话加速(TTL 1h,可丢) |
| 向量嵌入 | Milvus | 知识库语义检索 |
| 工作日志 | MySQL → 后续汇入 ClickHouse | 操作审计与分析 |
没有为消息引入任何新组件,MySQL 不需要扩容,Redis 回归到它最擅长的"热缓存"角色,ClickHouse 从"规划中的分析库"变成了"分析 + 高吞吐 append-only 存储"双职责。
9. 经验总结
- 数据形态决定存储,而不是"系统里已有哪个库"。 对话消息是 append-only + 范围读 + 聚合分析,这是列存 OLAP 的主场。塞进 MySQL 只是把分库分表的账单延后。
- OLAP 场景里冗余胜过 JOIN,跨异构库 JOIN 是红线。 评审时"这个字段能 JOIN 出来"不是删列的充分理由——先问分析查询要不要它。
- ClickHouse 写入要攒批,但持久化正确性优先。 MVP 阶段同步 fire-and-forget 是对的:先把"重启不丢数据"拿到,再优化连接开销。
- Flyway 不是银弹,幂等 DDL + ApplicationRunner 对无事务的存储更实在。
- 把存储藏在 SPI 后面。 上层只看见
ChatMessageGateway,ClickHouse 是 infra 层的实现细节——存储选型的推翻成本因此趋近于零(这次改造本身就是证据:从 Redis-only 到 ClickHouse+Redis,业务层代码几乎没动)。