☰
FastGPT 流恢复服务重启后生成态快速重置:基于 Redis Stream Activity 的 2 分钟 stale 检测设计
2026/10/11 9:20:12 网站建设 项目流程

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 的最近活动时间作为"服务是否还在持续生成"的更精确依据,从而大幅缩短异常状态修正的响应时间。

问题分析:为什么不能简单加大扫描频率

设计文档明确指出了四个核心矛盾:

  1. 旧逻辑只看MongoChat.updateTime是否超过 30 分钟,修正速度慢;
  2. 服务重启后 Mongo 里的generating状态可能残留,但 Redis Stream 不再有新数据或心跳,两者形成信息差;
  3. 直接把 cron 改成频繁扫描 Mongo updateTime 无法解决问题,因为无法区分"正常的长耗时生成"(例如多轮工具调用、长文档 RAG)和"异常中断"——长耗时任务的 updateTime 同样很久没变;
  4. Redis 可能短暂异常,不能因为一次 Redis 读失败就误把正在生成的会话改成done,否则会造成真实生成中的会话被错误终止。

这些约束决定了方案必须引入"第二信息源"(Redis activity)做交叉验证,同时保留原 30 分钟逻辑作为降级兜底。

最终方案总览

方案由五个部分组成,相互配合形成"快速判定 + 兜底修正 + 异常降级"的完整闭环:

  1. Redis 记录 stream activity(活动时间戳);
  2. 2 分钟无活动视为异常中断(快速判定);
  3. 保留 30 分钟 Mongo 兜底(最终一致性);
  4. Redis 异常时跳过快速修正(防误杀);
  5. 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_SECONDS300(5 分钟)生成中 Redis 流式镜像的续期 TTL(秒)
STREAM_RESUME_POST_COMPLETE_TTL_SECONDS30流结束后缩短 TTL,便于回收(秒)
STREAM_RESUME_REDIS_MAXMEMORY_RATIO0.5Redis 已用内存/maxmemory 达到该阈值时,停止为新请求创建流恢复镜像
STREAM_RESUME_REDIS_MEMORY_CHECK_INTERVAL_MS5000Redis 内存水位检测缓存时长(毫秒),避免每个流请求都调用 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.tsRedis DAL 层:key 生成、touchState刷新 activity、clearMirror、shrinkTTLAfterComplete、getActive读取

七、验证点与测试覆盖

设计文档列出了 5 个验证点,全部在测试中有对应覆盖:

  1. stream 写入会刷新 active key——见 resume.test.ts,覆盖 stream mirror active key 刷新、shrinkTTLAfterComplete后设置短 TTL、清理时删除keyOfActive等场景;
  2. active 超过 2 分钟未更新时,generating会话被修正为done——cleanStaleGeneratingChats.test.ts 中stale-chat场景(updatedAt早于now - STREAM_RESUME_INACTIVE_MS);
  3. active 仍新鲜时,不修正生成态——同一测试中的active-chat场景(updatedAt = now - INACTIVE_MS + 1000),断言MongoChat.updateOne未被调用;
  4. Redis 读取异常时,不执行 2 分钟快速修正——redis.get.mockRejectedValueOnce(new Error('redis down'))模拟 Redis 故障,只有fallbackChat(超过 30 分钟)被修正,inactive-only-chat被跳过;
  5. 超过 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),仅供参考

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

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

立即咨询