FastGPT 流恢复服务重启后生成态快速重置:基于 Redis Stream Activity 的 2 分钟 stale 检测设计
【免费下载链接】FastGPTFastGPT is a knowledge-based platform built on the LLMs, offers a comprehensive suite of out-of-the-box capabilities such as data processing, RAG retrieval, and visual AI workflow orchestration, letting you easily develop and deploy complex question-answering systems without the need for extensive setup or configuration.项目地址: https://gitcode.com/GitHub_Trending/fa/FastGPT
流恢复(Stream Resume)是 FastGPT 在对话生成过程中将 SSE 输出镜像写入 Redis Stream,并在断线、刷新后重放历史流的能力。本设计解决的是服务崩溃或重启时,chatGenerateStatus停留在generating而无法及时恢复为done的问题——通过将 Redis Stream 的最近活动时间作为"服务是否仍在持续生成"的精确依据,把异常会话的状态修正时间从最长 30 分钟缩短到 2 分钟,同时用 30 分钟 Mongo 兜底保证极端场景下的最终一致性。读完本文,你将掌握该方案的问题建模、双通道判定策略、Redis 异常降级逻辑,以及对应的源码实现与测试覆盖。
背景:为什么 30 分钟清理不够快
FastGPT 的对话记录通过MongoChat.chatGenerateStatus字段标记生成状态(generating/done)。正常情况下生成结束时状态会被写回done,但当服务崩溃或重启时,正在生成的对话来不及完成这次状态回写,generating就会残留在 Mongo 中。
旧清理逻辑只依赖一个条件:MongoChat.updateTime是否超过 30 分钟。这意味着:
- 崩溃后侧栏和恢复逻辑会在较长时间内继续认为该会话"还在生成",用户等待很久才能看到状态恢复正常;
- 30 分钟是一个过于保守的时间窗,异常恢复速度慢。
与此同时,流恢复的 stream 模式会持续向 Redis Stream 写入数据,并通过心跳维持连接(XREAD 阻塞读取)。因此可以把Redis Stream 的最近活动时间作为"服务是否还在持续生成"的更精确依据,从而大幅缩短异常状态修正的响应时间。
问题分析:为什么不能简单加大扫描频率
设计文档明确指出了四个核心矛盾:
- 旧逻辑只看
MongoChat.updateTime是否超过 30 分钟,修正速度慢; - 服务重启后 Mongo 里的
generating状态可能残留,但 Redis Stream 不再有新数据或心跳,两者形成信息差; - 直接把 cron 改成频繁扫描 Mongo updateTime 无法解决问题,因为无法区分"正常的长耗时生成"(例如多轮工具调用、长文档 RAG)和"异常中断"——长耗时任务的 updateTime 同样很久没变;
- Redis 可能短暂异常,不能因为一次 Redis 读失败就误把正在生成的会话改成
done,否则会造成真实生成中的会话被错误终止。
这些约束决定了方案必须引入"第二信息源"(Redis activity)做交叉验证,同时保留原 30 分钟逻辑作为降级兜底。
最终方案总览
方案由五个部分组成,相互配合形成"快速判定 + 兜底修正 + 异常降级"的完整闭环:
- Redis 记录 stream activity(活动时间戳);
- 2 分钟无活动视为异常中断(快速判定);
- 保留 30 分钟 Mongo 兜底(最终一致性);
- Redis 异常时跳过快速修正(防误杀);
- cron 从每 5 分钟调整为每 1 分钟(缩短恢复时间)。
一、Redis 记录 Stream Activity
设计要点
流恢复写入 Redis stream 时,同步刷新活动状态 key。对应源码在 streamResume.ts 中定义了 key 生成规则(getKeys),三类 key 同属stream:resume命名空间:
keyOfStream: `stream:resume:data:{teamId}:{sourceType}:{sourceId}:{chatId}` keyOfUnavailable: `stream:resume:unavailable:{teamId}:{sourceType}:{sourceId}:{chatId}` keyOfActive: `stream:resume:active:{teamId}:{sourceType}:{sourceId}:{chatId}`activity key 存储的数据结构由StreamResumeActiveStateSchema约束,只记录一个字段:
export const StreamResumeActiveStateSchema = z.object({ updatedAt: PositiveSafeIntegerSchema // 毫秒时间戳 });即keyOfActive的 value 为{"updatedAt": <毫秒时间戳>}。
写入策略:与 TTL touch 绑定,避免额外压力
关键设计是不对每个 chunk 都强制写 Redis,而是按既有 touch 间隔刷新。在 streamResume.ts 的enqueueRaw实现中可以看到,写入 stream 后通过lastTouchedAt变量节流:
const now = Date.now(); if (lastTouchedAt === 0 || now - lastTouchedAt >= this.parsedTtlTouchIntervalMs) { await this.touchState(keys); lastTouchedAt = now; }touchState同时做两件事(Promise.all并行):
expireStream续期 stream key 的 TTL;set刷新 active key,值为{"updatedAt": Date.now()},TTL 与 stream 保持一致。
ttlTouchIntervalMs在 resume.ts 中固定为1000(毫秒),即最多每秒刷新一次 activity,既保证了活动信号的新鲜度,又避免高频 Redis 写入。
生命周期管理
- 流完成时:
shrinkTTLAfterComplete会把 stream key 与 active key 的 TTL 从生成期 TTL 缩短为postCompleteTtlSeconds(默认 30 秒),便于快速回收; - 新镜像创建前:
createMirror开头的clearMirror会一并删除 unavailable、stream、active 三个 key,避免上一次会话的残留状态污染新会话; - 清理 mirror key 时:同样会删除 active key。
相关环境变量
在 env.ts 中定义了完整的配置项,默认值如下:
| 环境变量 | 默认值 | 说明 |
|---|---|---|
STREAM_RESUME_TTL_SECONDS | 300(5 分钟) | 生成中 Redis 流式镜像的续期 TTL(秒) |
STREAM_RESUME_POST_COMPLETE_TTL_SECONDS | 30 | 流结束后缩短 TTL,便于回收(秒) |
STREAM_RESUME_REDIS_MAXMEMORY_RATIO | 0.5 | Redis 已用内存/maxmemory 达到该阈值时,停止为新请求创建流恢复镜像 |
STREAM_RESUME_REDIS_MEMORY_CHECK_INTERVAL_MS | 5000 | Redis 内存水位检测缓存时长(毫秒),避免每个流请求都调用 INFO MEMORY |
二、2 分钟无活动视为异常中断
常量定义
新增常量STREAM_RESUME_INACTIVE_MS = 2 * 60 * 1000(2 分钟),定义于 resume.ts:
export const STREAM_RESUME_INACTIVE_MS = 2 * 60 * 1000;同时暴露两个关键判定函数:
export const getStreamResumeActiveState = async (params: StreamResumeRedisKeysParams) => { return streamResumeCache.getActive(params); }; export const isStreamResumeActiveStale = ( state: StreamResumeActiveState | undefined, now = Date.now() ) => !state || now - state.updatedAt > STREAM_RESUME_INACTIVE_MS;isStreamResumeActiveStale的判定逻辑非常简洁:
- activity 不存在(
!state):视为 stale; - activity 存在但
now - updatedAt > 2min:视为 stale; - activity 仍新鲜:跳过,认为生成仍活跃。
为什么是 2 分钟
设计文档明确说明了依据:stream 模式会每分钟推送一次心跳,2 分钟给了一次心跳延迟留出缓冲。也就是说,正常情况下 activity 的最长更新时间间隔为 1 分钟,2 分钟的阈值可以容忍一次心跳的抖动而不误判。
清理任务的候选筛选
在 cleanStaleGeneratingChats.ts 中,清理任务先通过 Mongo 查询做第一层过滤:
const generatingChats = await MongoChat.find( { chatGenerateStatus: ChatGenerateStatusEnum.generating, updateTime: { $lt: inactiveThreshold } // inactiveThreshold = now - 2min }, { _id: 1, teamId: 1, sourceType: 1, appId: 1, chatId: 1, updateTime: 1 } ) .lean() .exec();查询条件为"生成态 + updateTime 早于 2 分钟前",只取出候选会话,再逐个检查 Redis activity。这样把 Redis 读取量限制在候选集内,而非全量会话高频扫描。
注意查询条件中的updateTime: { $lt: now - 2min }与 activity 判定的关系:候选集天然排除了刚更新过 Mongo 的会话,与 activity 的 2 分钟阈值保持一致。
三、保留 30 分钟 Mongo 兜底
当会话的updateTime已超过 30 分钟(STALE_GENERATING_CHAT_MINUTES = 30)时,直接按旧逻辑修正为done,不再依赖 Redis 判定。核心实现:
const fallbackThreshold = subMinutes(now, STALE_GENERATING_CHAT_MINUTES); const shouldUseFallback = !!chat.updateTime && chat.updateTime < fallbackThreshold; if (shouldUseFallback) { const currentModifiedCount = await markChatAsDone(chat, now); modifiedCount += currentModifiedCount; fallbackCount += currentModifiedCount; continue; }兜底的作用有两层:
- Redis activity key 被提前清理或不存在时,长期异常仍能被修正(例如 stream TTL 到期、active key 被误删、或该会话从未开启过 stream resume 镜像);
- Redis 读异常时,不会立刻误判短时生成会话——因为兜底只作用于超过 30 分钟的会话,短时生成会话会安然跳过。
markChatAsDone是一个带条件守卫的原子更新,只有状态仍为generating时才会更新:
const result = await MongoChat.updateOne( { _id: chat._id, chatGenerateStatus: ChatGenerateStatusEnum.generating }, { $set: { chatGenerateStatus: ChatGenerateStatusEnum.done, updateTime: now, hasBeenRead: false } } ); return result.modifiedCount ?? 0;注意hasBeenRead: false的细节:状态修正后主动把会话标记为未读,确保侧栏能重新展示该会话的最终状态。
四、Redis 异常时跳过快速修正
这是防误杀的关键设计。在遍历候选会话的过程中,如果读取 Redis activity 抛错,则本轮清理记录 warn 日志并设置redisFailed = true:
try { const activeState = await getStreamResumeActiveState({ ... }); if (isStreamResumeActiveStale(activeState, now.getTime())) { const currentModifiedCount = await markChatAsDone(chat, now); modifiedCount += currentModifiedCount; inactiveCount += currentModifiedCount; } } catch (error) { redisFailed = true; logger.warn('cleanStaleGeneratingChats: failed to inspect stream resume activity', { error }); }随后,对于剩余未处理的候选会话,一旦redisFailed为 true 就直接跳过 Redis 快速判定:
if (redisFailed) { continue; }这样Redis 短暂不可用时,不会把真实仍在生成的会话误改成 done;只有超过 30 分钟的会话仍能通过兜底逻辑修正。这是一种典型的"故障时保守降级"策略——宁可延长异常恢复时间,也不牺牲正在生成会话的正确性。
五、cron 调整为每分钟
清理任务的调度位于 cron.ts:
const cleanStaleGeneratingChatCron = () => { setCron('*/1 * * * *', async () => { if ( await checkTimerLock({ timerId: TimerIdEnum.cleanStaleGeneratingChat, lockMinuted: 1 }) ) { await cleanStaleGeneratingChats(); } }); };调整内容:
- 执行频率从每 5 分钟改为每 1 分钟(cron 表达式
*/1 * * * *); - 定时锁时间从 4 分钟同步缩短为1 分钟(
lockMinuted: 1),避免上一轮任务未结束时下一轮并发执行。
这里有一个需要理解的性能前提:虽然频率提升了 5 倍,但由于候选查询先限制generating且updateTime < now - 2min,再按候选逐个检查 Redis activity,频率提升主要用于缩短异常恢复时间,而非全量高频扫描。候选集通常是极小规模(正常情况几乎没有卡在 generating 的会话),因此成本可控。
六、涉及文件与实现证据汇总
| 文件 | 改动/职责 |
|---|---|
| resume.ts | 增加keyOfActive(经getStreamResumeRedisKeys透出)、STREAM_RESUME_INACTIVE_MS;stream 写入时刷新 active state;完成后同步缩短 stream key 与 active key TTL;清理 mirror key 时删除 active key;暴露getStreamResumeActiveState与isStreamResumeActiveStale |
| cleanStaleGeneratingChats.ts | 从 30 分钟 updateTime 单条件清理,改为 Redis activity 快速判定 + 30 分钟兜底;返回modifiedCount、inactiveCount、fallbackCount便于观察修正来源 |
| cron.ts | 清理任务执行频率从 5 分钟改为 1 分钟,定时锁从 4 分钟改为 1 分钟 |
| streamResume.ts | Redis DAL 层:key 生成、touchState刷新 activity、clearMirror、shrinkTTLAfterComplete、getActive读取 |
七、验证点与测试覆盖
设计文档列出了 5 个验证点,全部在测试中有对应覆盖:
- stream 写入会刷新 active key——见 resume.test.ts,覆盖 stream mirror active key 刷新、
shrinkTTLAfterComplete后设置短 TTL、清理时删除keyOfActive等场景; - active 超过 2 分钟未更新时,
generating会话被修正为done——cleanStaleGeneratingChats.test.ts 中stale-chat场景(updatedAt早于now - STREAM_RESUME_INACTIVE_MS); - active 仍新鲜时,不修正生成态——同一测试中的
active-chat场景(updatedAt = now - INACTIVE_MS + 1000),断言MongoChat.updateOne未被调用; - Redis 读取异常时,不执行 2 分钟快速修正——
redis.get.mockRejectedValueOnce(new Error('redis down'))模拟 Redis 故障,只有fallbackChat(超过 30 分钟)被修正,inactive-only-chat被跳过; - 超过 30 分钟的旧会话仍能通过兜底逻辑修正——
fallbackChat的updateTime早于 30 分钟,走兜底分支,fallbackCount为 1。
测试还验证了返回值统计的准确性:第一个用例中两个会话被修正,断言modifiedCount: 2, inactiveCount: 2, fallbackCount: 0;Redis 故障用例中断言modifiedCount: 1, inactiveCount: 0, fallbackCount: 1。这些统计值配合日志输出,可以在线上观测修正到底来自快速判定还是兜底逻辑,便于进一步调参。
八、方案权衡总结
| 场景 | 修正耗时 | 依据 |
|---|---|---|
| 正常生成(activity 新鲜) | 不修正 | Redis activityupdatedAt在 2 分钟内持续刷新 |
| 崩溃/重启(activity 缺失或超时) | ≤ 2 分钟 + cron 周期(≤ 1 分钟) | Redis activity stale 判定 |
| Redis 短暂故障 | 仍为 30 分钟 | 降级到 Mongo 兜底 |
| Redis activity key 被提前清理 | 仍为 30 分钟 | Mongo 兜底保证最终一致 |
这套设计的关键价值在于:用 Redis Stream 的天然心跳信号替代"猜测",在不增加全量扫描成本的前提下,将异常会话状态恢复时间从 30 分钟级缩短到分钟级,同时对 Redis 故障保持保守降级,兼顾了恢复速度与状态正确性。
【免费下载链接】FastGPTFastGPT is a knowledge-based platform built on the LLMs, offers a comprehensive suite of out-of-the-box capabilities such as data processing, RAG retrieval, and visual AI workflow orchestration, letting you easily develop and deploy complex question-answering systems without the need for extensive setup or configuration.项目地址: https://gitcode.com/GitHub_Trending/fa/FastGPT
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考