AI应用对话数据层设计:从流式落库到多会话隔离
2026/9/8 16:35:33 网站建设 项目流程

做 AI 应用的人都有一个共同体验:第一版 Demo 跑通特别快,把大模型接口一调、前端一个流式渲染加上去,看起来就"能用了"。但一旦想把它变成真正能落地的产品,比如一个支持多会话、流式输出的 AI 机器人套件,你会发现最花时间的根本不是调模型接口,而是怎么把对话数据这一层组织好。

TinyRobot Kit 就是在这个背景下折腾出来的东西。它不是又一个 LLM 封装库,而是一个偏"数据层"的对话机器人套件,核心就解决三件事:流式消息怎么可靠落库、多会话之间怎么隔离与切换、模型需要的上下文档怎么从库里高效组装。这篇文章按我实际设计和迭代这套 Kit 的过程,把数据层的组织思路、核心表结构、流式写入方案、多会话并发处理、以及我在生产环境踩过的坑完整复盘一遍。正在做对话类 AI 应用的开发者,尤其是从"调通接口"往"做产品"过渡阶段的同学,应该能从里面找到不少可以直接抄走的方案。

1. 对话数据层:AI 应用里最容易被低估的一层

1.1 从"调用模型"到"组织对话"的转折点

大多数 AI 项目的演进路径是高度相似的:先调通一个 Chat Completion 接口,用流式把 token 吐到前端,做一个"打字机"效果,Demo 就成立了。但 Demo 和产品之间隔着一整层东西——数据。

普通 CRUD 应用的数据是结构化的、静止的:用户表、订单表、商品表,一条记录就是一个完整的事实。对话数据完全不是这样。一条消息从生成到最终定稿,中间经历了"正在生成、生成了一部分、生成完毕、生成失败"等多个状态;这些状态还是高频率变化的,模型每秒钟可能吐出来几十个 token;聊到第三十轮的时候,用户想回看第一轮说了什么,数据还得按原样还原出来。

TinyRobot Kit 要做的第一件事,就是把"对话"这个概念从内存里搬到持久化存储里。模型调用是可以无状态的,但产品不能无状态。用户关掉页面再打开,会话还在;网络断了一下,已经生成的内容不能丢;两个会话轮流切换,每个会话的上下文不能串。

1.2 数据层必须回答的三个问题

我在设计 Kit 的数据层时,把所有需求收敛成了三个问题:

  • 一条消息的一生是怎样的?从用户按下发送,到模型返回完整内容,这条消息经历了哪些阶段?每个阶段的数据要存成什么样?
  • 多个会话如何并存?一个用户可能同时开着五个会话,后台可能还有一个定时任务在往会话里写数据。如何保证互不干扰、切换无损?
  • 模型的上下文从哪里来?每次请求模型前,Kit 都要把历史消息取出来拼成 prompt。这块取数逻辑如果做得糙,轮数一多就会出大问题。

这三个问题分别对应了消息生命周期管理、会话隔离与并发控制、上下文组装策略。下面的章节就按这条线展开。

2. TinyRobot Kit 的会话模型:Session 与 Message 的边界到底怎么划

2.1 为什么不能只存一张消息表

最早我的想法特别朴素:一张 message 表,字段是 id、会话 id、角色、内容、时间,完事了。但用起来很快就发现不够。最直接的问题是——"会话"本身也是需要存东西的。

每个会话都有标题、创建时间、最后活跃时间、关联的用户、甚至还有模型参数(temperature、max_tokens)。这些字段如果冗余到每条消息里,浪费且容易不一致;如果完全不带,会话列表页就没法做了。所以第一步就是把数据拆成两层:sessionsmessages

sessions表存的是会话的"骨架",它回答的是"有哪些会话、每个会话属于谁、处于什么状态":

CREATE TABLE sessions ( id TEXT PRIMARY KEY, user_id TEXT NOT NULL, title TEXT DEFAULT '新会话', status TEXT NOT NULL DEFAULT 'active', -- active / archived / deleted system_prompt TEXT, model_config JSONB, -- temperature、max_tokens 等参数 last_message_at TIMESTAMPTZ, message_count INTEGER DEFAULT 0, created_at TIMESTAMPTZ NOT NULL DEFAULT now(), updated_at TIMESTAMPTZ NOT NULL DEFAULT now() ); CREATE INDEX idx_sessions_user_time ON sessions (user_id, updated_at DESC);

