工作流事件与流式输出:怎么把"跑图"变成前端的实时直播
Day11 把数据存进了变量池,今天看工作流跑起来时怎么和前端"实时对话"。你在页面上看到的:节点一个个亮起、大模型一个字一个字往外蹦、最后弹出结果——这些都不是跑完才一次性返回的,而是边跑边推。弄清三件事:①引擎产出的几十种"事件"(Queue 事件)长什么样;②主循环怎么从队列里一条条拉事件;③每种事件怎么被分发给对应的处理器、转成 SSE 流推给浏览器。主角是 core/app/apps/advanced_chat/generate_task_pipeline.py。
痛点:一个跑 10 秒的工作流,怎么让用户不干等?
Queue*Event;task_pipeline(消费者)用一个 for 循环 listen() 队列,来一条处理一条,转成 StreamResponse yield 出去,最外层再包成 SSE(Server-Sent Events)流给浏览器。生产者只管"发生了什么",消费者只管"怎么讲给前端",中间靠队列传递——这就是经典的生产者-消费者模式配上事件分发。本质:三个角色——队列、事件、分发器
| 角色 | 是什么 | 职责 |
|---|---|---|
| Queue 事件 | QueueXxxEvent(core/app/entities/queue_entities.py) | 描述"发生了一件事"的数据对象:节点开始、文字块、工作流成功… |
| 队列 + 主循环 | queue_manager.listen() | 生产者塞、消费者取;for 循环逐条拉出事件 |
| 分发器 | _dispatch_event + 处理器表 | 按事件类型查表,交给对应的 _handle_xxx 处理器 |
| 处理器 | _handle_xxx_event(十几个) | 把一个事件转成 0 个或多个 StreamResponse yield 出去 |
listen() 一条条拉 → _dispatch_event 查表找处理器 → _handle_xxx 把它翻译成前端能懂的流响应 → yield 出去 → 变成 SSE。后面几节就是把这条链拆开。注意:advanced_chat(高级对话)和 workflow(纯工作流)各有一个自己的 pipeline,今天以功能更全的 advanced_chat 为例。QueueEvent:一张"能发生什么事"的清单
所有事件类型先被一个枚举钉死 QueueEvent(core/app/entities/queue_entities.py:17):
# core/app/entities/queue_entities.py:17
class QueueEvent(StrEnum):
LLM_CHUNK = "llm_chunk" # 基础模式:模型输出块
TEXT_CHUNK = "text_chunk" # 工作流:一段文字块(打字机就靠它)
WORKFLOW_STARTED = "workflow_started" # 工作流开始
WORKFLOW_SUCCEEDED = "workflow_succeeded"
WORKFLOW_FAILED = "workflow_failed"
NODE_STARTED = "node_started" # 某节点开始
NODE_SUCCEEDED = "node_succeeded" # 某节点成功
NODE_FAILED = "node_failed"
ITERATION_START = "iteration_start" # 迭代/循环相关
RETRIEVER_RESOURCES = "retriever_resources" # 检索命中的知识来源
ERROR = "error"
PING = "ping" # 心跳:防止连接被中间层掐断
STOP = "stop" # 用户中止
...
每个类型对应一个数据类。比如最关键的"文字块" QueueTextChunkEvent(core/app/entities/queue_entities.py:186):
# core/app/entities/queue_entities.py:186
class QueueTextChunkEvent(AppQueueEvent):
event: QueueEvent = QueueEvent.TEXT_CHUNK
text: str # ★ 这次要吐的一小段文字
from_variable_selector: list[str] | None = None # 这段文字来自哪个变量(Day11 的 selector!)
in_iteration_id: str | None = None # 若在迭代里,属于哪个迭代
in_loop_id: str | None = None
StrEnum事件类型是字符串枚举——序列化成 JSON 时直接就是 "text_chunk" 这种字符串,前端好认。event 字段有默认值每个具体事件类都把 event 预设成自己的类型,所以造一个 QueueTextChunkEvent(text="你") 就自带 event=TEXT_CHUNK 标签,不用手填。from_variable_selector★还记得 Day11 的 selector 吗?一段流式文字知道自己来自哪个变量,前端就能把它拼到正确的位置。事件系统和变量系统在这里握手。queue_entities.py 里定义了 30 多个 QueueXxxEvent 类(从 QueueLLMChunkEvent 到 QueueHumanInputFormTimeoutEvent)。你不用背,只要知道:工作流里能发生的每一类事,都对应一个明确的数据类——这让"发生了什么"变成了强类型、可分发的对象,而不是一堆裸字典。主循环:listen() 一条条拉事件
消费端的心脏是 _process_stream_response(core/app/apps/advanced_chat/generate_task_pipeline.py:981)。它就是一个大 for 循环:
# core/app/apps/advanced_chat/generate_task_pipeline.py:981
def _process_stream_response(self, tts_publisher=None, trace_manager=None):
for queue_message in self._base_task_pipeline.queue_manager.listen(): # ① 阻塞式逐条拉
event = queue_message.event
match event: # ② 对"终局/特殊"事件特判
case QueueWorkflowStartedEvent():
self._resolve_graph_runtime_state()
yield from self._handle_workflow_started_event(event)
case QueueErrorEvent():
yield from self._handle_error_event(event)
break # ★ 出错:处理完就跳出循环
case QueueWorkflowFailedEvent():
yield from self._handle_workflow_failed_event(event, trace_manager=trace_manager)
break
case QueueWorkflowSucceededEvent():
yield from self._handle_workflow_succeeded_event(event, trace_manager=trace_manager)
break # ★ 成功:也跳出
case QueueStopEvent():
yield from self._handle_stop_event(event, graph_runtime_state=None, trace_manager=trace_manager)
break # ★ 用户停止:跳出
case _: # ③ 其余全部走通用分发
if responses := list(self._dispatch_event(event, ...)):
yield from responses
queue_manager.listen()★这是个"阻塞式生成器":队列有事件就吐一个,没有就等着。所以这个 for 会一直转到工作流出结果为止——它就是"直播不断线"的那根主线。match / casePython 3.10 的结构化模式匹配。这里把会结束直播的几种"终局"事件(成功/失败/出错/停止)单独拎出来处理——因为它们处理完要 break 退出循环。break 的位置★成功、失败、出错、停止之后都 break。这保证"终场哨响了,直播就收播",不会在工作流结束后还傻等下一个永远不来的事件。case _剩下的"过程中"事件(节点开始、文字块、迭代…)全部甩给 _dispatch_event 通用分发(L05)。这一条 case _ 兜住了大多数事件。分发:一张类型→处理器的对照表
通用分发靠"一张表 + 一次查表"。表在 _get_event_handlers(core/app/apps/advanced_chat/generate_task_pipeline.py:902):
# core/app/apps/advanced_chat/generate_task_pipeline.py:902
def _get_event_handlers(self) -> dict[type, Callable]:
return {
QueuePingEvent: self._handle_ping_event,
QueueTextChunkEvent: self._handle_text_chunk_event, # 文字块 → 打字机
QueueWorkflowStartedEvent: self._handle_workflow_started_event,
QueueNodeStartedEvent: self._handle_node_started_event, # 节点开始 → "xx 中…"
QueueNodeSucceededEvent: self._handle_node_succeeded_event,
QueueIterationStartEvent: self._handle_iteration_start_event,
QueueRetrieverResourcesEvent: self._handle_retriever_resources_event,
... # 二十多个映射
}
查表 + 调用在 _dispatch_event(core/app/apps/advanced_chat/generate_task_pipeline.py:940):
# core/app/apps/advanced_chat/generate_task_pipeline.py:940
def _dispatch_event(self, event, *, tts_publisher=None, trace_manager=None, queue_message=None):
handlers = self._get_event_handlers()
event_type = type(event) # ① 拿事件的"类型标签"
if handler := handlers.get(event_type): # ② 表里查处理器
yield from handler(event, tts_publisher=tts_publisher, ...) # ③ 命中:交给它
return
if isinstance(event, (QueueNodeFailedEvent, QueueNodeExceptionEvent)): # 少数用 isinstance 兜
yield from self._handle_node_failed_events(event, ...)
return
return # ④ 没人认领:直接放过(不报错)
type(event)★用事件对象的"精确类型"当查表 key。这是 O(1) 字典查找——不管几十种事件,都是一次 dict.get,比几十个 if/elif 又快又清爽。源码注释也点明这是从"57 个 if"重构来的。:= 海象运算符if handler := handlers.get(event_type):查表和判空一步做完,命中就调用。isinstance 兜底节点失败有两种(Failed / Exception),用 isinstance 一起接住——精确查表 + 家族兜底,两种匹配策略配合。最后 return(放过)★没登记处理器的事件静默忽略,不报错。因为不是每个内部事件都需要推给前端——"没人认领就丢掉"是刻意设计,不是 bug。文字块处理器:打字机效果是怎么来的
最能体现"流式"的就是文字块处理器 _handle_text_chunk_event(core/app/apps/advanced_chat/generate_task_pipeline.py:516):
# core/app/apps/advanced_chat/generate_task_pipeline.py:516
def _handle_text_chunk_event(self, event, *, tts_publisher=None, queue_message=None, **kwargs):
delta_text = event.text
if delta_text is None:
return
should_direct_answer = self._handle_output_moderation_chunk(delta_text) # ① 内容审核
if should_direct_answer:
return
current_time = time.perf_counter()
if self._task_state.first_token_time is None and delta_text.strip():
self._task_state.first_token_time = current_time # ② 记录首字时间(TTFT 指标)
self._task_state.is_streaming_response = True
if tts_publisher and queue_message:
tts_publisher.publish(queue_message) # ③ 顺带喂给语音合成
self._task_state.answer += delta_text # ④ ★累加到完整答案
yield self._message_cycle_manager.message_to_stream_response( # ⑤ ★把这一小段 yield 出去
answer=delta_text, message_id=self._message_id, from_variable_selector=event.from_variable_selector)
delta_text★注意是 delta(增量):每次只推"新增的一小段",不是每次推全文。前端把这些增量一段段拼起来,就是打字机效果。first_token_time记录"首个非空字"的时间——这就是大家常说的 TTFT(首字延迟)指标,流式体验好不好全看它。self._task_state.answer += delta_text★一边推给前端、一边在服务端把完整答案累加起来。为什么要留一份?因为流结束后要把完整消息落库(历史记录、审计),不能只靠前端拼。yield message_to_stream_response★这一段被包成 StreamResponse yield 出去,沿着 L04 的 for 一路冒到最外层,变成一条 SSE data: 帧发给浏览器。_task_state.answer 攒一份完整答案——代价是多占一点内存,收益是流一旦断了/要落库/要过输出审核,服务端手里始终有全量文本。_handle_output_moderation_chunk 甚至能在中途发现违规内容时切换成"直接返回替换文案"。流式是体验,攒全量是正确性与安全的底牌——两个都要。节点输出裁剪 + 今日小结
最后补一个细节:节点内部产出的东西五花八门,但推给前端/落库的"工作流运行输出"要有个统一契约。这层裁剪由 project_node_outputs_for_workflow_run(core/workflow/workflow_run_outputs.py:7)做:
# core/workflow/workflow_run_outputs.py:7
def project_node_outputs_for_workflow_run(*, node_type, inputs, outputs) -> dict[str, Any]:
"""Project internal node outputs onto the workflow-run public contract."""
if node_type == BuiltinNodeTypes.START:
return dict(inputs) # ★ 开始节点:对外暴露的是它的"输入"(用户填的表单)
return dict(outputs) # 其它节点:暴露"输出"
START 节点特判★开始节点没有真正的"输出"——它的价值就是把用户填的表单/问题收进来。所以对外契约里,开始节点暴露的是它的 inputs。这解释了为什么工作流运行详情里"开始"那一格显示的是用户输入。其余节点正常暴露 outputs。一层薄薄的"投影",把内部结构映射成对外的公开契约——内外解耦。QueueWorkflowStartedEvent(前端显示"运行中")→ QueueNodeStartedEvent(node="知识检索")(显示"检索中…")→ QueueRetrieverResourcesEvent(带回命中的知识来源,挂到消息上)→ QueueNodeSucceededEvent → QueueNodeStartedEvent(node="LLM") → 一连串 QueueTextChunkEvent(text="北")、(text="京")、(text="今天")…(打字机逐字出现,服务端 answer 同步累加)→ QueueWorkflowSucceededEvent(主循环处理后 break 收播,完整消息落库)。页面上流畅的"边想边答",就是这几十条事件被逐条分发、翻译、推流的结果。👶 小白:这些事件走的是 WebSocket 吗?为什么要有个 PING 事件?
👨🏫 老师:不是 WebSocket,是 SSE(Server-Sent Events)——基于普通 HTTP 的"服务器单向持续推送",比 WebSocket 轻,正好适配"服务端一直往前端推"的场景。QueuePingEvent(心跳)的作用是:当工作流某一步跑得久、长时间没有内容要推时,定期发个空的 ping 帧,告诉中间的代理/负载均衡"连接还活着,别掐"。否则一些网关会因为"太久没数据"而主动断开这条流。心跳就是"我还在,别挂电话"。
🧠 今天你应该能回答
- 工作流怎么和前端实时对话?(生产者-消费者队列 + 事件分发 + SSE)
- 事件长什么样?(
QueueXxxEvent强类型数据类,30 多种,QueueEvent枚举钉死类型) - 主循环在哪、靠什么驱动?(
_process_stream_response的for ... listen()) - 为什么用字典分发而不是一堆 if?(
type(event)O(1) 查表,清爽且快) - 打字机效果怎么来的?(
TextChunkEvent逐段 delta,前端拼接,服务端同步累加) - PING 事件干嘛的?(心跳,防止长时间无数据时连接被中间层掐断)
✋ 10 分钟动手
cd /Users/bitmart/work/codes/github/AI_WORK/dify/api
# 1. 事件清单:数一数有多少种 Queue 事件
grep -c "^class Queue" core/app/entities/queue_entities.py
sed -n '17,55p' core/app/entities/queue_entities.py # QueueEvent 枚举
sed -n '186,199p' core/app/entities/queue_entities.py # QueueTextChunkEvent
# 2. 主循环 + 分发
sed -n '981,1031p' core/app/apps/advanced_chat/generate_task_pipeline.py # 主循环
sed -n '902,979p' core/app/apps/advanced_chat/generate_task_pipeline.py # 处理器表 + 分发
# 3. 打字机 + 输出裁剪
sed -n '516,549p' core/app/apps/advanced_chat/generate_task_pipeline.py # 文字块处理器
cat core/workflow/workflow_run_outputs.py # 输出投影(18 行)
core/indexing_runner.py:一份 PDF/文档,是怎么被抽取 → 清洗 → 切分 → 向量化 → 建索引,变成知识库里一条条可检索的片段的。