请求生命周期、SSE 流式与队列扩展
30 秒导读: 前三章讲的是"图怎么被表示、怎么被编译成流水线"。本章把它们串成一次真实的在线请求——一条 HTTP
POST /prediction/:id从进门、鉴权、分流到引擎,再把模型吐出的每个 token 通过 SSE(Server-Sent Events,服务端单向推流) 实时送回浏览器;最后讲清楚生产环境怎么用 BullMQ 队列 + Redis 发布/订阅把这套东西水平扩展到多台机器。读者是运维和后端。引擎内部算法不在这里重复,见 经典引擎 和 AgentFlow V2 引擎。
1. 这一章解决什么问题
前面几章是"静态"视角:一张图长什么样、怎么被拓扑排序、怎么被解释执行。但线上跑起来时,你面对的是一串很实际的工程问题:
- 一个
POST请求进来,从哪个路由到哪段代码,中间做了哪些鉴权和分流? - 用户要的是"打字机效果"(逐字出字),服务端怎么在一个 HTTP 连接上持续推数据?
- 多轮对话怎么记住上文?用户点"停止生成"时怎么真的把 LLM 掐断?
- 单进程扛不住并发了,怎么加机器?加了机器以后,worker 在 A 机器上跑, 浏览器连在 B 机器上,流式事件怎么跨进程送回去?
本章按这个顺序讲。核心结论先给你:
| 关切 | 一句话机制 | 关键角色 |
|---|---|---|
| 入口分流 | 路由 → 服务 → utilBuildChatflow 取图、鉴权、按 type 分流 | buildChatflow.ts:990 |
| 流式回传 | 一个长连接 + SSEStreamer 把每类事件写成 SSE 帧 | SSEStreamer.ts |
| LangChain→SSE | LangChain 回调处理器 CustomChainHandler 在 handleLLMNewToken 里调 streamTokenEvent | components/src/handler.ts:344 |
| 多轮记忆 | 从图里找 Memory 节点 → 定 sessionId → 拉历史 | utils/index.ts:1905/1823/1869 |
| 中断 | AbortControllerPool 按 chatflowid_chatId 存 AbortController,一 abort 全链路取消 | AbortControllerPool.ts |
| 水平扩展 | MODE=QUEUE 时把执行丢进 BullMQ,worker 消费,Redis pub/sub 转发 SSE 事件 | queue/* |
2. 顶层全景:一次流式请求怎么走
先看单进程模式(默认,MODE 未设为 queue)。一次带 streaming: true 的外部预测请求,链路是这样的:
浏览器 / 嵌入组件
│ POST /api/v1/prediction/:id body={ question, streaming:true, chatId? }
▼
┌─────────────────────────────────────────────────────────────┐
│ Express 应用 (App) │
│ ① 路由挂载 routes/index.ts → /prediction │
│ ② 控制器 controllers/predictions createPrediction │
│ · 校验来源域名 allowedOrigins │
│ · 打开 SSE:setHeader(text/event-stream) + flushHeaders() │
│ · sseStreamer.addExternalClient(chatId, res) ← 记住这条连接 │
│ ③ 服务 services/predictions buildChatflow │
│ └─ utilBuildChatflow(req) buildChatflow.ts:990 │
│ · 按 id 取 ChatFlow(数据库) │
│ · 非内部请求校验 API Key │
│ · 查 workspace / org / 配额 │
│ · 按 chatflow.type 分流 ↓ │
└──────────────┬───────────────────────────┬──────────────────┘
│ type=AGENTFLOW │ 其它(CHATFLOW/MULTIAGENT)
▼ ▼
executeAgentFlow executeFlow buildChatflow.ts:301
(AgentFlow V2 引擎) (经典引擎 / 老多智能体)
│ │
│ 两条路都拿到 sseStreamer,执行中不断调它推事件
▼ ▼
sseStreamer.streamXxxEvent(chatId, data) ← token / 工具 / 来源文档 …
│
▼
SSE 帧写回同一个 res 连接: message:\ndata:{"event":"token","data":"你"}\n\n
│
▼
执行完 → removeClient(chatId) 写 end 帧并 res.end()
怎么读这张图: 从上到下是时间顺序。关键点是 ② 那步——控制器先把 HTTP 响应对象 res 存进 SSEStreamer,后面引擎在深处任何地方只要拿着 chatId 就能往这条连接推数据,不必层层传 res。
图里"其它(CHATFLOW/MULTIAGENT)→ executeFlow"是粗粒度的:
MULTIAGENT类型的图进executeFlow后,还要靠端点节点类别再决定走不走buildAgentGraph——细节见 §3.3,别把"type=MULTIAGENT"当成buildAgentGraph的直接开关。
各部件一句话职责:
| 部件 | 干什么 | 在哪 |
|---|---|---|
| 路由 | 把 URL 挂到控 制器 | routes/index.ts:112(/prediction)、:97(/internal-prediction) |
| 控制器 | 开 SSE 连接、登记客户端、调服务 | controllers/predictions/index.ts createPrediction |
| 服务 | 薄封装,转调工具函数 | services/predictions/index.ts:10 buildChatflow |
utilBuildChatflow | 取图、鉴权、配额、按 type 分流、单进程 vs 队列决策 | utils/buildChatflow.ts:990 |
executeFlow | 组装并跑经典引擎的结束节点 | utils/buildChatflow.ts:301 |
SSEStreamer | 管理所有 SSE 连接,把事件写成 SSE 帧 | utils/SSEStreamer.ts:13 |
3. 入口链路:从路由到分流
3.1 两个入口:外部 vs 内部
Flowise 的预测有两个几乎对称的入口,区别只在鉴权和客户端类型:
| 入口 | 路由挂载 | 控制器 | 客户端类型 | 谁用 |
|---|---|---|---|---|
| 外部预测 | routes/index.ts:112 /prediction | controllers/predictions | EXTERNAL(addExternalClient) | API/嵌入网页,要校验 API Key 和来源域名 |
| 内部预测 | routes/index.ts:97 /internal-prediction | controllers/internal-predictions | INTERNAL(addClient) | Flowise 自己的画布"测试对话",走登录态 |
两者最终都调同一个工具函数 utilBuildChatflow(services/predictions/index.ts:3,10),差别通过 isInternal 参数体现——外部入口传 false,内部入口传 true。
3.2 控制器怎么开 SSE
createPrediction 里,只有当**"这条流可以流式" 且 "请求方要流式"** 同时成立,才走 SSE 分支(controllers/predictions/index.ts:57-59)。SSE 分支的动作顺序很关键:
// 示意,非源码:createPrediction 的 SSE 分支骨架
sseStreamer.addExternalClient(chatId, res) // 1. 先登记连接
res.setHeader('Content-Type', 'text/event-stream') // 2. 声明这是 SSE
res.setHeader('Cache-Control', 'no-cache')
res.setHeader('X-Accel-Buffering', 'no') // 关掉 nginx 缓冲,否则流被攒住
res.flushHeaders() // 3. 立刻把响应头发出去,连接就"开着"了
if (isQueueMode) await redisSubscriber.subscribe(chatId) // 4. 队列模式:订阅这个 chatId 的 Redis 频道
const apiResponse = await predictionsServices.buildChatflow(req) // 5. 真正执行(执行中会不断推 token)
sseStreamer.streamMetadataEvent(apiResponse.chatId, apiResponse) // 6. 收尾推一条 metadata
// finally: removeClient(chatId) 关连接
重点看第 1 步和第 5 步:连接先登记好,执行才开始。执行是 await 的,但执行内部会同步地边算边调 sseStreamer.streamTokenEvent(...),数据就已经顺着 res 流出去了——buildChatflow 这个 await 返回时,流早已推完,第 6 步只是补一条元数据。
真实的头部设置与客户端登记见 controllers/predictions/index.ts:69-81;内部入口对应逻辑在 controllers/internal-predictions/index.ts:39(用 addClient)。
3.3 utilBuildChatflow:分流枢纽
utilBuildChatflow(utils/buildChatflow.ts:990)是所有预测的必经关口。它按顺序做这几件事:
- 取图:按
req.params.id从数据库查ChatFlow,查不到抛 404(:996)。 - 鉴权:外部请求(
!isInternal)校验 API Key,失败抛 401(:1030-1035validateFlowAPIKey)。 - 定上下文:查 workspace / organization,取订阅与产品 id,做配额检查
checkPredictions(:1059)。 - 组装执行参数
executeData,把sseStreamer、telemetry、componentNodes等一股脑塞进去(:1061)。 - 决定单进程还是队列(
:1084):MODE=QUEUE→ 把executeData丢进 prediction 队列,等 job 完成(:1085-1090)。- 否则 → 当场
new AbortController(),直接await executeFlow(executeData)(:1100-1104)。
真正的引擎分流在 executeFlow 里,分两层——先按 chatflow.type 岔出 V2,再在经典引擎内部按端点节点类别决定走不走多智能体:
| 走哪条引擎 | 判定条件 | 判定层级 | 代码 |
|---|---|---|---|
executeAgentFlow(V2 引擎) | chatflow.type === 'AGENTFLOW' | 引擎级:看持久化的 type | isAgentFlowV2(:481) → :483 |
buildAgentGraph(老多智能体) | 端点节点 category 含 Multi Agents / Sequential Agents | 经典引擎内子路径:看端点类别,不是看 type | isAgentFlow(:539) → :605-607 |
结束节点 endingNodeInstance.run(...)(普通链) | 以上都不命中 | 经典引擎默认路径 | :788 |
两层分流,别只看 type。 第一层
isAgentFlowV2(:481)在executeFlow函数体最前面按type岔走 V2——AgentFlow V2 是一等公民,一进来就return。第二层的isAgentFlow(:539)才是buildAgentGraph的真正网关,它数的是"端点节点属不属于Multi Agents/Sequential Agents类别",与type无关(与 经典引擎getEndingNodes的端点校验一致)。警惕同名变量陷阱。
utilBuildChatflow里另有一个同名不同义的isAgentFlow = chatflow.type === 'MULTIAGENT'(:1003),它只喂给成功/失败指标计数器(incrementSuccessMetricCounter/incrementFailedMetricCounter,:1096/:1108/:1114),不参与引擎路由。两个isAgentFlow别混::539那个按端点类别、管路由;:1003那个按type、只管打点。
4. 流式回传:SSE 是怎么推的
4.1 一个连接管理器:SSEStreamer
SSEStreamer(utils/SSEStreamer.ts:13)是个进程内单例,核心是一张 Map<chatId, Client>——chatId 是连接的钥匙。它对外暴露一大堆 streamXxxEvent 方法,每个对应一种前端要处理的事件类型:
| 事件方法 | event 字段 | 什么时候推 |
|---|---|---|
streamStartEvent | start | 开始生成(带去重,只发一次,:151) |
streamTokenEvent | token | 每个 LLM token(打字机效果的来源,:165) |
streamAgentReasoningEvent | agentReasoning | 智能体推理步骤 |
streamSourceDocumentsEvent | sourceDocuments | RAG 命中的来源文档 |
streamUsedToolsEvent | usedTools | 调用过的工具 |
streamAgentFlowEvent 等 | agentFlowEvent… | AgentFlow V2 的节点级事件 |
streamMetadataEvent | metadata | 收尾元数据(chatId、sessionId 等,:289) |
streamAbortEvent | abort | 用户中断 |
streamErrorEvent | error | 出错 |
所有方法最后都收敛到一个私有 函数 safeWrite(:66)。它做两件事:
- 写给客户端:
client.response.write(帧);写失败(连接已断)就把这个chatId从 Map 删掉,防止对死连接反复写。 - 扇出给观察者:如果有别的连接注册成了这个
chatId的"观察者"(observers,:19),同一帧也复制一份写给它们——用于多标签同步、管理端旁观、测试探针等。
4.2 SSE 帧长什么样
Flowise 的 SSE 帧格式是固定的字符串拼接(见任一 streamXxxEvent):
message:\ndata:{"event":"token","data":"你好"}\n\n
即:一行 message:,一行 data: 后跟 JSON,再以空行 \n\n 结束一帧。事件类型不放在 SSE 的 event: 字段里,而是塞进 JSON 的 event 键——前端统一在 data 的 JSON 里读 event 来分派。
两个运维相关的细节:
- 心跳:
startHeartbeat(:372,默认 30 秒)对每个活连接写一行 SSE 注释:heartbeat\n\n。客户端会忽略它,但它能穿 过 ALB / 反向代理的空闲超时,防止长连接被中间件掐断。App 启动时就开了心跳(index.ts:133)。 - 收尾:
removeClient(:99)先写一帧{"event":"end","data":"[DONE]"}再res.end(),前端据此知道"这条回答结束了"。
4.3 LangChain 回调怎么变成 SSE(一处需要澄清)
引擎跑的是 LangChain 链路。LangChain 通过**回调(callback)**在"新 token 生成""链开始/结束"等时机通知外部。把这些回调翻译成 SSE 事件的,是 LangChain 的回调处理器 CustomChainHandler(components/src/handler.ts:344):它在 handleLLMNewToken(:366)里拿到每个新 token,直接调 this.sseStreamer.streamTokenEvent(this.chatId, token)(:389、:419)。
这套之所以能通:executeFlow 组装运行参数时,只在流式有效时才把 sseStreamer 和 shouldStreamResponse 传进节点的 run(buildChatflow.ts:781);节点内部据此构造带 sseStreamer 的 CustomChainHandler 挂到 LangChain 上。
一处澄清(诚实优先): 名字容易误导——
utils/callbackDispatcher.ts里的dispatchCallback(:12)不是 LangChain→SSE 的桥,它是一个外发 webhook 分发器:给 webhook 请求体做 HMAC 签名、带指数退避重试(RETRY_DELAYS = [0, 3000, 6000])。真正把 LangChain 回调转成 SSE 的是上面的CustomChainHandler。两者别混。
4.4 什么样的流才允许流式:isStreamValid 门槛
不是所有图都能流式。executeFlow 里先算一个 isStreamValid(buildChatflow.ts:740-746),规则:
- 开了**后处理(postProcessing)**的,一律不流式(
:743)——因为要拿到完整结果再加工。 - 否则调
checkIfStreamValid(:943),它再委托isFlowValidForStream(utils/index.ts:1470):结束节点里若有"自定义函数结束节点"(EndingNode),不许流式;并且要求请求里streaming为真。
AgentFlow V2 是例外——buildAgentGraph 分支里 shouldStreamResponse: true 写死(buildChatflow.ts:622,注释直言 "agentflow is always streamed")。
5. 会话与记忆:多轮怎么串起来
多轮对话的关键是同一个 sessionId 下的历史。Flowise 把它拆成四个工具函数(都在 utils/index.ts):
一次请求进来
│
▼
findMemoryNode(nodes, edges) :1905 从图里找那个连着下游的 Memory 节点
│ (一张图里约定只有 1 个 Memory 节点)
▼
getMemorySessionId(memoryNode, ...) :1823 定这次用哪个 sessionId
│ 优先级:API 传的 overrideConfig.sessionId
│ > API 传的 chatId
│ > UI 里写死的 sessionId
│ > 兜底用 chatId
▼
getSessionChatHistory(sessionId, ...) :1869 用这个 sessionId 初始化 Memory 节点,拉历史消息
│
▼
执行引擎(带着历史当上下文)
三个要点:
- sessionId ≠ chatId。
chatId是"这一轮对话"的标识,sessionId是"记忆归属"的标识。外部 API 可以显式传overrideConfig.sessionId把多个 chatId 归到同一记忆下;不传就退化成用chatId(getMemorySessionId的四级优先级,:1829-1855)。类型不是字符串会直接抛 400。 - 历史从哪来:
getSessionChatHistory(:1869)不是简单查表——它按 Memory 节点的类型(Buffer、Redis、第三方等)import对应节点类,init出实例,再调getChatMessages(sessionId, ...)。也就是说记忆的存取实现由图里选的 Memory 节点决定,这个函数只是统一入口。 - 清记忆:
clearSessionMemory(:768)遍历所有 Memory / OpenAIAssistant 节点逐个清。从"查看消息"弹窗触发时(isClearFromViewMessageDialog)只清匹配memoryType的那一个(:782)。
6. 中断:用户点"停止"时发生了什么
6.1 一个池子:AbortControllerPool
AbortControllerPool(AbortControllerPool.ts:4)就是一张 Record<string, AbortController>,键是 chatflowid_chatId。它只有四个方法:add / remove / get / abort(:38 的 abort 会 .abort() 后再删)。
链路上,单进程模式在 utilBuildChatflow 里执行前 add,执行完 remove(buildChatflow.ts:1100-1106)。这个 AbortController 的 signal 随 executeData 一路传到引擎和 LangChain 调用——一 abort,底层 HTTP 请求(如 OpenAI 调用)就被取消。
6.2 单进程 vs 队列:中断路径不同
停止请求打到 abortChatMessage(services/chat-messages/index.ts:190),按模式分两条路:
用户点"停止" → abortChatMessage(chatId, chatflowid) id = `${chatflowid}_${chatId}`
│
├── 单进程: abortControllerPool.abort(id) 直接掐本进程的 AbortController
│ services/chat-messages/index.ts:201
│
└── 队列模式: publishEvent({ eventName:'abort', id }) 发一条 BullMQ 事件
services/chat-messages/index.ts:196
│ (因为执行在别的 worker 进程,本进程手上没有那个 AbortController)
▼
worker 监听到 'abort' 事件 → abortControllerPool.abort(id)
commands/worker.ts:48-50
为什么队列模式要绕一圈: 收停止请求的 web 进程和真正在执行的 worker 进程不是同一个进程,web 进程手里没有那个 AbortController。所以它把 abort 变成一条 BullMQ 队列事件广播出去,由持有 controller 的 worker 收到后本地执行 abort。这是理解队列模式的一把钥匙:凡是"跨进程找对象"的地方,都改成"发消息"。
7. 队列模式:水平扩展
7.1 为什么要队列
单进程模式下,预测在收请求的那个进程里同步跑完。并发一高,CPU 密集的图执行会把 event loop 堵死,连 SSE 心跳都发不出去。队列模式(MODE=QUEUE)把"收请求"和"跑执行"拆成两种进程:
┌──────────────┐ ┌─────────────┐
浏览器 ──SSE长连接──►│ web 进程 ×N │ │ worker ×M │
▲ │ (main/start) │ │(command:worker)│
│ └──────┬───────┘ └──────┬──────┘
│ │ addJob │ 消费 job
│ ▼ ▼
│ ┌───────────────── Redis / BullMQ ─────────────────┐
│ │ prediction 队列 upsert 队列 schedule 队列 │
│ └───────────────────────┬─────────────────────────┘
│ │ 执行中产生的流式事件
│ ④ SSEStreamer 写回本地连接 │ publish 到 channel=chatId
└────────────◄── RedisEventSubscriber ◄── Redis pub/sub ◄── RedisEventPublisher
(web 进程,订阅 chatId) (worker 进程)
怎么读: web 进程只负责收请求、开 SSE、把活儿 addJob 丢进 Redis;worker 进程把活儿捞出来真跑。跑出来的 token 事件在 worker 侧 publish 到以 chatId 命名的 Redis 频道,web 侧订阅了这个频道,收到就写回它手上那条浏览器连接。
7.2 关键类和职责
| 类 | 职责 | 代码 |
|---|---|---|
QueueManager | 单例;建三条队列(prediction/upsert/schedule),管连接、BullBoard 面板 | queue/QueueManager.ts:22 |
BaseQueue | 抽象基类;封装 BullMQ 的 Queue/Worker/QueueEvents,提供 addJob/createWorker | queue/BaseQueue.ts:12 |
PredictionQueue | 预测队列;processJob 里重新注水依赖并调 executeFlow | queue/PredictionQueue.ts:34 |
RedisEventPublisher | 实现 IServerSideEventStreamer,把每个 stream 事件 publish 到 Redis 频道 | queue/RedisEventPublisher.ts:6 |
RedisEventSubscriber | 订阅 chatId 频道,收到消息还原成 sseStreamer.streamXxxEvent | queue/RedisEventSubscriber.ts:6 |
7.3 入队:哪些东西能丢进队列
utilBuildChatflow 在队列模式下这样入队(buildChatflow.ts:1084-1090):
// 示意,非源码
const predictionQueue = appServer.queueManager.getQueue('prediction')
const job = await predictionQueue.addJob(omit(executeData, OMIT_QUEUE_JOB_DATA)) // 剔掉不能序列化的字段
const result = await job.waitUntilFinished(queueEvents) // 阻塞等 worker 跑完
重点看 omit(executeData, OMIT_QUEUE_JOB_DATA): job 要经过 Redis 序列化,而 executeData 里有一堆不能序列化的活对象——Express 的 res、TypeORM 的 DataSource、sseStreamer 实例等。所以入队前把它们剔掉,只留纯数据。
那 worker 拿到没有这些依赖的 job 怎么跑?靠 PredictionQueue.processJob 重新注水(queue/PredictionQueue.ts:65-74):
// 示意,非源码:worker 侧把活依赖重新塞回去
if (this.appDataSource) data.appDataSource = this.appDataSource
if (this.componentNodes) data.componentNodes = this.componentNodes
if (this.redisPublisher) data.sseStreamer = this.redisPublisher // ← 关键:sseStreamer 换成 Redis 发布器
...
return await executeFlow(data)
这行 data.sseStreamer = this.redisPublisher 是整个队列流式的枢纽:引擎代码完全不知道自己在 worker 里,它照常调 sseStreamer.streamTokenEvent(...);只不过这个 sseStreamer 被悄悄换成了 RedisEventPublisher,于是"写回连接"变成了"publish 到 Redis 频道"(RedisEventPublisher.ts:52 safePublish,频道名就是 chatId)。这是依赖注入换实现的漂亮用法——同一份引擎代码,单进程时写 socket,队列时写 Redis。
7.4 出队与回传
- worker 消费:
worker命令(commands/worker.ts:17)对三条队列各createWorker()(:41)。BaseQueue.createWorker(:59)建 BullMQWorker,并发上限来自WORKER_CONCURRENCY(:8,默认极大的 100000,实际靠机器和 Redis 兜底)。 - 事件回传:web 进程在打开 SSE 时
redisSubscriber.subscribe(chatId)(controllers/predictions/index.ts:77)。worker 侧 publish 的消息被RedisEventSubscriber.handleEvent(:105)按eventType还原成对应的sseStreamer.streamXxxEvent,写回本地那条浏览器连接。 - 收尾清理:请求结束
unsubscribe(chatId);另有startPeriodicCleanup(:91,默认 60 秒)扫掉没有活客户端的僵尸订阅,防止 Redis 频道订阅泄漏。
7.5 MODE 决定进程形态
同一份代码,靠环境变量 MODE 长成不同角色:
MODE 未设 / 非 queue MODE=queue
───────────────────── ─────────────────────────────────
start() index.ts:389 web 进程: start() 里额外
└ App.initDatabase :396 · QueueManager.setupAllQueues index.ts:137-151
└ 单进程: · new RedisEventSubscriber index.ts:153
收请求即 executeFlow · (不 建 worker,只入队)
worker 进程: 单独跑 `flowise worker`
commands/worker.ts → createWorker 消费
App 启动 (start → initDatabase) 里,MODE=QUEUE 分支才会 setupAllQueues 并建 RedisEventSubscriber(index.ts:137-156);worker 是另一个进程(flowise worker 命令)。也就是说生产上你至少跑两类进程:一批 web(收请求 + 持 SSE 连接),一批 worker(跑执行),中间靠 Redis 联通。
8. 单进程 vs 队列:部署对比
| 维度 | 单进程(默认) | 队列(MODE=QUEUE) |
|---|---|---|
| 进程 | 一个 web,收请求即执行 | N 个 web + M 个 worker |
| 依赖 | 无需 Redis(除非用 Redis 记忆等) | 必须 Redis(BullMQ + pub/sub) |
| 执行位置 | 收请求的进程内 executeFlow | worker 进程 processJob → executeFlow |
sseStreamer 实现 | SSEStreamer(直接写 socket) | worker 侧换成 RedisEventPublisher(写 Redis) |
| 流式回传 | 同进程,直接写 res | worker publish → web subscribe → 写 res |
| 中断 | abortControllerPool.abort 直接掐 | 发 BullMQ abort 事件,worker 收到再掐 |
| 扩展性 | 垂直(加大机器) | 水平(加 worker) |
| 运维面板 | 无 | BullBoard /admin/queues(路由挂载 index.ts:350、basePath index.ts:140;面板路由建于 QueueManager.ts:173) |
选择建议(依据代码事实,非投资/业务建议):低并发或自托管小实例用单进程最省事;要水平扩展、隔离执行、削峰才上队列,代价是引入 Redis 和多进程运维。
9. 边界与坑
- SSE 是单向的:服务端 → 客户端。用户的"停止"不是走这条流,而是另发一个 HTTP 请求到
abortChatMessage(§6.2)。 - 代理缓冲会毁掉流式:控制器显式设了
X-Accel-Buffering: no(controllers/predictions/index.ts:73)关 nginx 缓冲;若你的反代/ALB 仍在攒包,前端会看到"一次性蹦出全部文字"而非逐字。心跳(§4.2)也是为穿透代理空闲超时而设。 - 队列模式下事件靠 chatId 频道:如果
chatId在多处不一致,或忘了subscribe,worker 的 token 会 publish 到没人听的频道,前端就"卡住不出字"。 removeOnComplete默认不清:BaseQueue.addJob(:40-56)默认只在失败时删 job,完成的 job 默认保留,除非配REMOVE_ON_AGE/REMOVE_ON_COUNT——长期跑要注意 Redis 里 job 堆积。callbackDispatcher名不副实:它是 webhook 外发器,不是流式桥,别在排查"前端不出字"时去看它(§4.3)。- 两个同名
isAgentFlow别混(§3.3)::539按端点类别、决定走不走buildAgentGraph;:1003按type==='MULTIAGENT'、只喂指标计数。排查"为什么进/没进多智能体路径"要看前者。
10. 代码地图(导航索引)
| 主题 | 文件路径 | 符号名 |
|---|---|---|
| 路由挂载(外部/内部预测) | packages/server/src/routes/index.ts | /prediction(:112)、/internal-prediction(:97) |
| 外部预测控制器 / 开 SSE | packages/server/src/controllers/predictions/index.ts | createPrediction |
| 内部预测控制器 | packages/server/src/controllers/internal-predictions/index.ts | createInternalPrediction |
| 预测服务(薄封装) | packages/server/src/services/predictions/index.ts | buildChatflow(:10) |
| 分流枢纽:取图/鉴权/单进程 vs 队列 | packages/server/src/utils/buildChatflow.ts | utilBuildChatflow(:990) |
| 经典引擎入口 / 两层分流 | packages/server/src/utils/buildChatflow.ts | executeFlow(:301)、isAgentFlowV2(:481,按 type)、isAgentFlow(:539,按端点类别→buildAgentGraph) |
| 同名变量(仅打点,非路由) | packages/server/src/utils/buildChatflow.ts | isAgentFlow(:1003,type==='MULTIAGENT',喂指标计数 :1096/1108/1114) |
| 流式有效性门槛 | packages/server/src/utils/buildChatflow.ts / utils/index.ts | checkIfStreamValid(:943)、isFlowValidForStream(index.ts:1470) |
| SSE 连接管理器 | packages/server/src/utils/SSEStreamer.ts | SSEStreamer(:13)、safeWrite(:66)、startHeartbeat(:372) |
| LangChain 回调 → SSE 桥 | packages/components/src/handler.ts | CustomChainHandler(:344)、handleLLMNewToken(:366) |
| Webhook 外发器(易混,非流式桥) | packages/server/src/utils/callbackDispatcher.ts | dispatchCallback(:12) |
| 会话记忆:找节点/定 session/拉历史/清 | packages/server/src/utils/index.ts | findMemoryNode(:1905)、getMemorySessionId(:1823)、getSessionChatHistory(:1869)、clearSessionMemory(:768) |
| 中断池 | packages/server/src/AbortControllerPool.ts | AbortControllerPool(:4)、abort(:38) |
| 中断触发(单进程 vs 队列) | packages/server/src/services/chat-messages/index.ts | abortChatMessage(:190) |
| 队列管理器 | packages/server/src/queue/QueueManager.ts | QueueManager(:22)、setupAllQueues(:117) |
| 队列基类(BullMQ 封装) | packages/server/src/queue/BaseQueue.ts | BaseQueue(:12)、addJob(:37)、createWorker(:59) |
| 预测队列(注水 + 执行) | packages/server/src/queue/PredictionQueue.ts | PredictionQueue(:34)、processJob(:65) |
| SSE 事件跨进程发布 | packages/server/src/queue/RedisEventPublisher.ts | RedisEventPublisher(:6)、safePublish(:52) |
| SSE 事件跨进程订阅还原 | packages/server/src/queue/RedisEventSubscriber.ts | RedisEventSubscriber(:6)、handleEvent(:105) |
| App 启动 / MODE 分形 | packages/server/src/index.ts | App(:63)、initDatabase(:83)、start(:389) |
| worker 进程入口 | packages/server/src/commands/worker.ts | Worker(:17)、abort 事件监听(:48) |
同组其它章: Flowise 全景 · 数据模型与节点系统 · 经典执行引擎 · AgentFlow V2 引擎 · 前端画布