messages表存的是对话的"血肉",回答的是"这个会话里聊了什么、每条消息现在怎么样了":

CREATE TABLE messages ( id TEXT PRIMARY KEY, session_id TEXT NOT NULL REFERENCES sessions(id), role TEXT NOT NULL, -- system / user / assistant / tool content TEXT, -- 消息最终内容 content_delta TEXT, -- 流式增量写入的临时区域 status TEXT NOT NULL DEFAULT 'pending', -- pending / streaming / completed / failed seq INTEGER NOT NULL, -- 会话内递增序号 token_count INTEGER DEFAULT 0, parent_id TEXT, -- 用于将来支持分支会话 meta JSONB, -- 扩展字段,比如工具调用参数 created_at TIMESTAMPTZ NOT NULL DEFAULT now(), updated_at TIMESTAMPTZ NOT NULL DEFAULT now() ); CREATE INDEX idx_messages_session_seq ON messages (session_id, seq);

2.2 消息状态机:一条消息的一生

如果只需要存"最终结论",content一个字段就够了。但要支持流式,你得先把"生成中"这个状态显式建模出来。

我把消息的生命周期定义成了四个状态:

  • pending(等待中):请求刚创建,消息记录已插入,但模型还没有返回任何内容。
  • streaming(生成中):模型已经开始返回增量,content_delta不断追加。
  • completed(完成):流结束,content_delta合并进content,token 计数写入。
  • failed(失败):请求中断或异常,保留已生成的部分,打上失败标记。

这个状态机是整个数据层的核心。它的价值在于:任何时刻,你都能准确回答"这条消息现在到底怎么样了"。前端要显示 loading 还是完整内容、后台要决定是否重试、日志要记录失败原因,全靠它。

还有一个容易被忽略的点:seq字段。同一会话内的消息如果按created_at排序,在并发写入、时钟回拨、批量导入的场景下顺序是不可靠的。用session_id + seq作为排序键和唯一键,才能保证顺序的确定性。我把seq设计成会话内的单调递增序列,插入新消息时用max(seq) + 1计算,宁可多一条 SQL 查询,也不在排序上埋雷。

2.3 会话元数据与上下文的关系

设计完表结构,我意识到一个更深的问题:会话不只是消息的容器,它还隐含了"这个对话的背景"。

同一个 session 里,system_prompt是全局设定的、history是需要动态截取的、model_config是固定附带的。这三类数据如果散落各处,组装上下文的时候就要到处拼。所以我把system_promptmodel_config直接冗余进了sessions表,而不是做成关联表。

冗余在这里不是坏味道,而是刻意的取舍。理由是:会话表的变化频率极低,但读取频率极高——每次请求模型、每次打开会话列表都要读。与其每次都 JOIN 一张配置表,不如在写入会话时一次性冗余好,读取时一条 SQL 全拿出来。这个取舍在规模变大之后收益非常明显。

3. 流式消息落库:打字机效果的每一帧都不能丢

3.1 流式过程的数据视角拆解

大模型的流式返回,从数据视角看是这样一个过程:请求发出后,服务端在几秒到几十秒的时间里,陆续吐出一批又一批的文本增量。每批增量都是一个小的文本块,把它们按顺序拼接起来,才是一条完整的 assistant 消息。

实现流式落库,最笨的办法是每收到一个增量就 UPDATE 一次数据库。在模型返回慢、增量频率低的时候,这样做没问题。但现在的模型返回速度越来越快,某些情况下每秒钟会有几十帧数据过来。如果每一帧都触发一次 SQL UPDATE,数据库连接会被迅速打满,而且大量写入是多余的——用户根本不在乎中间那几十个中间态。

TinyRobot Kit 采用的方案是:前端走 SSE 实时收每一帧,后端落库走批量合并。SSE(Server-Sent Events)负责把增量即时推给前端渲染,保证"打字机"效果;数据库落库则用一个增量缓冲区,按时间窗口批量合并写入。

3.2 增量缓冲区的实现

我定义了一个流式写入器 StreamingMessageWriter,它的职责只有一个:接收增量,合并写入,控制落库频率。

