消息回传与实时通信:PaiFlow 的前后端 SSE 是如何实现的?
最早做 PaiAgent 的时候,工作流是等所有节点执行完成后,再一次性把结果返回给前端。
这种方式实现起来很简单,但体验并不好。
如果一个工作流需要执行十几秒甚至几十秒,用户看到的只有一个 Loading,很难知道:
- 工作流到底有没有开始?
- 当前执行到哪个节点?
- 大模型有没有开始生成?
- 是还在执行,还是已经卡死了?
所以到了 PaiFlow,我做了一个比较重要的改造:
工作流执行过程中的每一个关键状态,都实时推送给前端。
最终采用的方案,就是 事件驱动 + 消息队列 + SSE。
1. 为什么选择 SSE?
PaiFlow 的实时通信,本质上是一个非常典型的单向推送场景:
前端发起工作流
↓
后端执行节点
↓
节点不断产生状态和内容
↓
服务端持续推送给前端
前端并不需要持续往服务端发送消息,因此相比 WebSocket,SSE 会更轻量。
整体链路可以简化成:
WorkflowController
↓
SseEmitter
↓
SseStreamCallback
↓
WorkflowEngine
↓
WorkflowMsgCallback
↓
streamQueue
↓
异步消费线程
↓
SseEmitter.send()
↓
前端
前端发起一次 HTTP 请求以后,这条连接不会立即关闭。
后端一边执行工作流,一边把节点状态、LLM 流式内容持续写入这条连接。
2. 从 Controller 开始
工作流执行入口会先创建一个 SseEmitter:
SseEmitter emitter = new SseEmitter();
StreamCallback callback =
new SseStreamCallback(emitter);
workflowEngine.execute(
workflowDSL,
variablePool,
inputs,
callback
);
这里有一个比较重要的设计。
WorkflowEngine 并不依赖 SseStreamCallback,而是依赖:
StreamCallback
也就是说,工作流引擎并不知道底层到底使用 SSE、WebSocket,还是其他通信方式。
如果未来想换成 WebSocket,只需要增加:
WebSocketStreamCallback
工作流引擎和节点执行器基本不用改。
这其实就是把工作流执行和消息传输协议解耦了。
3. 工作流是怎么产生消息的?
进入 WorkflowEngine 后,会创建一个:
WorkflowMsgCallback
它负责监听整个工作流的生命周期。
例如:
onWorkflowStart
onNodeStart
onNodeProcess
onNodeEnd
onWorkflowEnd
执行一个 LLM 节点时,大概是这样的:
AbstractNodeExecutor
↓
onNodeStart()
↓
LLMNodeExecutor.executeNode()
↓
调用大模型
↓
每返回一个 Chunk
↓
onNodeProcess()
↓
模型执行结束
↓
onNodeEnd()
所以工作流并不是执行完一个节点以后才产生消息,而是执行过程中不断产生事件。
例如:
🟡 LLM 节点开始执行
📝 正在生成第一段内容……
📝 正在生成第二段内容……
✅ LLM 节点执行完成
这就是实时状态展示的基础。
4. 为什么中间还要加一个队列?
这里有一个非常关键的设计:
WorkflowMsgCallback 并不会直接调用
emitter.send()。
事件产生以后,会先转换成统一消息对象,然后放入:
streamQueue
大概是:
NodeExecutor
↓
WorkflowMsgCallback
↓
ChatCallBacks
↓
构造 LLMGenerate
↓
streamQueue.offer()
另外一个异步线程负责消费队列:
while (tag) {
LLMGenerate resp = streamQueue.poll();
if (resp != null) {
clientCallback.callback("stream", resp);
}
}
最终:
streamQueue
↓
SseStreamCallback
↓
SseEmitter.send()
↓
浏览器
这样做最大的好处,就是把消息生产和消息发送解耦。
节点执行器只负责:
我产生了一条消息
至于什么时候发送、怎么发送,则由另外一个线程负责。
特别是 LLM 流式输出时,短时间内可能产生大量消息。
如果每产生一个 Token 都同步等待网络发送,很容易反过来拖慢工作流执行。
加一层 Queue,相当于在执行引擎和网络之间增加了一个缓冲区。
5. 消息怎么统一?
如果每种节点自己定义一种返回结构,前端很快就会乱掉。
所以 PaiFlow 把工作流事件统一封装成:
LLMGenerate
它大致分成三部分。
第一部分:元信息
{
"code": 0,
"message": "Success",
"id": "workflow_xxx",
"created": 1764319971
}
主要告诉前端:
这是哪个工作流?
事件有没有异常?
什么时候产生的?
第二部分:workflow_step
{
"workflow_step": {
"node": {
"id": "node-llm-01",
"alias_name": "大模型节点",
"executed_time": 2.5
},
"progress": 0.4
}
}
它负责描述:
现在执行哪个节点?
执行了多久?
整个工作流进行到哪里?
第三部分:choices
{
"choices": [
{
"delta": {
"content": "你好",
"reasoning_content": ""
}
}
]
}
这一部分主要兼容大模型的流式输出格式。
LLM 每生成一段文本,就把增量内容放进:
choices
↓
delta
↓
content
前端不断拼接 content,就能实现我们常见的打字机效果。
6. SSE 最终是怎么发出去的?
真正负责网络发送的是:
SseStreamCallback
核心逻辑其实非常简单:
public void callback(String eventType, Object data) {
String jsonData = JSON.toJSONString(data);
emitter.send(
SseEmitter.event()
.name(eventType)
.data(jsonData)
);
}
最后发送出去的数据类似:
event: stream
data: {"content":"hello"}
event: stream
data: {"content":" world"}
直到工作流执行完成:
emitter.complete();
这条 SSE 连接才正式关闭。
7. 一条消息的完整生命周期
把前面的内容串起来,一条 LLM 消息的生命周期其实就是:
LLM 返回一个 Chunk
↓
LLMNodeExecutor
↓
onNodeProcess()
↓
WorkflowMsgCallback
↓
ChatCallBacks
↓
构造 LLMGenerate
↓
streamQueue.offer()
↓
异步线程 poll()
↓
SseStreamCallback.callback()
↓
SseEmitter.send()
↓
HTTP SSE
↓
前端 handleMessage()
↓
页面实时更新
这样一看,其实就很清楚了。
8. 前端怎么接收?
前端使用的是 fetchEventSource:
fetchEventSource('/workflow/chat', {
method: 'POST',
body: JSON.stringify(params),
signal: controller.signal,
onmessage(e) {
handleMessage(nodes, edges, e, get, set);
}
});
所有后端推送过来的消息,最终都会进入:
handleMessage
再根据:
workflow_step
choices
finish_reason
code
这些字段更新页面。
例如:
RUNNING
↓
节点黄色高亮
PROCESS
↓
不断追加 LLM 内容
SUCCESS
↓
节点变成绿色
ERROR
↓
展示节点错误
这样用户就能实时看到整个工作流的执行过程。
9. 最后的整体架构
PaiFlow 的实时通信,最终可以浓缩成四层:
第一层:事件产生
WorkflowEngine
NodeExecutor
LLMNodeExecutor
↓
第二层:事件标准化
WorkflowMsgCallback
ChatCallBacks
LLMGenerate
↓
第三层:事件缓冲
streamQueue
orderStreamResultQ
↓
第四层:事件传输
StreamCallback
SseStreamCallback
SseEmitter
它解决的其实不是单纯的“SSE 怎么用”,而是三个更重要的问题:
第一,工作流执行与网络通信解耦。
第二,节点事件统一成标准消息模型。
第三,通过消息队列实现生产和消费解耦。
最终无论底层跑的是 Java 工作流、Python 工作流,前端看到的消息结构都可以保持一致。
这也是 PaiFlow 能够实时展示:
工作流开始
↓
节点开始
↓
节点处理中
↓
LLM 流式输出
↓
节点完成
↓
工作流完成
整个执行过程的核心原理。
如果后面想继续深入,其实还可以再拆两个比较有意思的话题:
一个是工作流节点的调度和执行机制;另一个是 VariablePool 如何实现节点之间的数据传递。
这两个弄明白以后,整个 PaiFlow 工作流引擎的骨架基本上也就彻底串起来了。