1961 字
约 6 分钟
1
消息回传与实时通信:PaiFlow 的前后端 SSE 是如何实现的?
2026-09-16
2026-09-16

消息回传与实时通信: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 工作流引擎的骨架基本上也就彻底串起来了。

消息回传与实时通信:PaiFlow 的前后端 SSE 是如何实现的?
http://www.clxhxhhr.top/posts/642/
作者
clxstart
发布于
2026-09-16
许可协议
CC BY-NC-SA 4.0
评论
0 条
还没有评论,先写一条吧。