Day 12 / 共 20 天 · 阶段3 工作流引擎

工作流事件与流式输出:怎么把"跑图"变成前端的实时直播

Day11 把数据存进了变量池,今天看工作流跑起来时怎么和前端"实时对话"。你在页面上看到的:节点一个个亮起、大模型一个字一个字往外蹦、最后弹出结果——这些都不是跑完才一次性返回的,而是边跑边推。弄清三件事:①引擎产出的几十种"事件"(Queue 事件)长什么样;②主循环怎么从队列里一条条拉事件;③每种事件怎么被分发给对应的处理器、转成 SSE 流推给浏览器。主角是 core/app/apps/advanced_chat/generate_task_pipeline.py

📍 你在 20 天里的位置(阶段3:工作流引擎)
S1 起步 S2 模型运行时 D11 变量系统 D12 事件与流式 S4 RAG 知识库 S5 工具/Agent S6 收官
💡 先用两个类比兜住今天 类比一:整套机制像体育赛事直播。场上(工作流引擎)不停发生事件——开赛、进球、换人、终场;一个"解说员"(task_pipeline)盯着这些事件,逐条翻译成观众(前端)能懂的画面推出去。比赛没结束,直播已经在放了——这就是"流式"。类比二:事件分发像快递分拣中心的传送带。包裹(事件)在传送带上一个个过来,分拣员看包裹上的"类型标签"(事件类型),按一张对照表把它扔进对应的格口(处理器函数)。类型对不上的,直接放过——不是所有事件都要给观众看。
L01

痛点:一个跑 10 秒的工作流,怎么让用户不干等?

🤔 痛点一个稍复杂的工作流:知识库检索 2 秒 + LLM 生成 8 秒。如果等它整个跑完再把结果一次性返回,用户要盯着转圈 10 秒——体验极差。而且中途出错、或者用户想中途停止,也全都看不到、控制不了。我们真正想要的是:节点刚开始跑就告诉前端"检索中…",大模型一吐字就同步显示(打字机效果),出错立刻报错,用户点停就能停。可工作流引擎在后台跑、HTTP 响应在另一头等,它俩怎么实时对上话?
💡 本质:生产者-消费者 + 事件分发Dify 用一条消息队列把两头解耦:工作流引擎(生产者)每发生一件事,就往队列里塞一个 Queue*Eventtask_pipeline(消费者)用一个 for 循环 listen() 队列,来一条处理一条,转成 StreamResponse yield 出去,最外层再包成 SSE(Server-Sent Events)流给浏览器。生产者只管"发生了什么",消费者只管"怎么讲给前端",中间靠队列传递——这就是经典的生产者-消费者模式配上事件分发
L02

本质:三个角色——队列、事件、分发器

角色是什么职责
Queue 事件QueueXxxEventcore/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 为例。
L03

QueueEvent:一张"能发生什么事"的清单

所有事件类型先被一个枚举钉死 QueueEventcore/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"                           # 用户中止
    ...

每个类型对应一个数据类。比如最关键的"文字块" QueueTextChunkEventcore/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 类(从 QueueLLMChunkEventQueueHumanInputFormTimeoutEvent)。你不用背,只要知道:工作流里能发生的每一类事,都对应一个明确的数据类——这让"发生了什么"变成了强类型、可分发的对象,而不是一堆裸字典。
L04

主循环:listen() 一条条拉事件

消费端的心脏是 _process_stream_responsecore/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 _ 兜住了大多数事件。
L05

分发:一张类型→处理器的对照表

通用分发靠"一张表 + 一次查表"。表在 _get_event_handlerscore/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_eventcore/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。
事件从引擎到浏览器:一条传送带 + 一次查表 工作流引擎 发 Queue 事件 队列 listen() _dispatch_event type(event) 查表 _handle_xxx → StreamResponse 处理器表(dict:事件类型 → 处理函数) TextChunk → 打字机 NodeStarted → "xx中…" 未登记 → 放过 yield 出的 StreamResponse 被外层包成 SSE 推给浏览器 终局事件(成功/失败/出错/停止)→ 处理后 break 收播
图注:生产者发事件进队列,消费者逐条拉出、按类型查表分发、翻译成流响应,最后成 SSE。
L06

文字块处理器:打字机效果是怎么来的

最能体现"流式"的就是文字块处理器 _handle_text_chunk_eventcore/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: 帧发给浏览器。
💡 设计取舍:为什么"边推边攒"两件事一起做?纯流式的诱惑是"只管往外推、不留副本",省内存。但 Dify 选择同时在 _task_state.answer 攒一份完整答案——代价是多占一点内存,收益是流一旦断了/要落库/要过输出审核,服务端手里始终有全量文本_handle_output_moderation_chunk 甚至能在中途发现违规内容时切换成"直接返回替换文案"。流式是体验,攒全量是正确性与安全的底牌——两个都要。
L07

节点输出裁剪 + 今日小结

最后补一个细节:节点内部产出的东西五花八门,但推给前端/落库的"工作流运行输出"要有个统一契约。这层裁剪由 project_node_outputs_for_workflow_runcore/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(带回命中的知识来源,挂到消息上)→ QueueNodeSucceededEventQueueNodeStartedEvent(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_responsefor ... 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 行)
明日预告 · Day 13:工作流和模型运行时到这里告一段落。接下来两天进入 RAG(检索增强生成)——先看 core/indexing_runner.py:一份 PDF/文档,是怎么被抽取 → 清洗 → 切分 → 向量化 → 建索引,变成知识库里一条条可检索的片段的。
← Day 11 变量系统 Day 13 · RAG 总览与索引 →