class StreamingMessageWriter: def __init__(self, session_id: str, message_id: str, repo, flush_interval=2.0): self.session_id = session_id self.message_id = message_id self.repo = repo self.buffer = [] self.flush_interval = flush_interval self._last_flush = time.time() async def append(self, delta: str): self.buffer.append(delta) # 缓冲达到阈值,立即落库 if len("".join(self.buffer)) >= 512: await self.flush() async def flush(self): if not self.buffer: return text = "".join(self.buffer) self.buffer.clear() await self.repo.append_delta(self.session_id, self.message_id, text) self._last_flush = time.time()

这里有两个触发落库的条件:一是缓冲区文本达到 512 字符,二是定时器到了 2 秒。前者保证了大段文本不会长期滞留内存,后者保证了小碎块增量也会被周期性地持久化。也就是说不论增量大小,最迟 2 秒内的数据一定已经进库。

为什么是 512 和 2 秒这两个数?没有绝对标准,是我在实测中调出来的。设太小,比如 128 字符,写入频率还是偏高;设太大,比如 2048,在模型输出慢的场景下,用户关页面时容易丢失较多数据。2 秒的定时刷新,配合流结束时的强制 flush,可以做到最多丢 2 秒的增量——这个损失在"已生成内容"的视角下是完全可以接受的。

3.3 流结束时的定稿操作

流式返回结束之后,不能直接拍屁股走人。缓冲区里可能还有没落库的增量,消息状态还是 streaming,token 计数还没更新。我专门做了一个 finalize 流程:

async def finalize(self, token_count: int, extra_meta: dict = None): await self.flush() content = await self.repo.get_full_content(self.message_id) await self.repo.complete_message( message_id=self.message_id, content=content, token_count=token_count, meta=extra_meta, ) # 顺带更新会话的 last_message_at 和 message_count await self.repo.bump_session(self.session_id)

注意complete_message这一步,它做的不只是 UPDATE 状态,还把之前累积在content_delta里的所有碎片合并成完整文本写入content。我选择用一次读再写的方式,而不是在内存里维护完整内容,原因是:在分布式部署下,负责流式接收的进程和负责定稿的进程可能不是同一个,从库里拿全量最稳妥。

3.4 中断与恢复:让残片数据也能自洽

流式请求最讨厌的地方在于,它随时可能断。网络抖动、用户关页面、模型服务超时,随便一种情况都会让消息停在"生成了一半"的状态。

TinyRobot Kit 的处理策略是这样的:

  • 流中断时,写入器先尝试执行一次 flush,把已收到的增量落库。
  • 消息状态置为 failed,但保留content_delta里已有的部分文本。
  • 读取会话时,遇到 failed 的 assistant 消息,前端可以展示"(已中断)"以及已经生成的内容,而不是整个消息消失。
  • 重试时插入一条新的 assistant 消息,而不是覆盖失败的旧消息,避免并发重试互相污染。

这个策略的本质是:不要试图让失败的数据变得完美,而是让失败本身可观察、可恢复。用户和管理员都能清楚地看到哪条消息是失败的、失败前生成了什么,这比默默吞掉错误好得多。

4. 多会话并发:隔离、切换与上下文窗口管理

4.1 会话隔离:按 session_id 做分区

多会话的核心矛盾在于:一个进程里同时跑着多个用户、多个会话的流式请求,数据绝不能串。我见过不少新手项目,把所有消息放在一个全局数组里,session 用内存里的一个变量标记,一换会话就乱了。

TinyRobot Kit 的隔离策略简单但有效:一切数据操作都以session_id为第一维度。SQL 查询强制带 session_id,缓存 key 强制带 session_id,内存中的流式写入器实例也按 session_id 分桶存放。

class SessionRepository: def list_messages(self, session_id: str, after_seq: int = 0, limit: int = 50): return self.db.query( "SELECT * FROM messages WHERE session_id = ? AND seq > ?" " ORDER BY seq LIMIT ?", session_id, after_seq, limit, ) def get_context_messages(self, session_id: str, window_size: int): # 只取该 session 内的消息,天然隔离 return self.db.query( "SELECT * FROM messages WHERE session_id = ? AND status IN ('completed', 'user')" " ORDER BY seq DESC LIMIT ?", session_id, window_size, )

在内存层面,每个会话的写入器单独用一个asyncio.Lock保护。用户同时往同一个会话发两条消息的情况,锁会保证串行处理,避免两条流式响应在同一会话里交错写入。

4.2 上下文组装:滑动窗口 + token 预算

模型不是无限上下文,组装 prompt 必须有取舍。我的方案是"从尾部向前,按 token 预算截取"。

具体算法:

  1. 先固定放 system prompt,消耗一部分 token。
  2. 从会话最新的消息开始,倒序往历史走,逐条累加 token。
  3. 累加到剩余预算不足时停止,剩下的更早消息被截断。
  4. 把截取到的消息按正序排列,作为上下文送入模型。
TOKEN_BUDGET = 8000 # 目标模型上下文的一半,留余量给本次回复 def assemble_context(system_prompt: str, history, budget=TOKEN_BUDGET): used = estimate_tokens(system_prompt) selected = [] for msg in reversed(history): cost = estimate_tokens(msg.content) + 4 # 每条消息的 overhead if used + cost > budget: break selected.append(msg) used += cost selected.reverse() return [{"role": "system", "content": system_prompt}] + [ {"role": m.role, "content": m.content} for m in selected ]

几个值得注意的细节:

  • budget 不要设满。目标模型的上下文如果是 16k,我一般只给历史分配 8k,剩下的留给 system prompt 和本次回复的生成空间。
  • 消息级截断优于 token 级截断。宁可整体丢弃更早的一轮对话,也不要在一轮的中间切断,否则模型看到的语义是残缺的。
  • failed 状态的消息进不了上下文。组装之前必须过滤掉没有完成的消息,否则模型会被半截文本干扰。

4.3 会话切换与懒加载

多会话产品里最常见的交互是:用户在一个列表里点来点去,切换会话。我的做法是消息列表懒加载,而不是一次性把所有会话的消息全查出来。

会话列表只查sessions表,用updated_at DESC排序,展示标题和最后消息时间。点击某个会话时,才去查该会话最近的消息。再配合"加载更多"的分页查询,初始只取 50 条,下拉时再取更早的 50 条。这样即使某个会话聊了几百轮,打开列表页的性能也不会被拖垮。

懒加载还有一个附带好处:内存里的活跃会话数量可控。因为每次只渲染当前会话,系统维护的流式写入器、上下文缓存都只跟当前活跃的少数会话有关,不是全量。

4.4 同会话并发写冲突的处理

最容易被忽略的场景是:用户在同一个会话里,趁上一条回复还没生成完,就发了第二条消息。

这会在数据层引发两个问题:一是上下文组装时读到半成品 assistant 消息;二是两条流式写入同时更新同一条 session 记录,message_count可能算错。

处理方案分两步。第一步,组装上下文时只取status='completed'role='user'的消息,streaming 中的半成品直接跳过。第二步,对sessions表的message_countlast_message_at更新,改成"定稿时重算"而不是"追加时自增"。

async def bump_session(self, session_id: str): await self.db.execute( """ UPDATE sessions SET message_count = (SELECT count(*) FROM messages WHERE session_id = ?), last_message_at = now(), updated_at = now() WHERE id = ? """, session_id, session_id, )

有人会觉得 count(*) 慢,但一个会话的消息数量撑死在几千条,加上session_id上的索引,子查询代价非常小。用重算替代自增,从根上避免了并发下计数器漂移的问题。

5. 踩坑实录:数据层设计里我交过的学费

5.1 坑一:created_at 排序导致消息顺序错乱

第一版代码我用created_at给消息排序,测试时一切正常。直到有一次给消息表做批量导入,导入脚本用的时间戳是"当前时间批量写入",结果同一会话里的消息顺序全乱了,因为一批导入的消息 created_at 几乎相同,排序结果不确定。

后来我把seq作为排序键,created_at只做展示时间。所有新消息的seq都从会话当前最大值加一,导入脚本也必须显式指定 seq。这个改动之后再没出现过顺序问题。

5.2 坑二:流式中断留下的 content_delta 残片

某次线上故障排查时发现,有大量 assistant 消息的content是空的,但content_delta里有一段完整的文本。原因是流式请求在处理过程中被强制 kill,finalize 流程没跑,content_delta里的内容就永久停留在"临时区"了。

修复方案是双管齐下。一是让读取侧兼容:读取消息时,如果content为空但content_delta不为空,把 delta 当做已生成内容返回。二是加一个定时巡检任务,找出状态为 streaming 但超过 10 分钟没有更新的消息,把它们强制标记为 failed 并合并残片。这套机制上线后,再也没有出现过"消息消失了"的用户反馈。

5.3 坑三:token 统计口径不一致

上下文组装的时候,我用的是字符数除以 4的粗略估算;计费模块用的却是模型返回的usage字段。两个口径对不上,导致后台统计的 token 消耗和实际账单差了不少。

统一口径之后,messages.token_count一律以模型返回的 usage 为准(assistant 消息、user 消息分别记录),上下文组装时的估算只用于做截断判断,不进入任何计费逻辑。估算值和精确值的用途分开,问题自然消解。

5.4 坑四:同会话并发请求把上下文搞乱

还有一次压测时发现,同会话并发发两条消息,模型收到的历史里出现了"未来的消息"。追了半天发现,两个请求几乎同时组装上下文,第二个请求组装时,第一条 assistant 消息还没定稿入库,所以上下文里没有它;等第二条请求生成完,第一条才定稿,顺序看起来就像时间线错位了。

解决办法就是 4.4 里说的:上下文组装只认 seq 小于当前请求的已定稿消息,并且同一个会话的请求生成过程加锁串行化。宁可让用户的第二条消息排队等一秒,也不能让模型拿到一份逻辑错乱的历史。

5.5 踩坑小结

问题根因解法
消息顺序错乱created_at 排序不稳定引入会话内单调递增 seq 作为排序键
流式残片不可见finalize 流程未执行读取侧兼容 + 巡检任务兜底
token 统计偏差估算与计费口径混用估算只用于截断,计费以 usage 为准
上下文出现未来消息并发请求未串行化同会话加锁 + 只组装已定稿消息
计数器漂移自增更新存在竞态改为定稿时 count(*) 重算

6. 走向生产:索引、缓存与整个数据层的下一步

6.1 索引设计:查询模式决定索引

数据层跑稳之后,我开始盯性能。对话类应用的查询模式其实很固定:查会话列表、查某会话的消息分页、查某会话用于组装上下文。针对这三个模式,索引设计就三张:

  • sessions(user_id, updated_at DESC):支撑会话列表页。
  • messages(session_id, seq):支撑消息分页与上下文取数。
  • messages(session_id, status, seq):支撑"过滤未完成消息"的上下文组装。

不要试图给消息的 content 建全文索引,对话检索是另一个专题,把它和主链路的数据层混在一起会让索引膨胀失控。生产环境跑下来,TinyRobot Kit 在有几千个活跃会话时,最频繁的三条查询都在 10ms 以内。

6.2 缓存:只缓存"稳定"的数据

对话数据的缓存策略,我踩过不少弯路之后总结出一条原则:缓存只给稳定数据,流式数据不进缓存

会话列表可以缓存,因为它的变更频率低、读频率高,TTL 设为 30 秒,用户感觉不到延迟,数据库压力小一大截。消息的完整内容不缓存,因为消息一旦进入流式阶段,每帧都在变,缓存刷新跟不上的话,只会带来一致性问题。最终实现是:只有定稿的 completed 消息,才在读取时做一次短 TTL(比如 5 分钟)的缓存,且只缓存消息完整内容,不缓存列表页。

6.3 演进方向:分支会话、长期记忆与事件溯源

数据层稳定交付之后,TinyRobot Kit 的下一步我有三个方向:

一是分支会话messages.parent_id字段已经预留,后续可以让用户在某个历史节点上"重新开始",生成一个不同走向的会话分支。

二是长期记忆层。会话之间存在跨会话的稳定信息,比如用户偏好、未完成事项。光靠数据库表存聊天记录表达不了这种关系,需要单独提炼一个 memory 存储,与会话数据解耦。

三是事件溯源。目前 messages 表的 status 流转是通过 UPDATE 完成的,如果未来要做完整的审计和重放,需要把状态变更本身作为事件流存下来。这个改造比较大,属于"等到真有需求再动"的范畴。

就个人实践而言,TinyRobot Kit 到目前为止最有价值的不是某个单一的亮点,而是那一套"把对话当数据管理"的完整思维:每条消息有状态机、每个会话有边界、每次写入有缓冲、每次组装有预算。对话类 AI 应用,模型能力决定下限,数据层决定上限。把数据层想清楚,产品规模再大,心里也是踏实的。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询