Day 47 / 共 60 天 · 阶段7 Flow 事件驱动

对话式 Flow:把"事件流程"套上"多轮聊天"的壳

前六天的 Flow 是"跑一次就完"的批处理。但很多产品要的是多轮对话:用户说一句、Flow 回一句、记住上下文、下一句接着聊。CrewAI 把这套能力做成一层可选的对话式扩展_ConversationalMixin + conversation.py 工具函数),复用前面全部机制——对话就是"每收一句用户消息,就 kickoff 一轮,用持久化的 ChatState 串起多轮"。今天读 ChatState、一轮对话的归一化、消息读写、以及意图分类怎么接回 D44 的 router。

📍 你在 60 天里的位置(阶段7 Flow 事件驱动 · 共 8 天)
阶段6 记忆 D41 Flow 总览 D42 装饰器 D43 定义契约 D44 路由跳转 D45 状态持久化 D46 表达式 D47 对话式 D48 选型 阶段8 LLM
💡 先用一个类比兜住今天 对话式 Flow 就像客服工单系统的"会话视图"。底层还是一张张工单(每轮 kickoff),但系统给你一个"会话"外壳,把同一个用户的多张工单按时间串起来、记住聊过什么(ChatState.messages)。用户每发一句,系统开一张新工单(kickoff 一轮),先判断"这句是想问天气还是想退货"(意图分类),再按意图派给对应处理流程(router)。对话 = 持久化状态 + 每轮一次 kickoff + 意图路由,全是前几天的积木。
L01

痛点:一次性的流程,怎么变成能记住上下文的对话?

🤔 痛点Flow 默认 kickoff 一次跑完就结束、状态清空。可对话要"我上一句问了 A,这一句说 '那 B 呢',你得知道 B 是接着 A 问的"。多轮之间怎么保留上下文?每轮又怎么判断用户到底想干嘛、该走哪条处理分支?
💡 一句话本质 对话式 Flow 不引入新机制,全靠组合旧的:① 用 D45 的持久化 + ChatState 跨轮保留 messages;② 每收到一句用户消息,kickoff 一轮(用同一个 session id 恢复上轮状态);③ 轮内先做意图分类,把用户这句归到某个 intent,再用 D44 的 router 按 intent 分流。对话 = 持久化 + 每轮 kickoff + 意图路由的组装。
大白话没有"对话引擎"这种新东西。就是把"存档(D45)+ 分支(D44)"拼一拼:存档让它记住聊过啥,分支让它按你这句的意图走不同的应答。
L02

Mixin:对话能力是"叠"上去的,不是"改"进去的

回顾 D41 的公开 Flow 组合(flow/flow.py:33):

# flow/flow.py:33
class Flow(_ConversationalMixin, RuntimeFlow[T]):
    """Public Flow class with experimental conversational extension behavior."""

crewai/flow/__init__.py 导出的对话式相关符号(flow/__init__.py:7):

# flow/__init__.py:7
from crewai.flow.conversation import (
    ChatState,               # 推荐的对话状态结构
    ConversationalConfig,    # 类级默认(意图、退出词等)
    ConversationalInputs,    # kickoff 约定的输入键
)
_ConversationalMixin 在前★放在基类列表最左:Python MRO(方法解析顺序)里它优先。对话相关方法(如 classify_intentask)由它提供,能覆盖/补充引擎行为。
experimental它来自 crewai.experimental.conversational_mixin——标注"实验性"。核心引擎稳定,对话层还在演进,用 Mixin 隔离,改它不动引擎。
conversation.py 是工具函数集不是一个大类,而是一堆纯函数normalize_kickoff_inputsappend_message 等),Mixin 和引擎按需调用。函数式、易测。
💡 设计取舍①:为什么对话做成 Mixin + 纯函数,而不是内建进 Flow? 大多数 Flow 是批处理,用不到对话。如果把 messages、意图分类这些硬塞进 Flow 基类,就是让 99% 不聊天的 Flow 背着一堆对话字段和逻辑。做成可叠加的 Mixin + 独立的纯函数模块,则"要对话才组合进来、不要就零负担",而且对话逻辑能单独测试、单独演进(标 experimental)。代价是多一层组合、读代码要跨文件。"可选能力做成可插拔层"是控制核心复杂度的关键手段——和 D45 持久化后端、D41 三层拆分一脉相承。
L03

ChatState:一份"对话进度"的推荐状态结构

对话状态的推荐形状(flow/conversation.py:51):

# flow/conversation.py:51
class ChatState(BaseModel):
    """Recommended persisted state shape for multi-turn flows."""
    id: str = Field(default_factory=lambda: str(uuid4()))   # 会话 id(= 持久化键)
    messages: list[LLMMessage] = Field(default_factory=list) # ★多轮消息历史
    last_user_message: str | None = None                     # 最近一句用户输入
    last_intent: str | None = None                           # 最近一次意图分类结果
    session_ready: bool = False                              # 会话是否已初始化
id = 会话 id★和 D45 的 state.id 是同一个东西。对话里它就是"会话/session id"——同一个 id 的多轮 kickoff 共享历史。
messages★对话的核心:一个消息列表,累积 user/assistant 每轮的话。跨轮靠持久化保留。
last_user_message / last_intent缓存最近一句和它的意图——router 读 last_intent 就能分流(L06)。
ChatState 只是"推荐"你也可以用自己的状态模型,只要有 messages 字段即可(L05 的读写函数会找 messages)。ChatState 是省事的默认。
数据结构:ChatState 跨轮累积 ChatState (session=X) messages: [ {user:"上海天气?"}, {assistant:"18°C"}, {user:"那北京呢?"} ] last_intent: "weather" 第1轮 kickoff:存两条 第2轮:读档+追加,接着聊
图注:同一 session id 的多轮 kickoff 通过持久化共享 messages,"那北京呢"才有上下文可依。
L04

归一化:把各种"一句话"入参收敛成一轮 kickoff 输入

对话入参可能来自不同地方,先归一(flow/conversation.py:71):

# flow/conversation.py:71
def normalize_kickoff_inputs(inputs, *, user_message=None, session_id=None):
    if inputs is None and user_message is None and session_id is None:
        return None                          # ★都没给 → 返回 None(保持无状态语义,不硬造空 dict)
    merged: dict[str, Any] = dict(inputs or {})
    if session_id is not None:
        merged["id"] = session_id            # 会话 id → 复用 D45 的恢复机制(inputs["id"])
    if user_message is not None:
        merged["user_message"] = user_message
    return merged
三种入参来源用户可能传 inputs={...}、或 user_message=、或 session_id=。归一成一个 inputs dict,后面统一处理。
session_id → inputs["id"]★对话的"续聊"直接复用 D45:把 session_id 塞进 inputs["id"],kickoff 就会 load_state 恢复上一轮的 messages。没造新机制,接了老接口。
都没给返回 None★细节:不硬造空 dict。docstring 说明这是为了让无状态 Flow 的 FlowStartedEvent.inputs 保持 None,语义干净。
大白话用户可能用好几种方式告诉框架"这是哪个会话、说了啥"。这个函数把它们统一成一种标准格式,其中"会话 id"直接当成 D45 的存档编号——续聊就是"读上次的档"。
L05

消息读写:优先写 state.messages,否则用后备缓冲

追加一条消息,兼容 dict / BaseModel / 无状态三种情况(flow/conversation.py:116):

# flow/conversation.py:116
def append_message(flow, role, content, **extra) -> None:
    message: LLMMessage = {"role": role, "content": content}
    for key, value in extra.items():
        if key in ("tool_call_id", "name", "tool_calls", "files"):
            message[key] = value
    state = getattr(flow, "_state", None)
    if state is not None:
        if isinstance(state, dict):
            messages = state.get("messages")
            if isinstance(messages, list):
                messages.append(message); return          # ① dict 状态有 messages → 写进去
        elif isinstance(state, BaseModel) and hasattr(state, "messages"):
            ...
            state.messages.append(message); return        # ② 模型状态有 messages → 写进去
    # ③ 状态里没 messages → 写进 flow 的后备缓冲
    if not hasattr(flow, "_conversation_messages"):
        object.__setattr__(flow, "_conversation_messages", [])
    flow._conversation_messages.append(message)
优先写 state.messages★消息进 state → 自动跟着 D45 持久化,跨轮保留。这是首选路径。
后备缓冲 _conversation_messages状态里没 messages 字段(比如你用了自定义状态又没加 messages)→ 退到 flow 实例上的临时缓冲。能用但不跨轮持久化
extra 只放白名单键tool_call_id / name / tool_calls / files 才允许进消息——防止乱七八糟的字段污染消息结构。
get_conversation_messages 对称读消息(flow/conversation.py:97)同样"先 state.messages,否则后备缓冲",读写路径一致。
这个"优先 state、否则后备"的双路径,让对话工具函数对状态形状很宽容:你用 ChatState 最省事(自动持久化),用别的也不崩(退化到内存缓冲)。框架用"渐进增强"而非"强制约定"来降低使用门槛。
L06

意图分类:把用户这句归到某个 intent,接回 router

收到用户消息,可选地做意图分类(flow/conversation.py:160):

# flow/conversation.py:160
def receive_user_message(flow, text, *, outcomes=None, llm=None) -> str:
    append_message(flow, "user", text)                # 先记进历史
    set_state_field(flow, "last_user_message", text)
    if outcomes and llm is not None:
        intent = flow.classify_intent(                 # ★让 LLM 从候选 outcomes 里挑一个意图
            text, outcomes, llm=llm, context=get_conversation_messages(flow))
        set_state_field(flow, "last_intent", intent)   # 写进 state.last_intent
        return intent
    return text

整轮的准备工作串起来(flow/conversation.py:184):

# flow/conversation.py:184
def prepare_conversational_turn(flow, *, user_message=None, intents=None, intent_llm=None, config=None):
    if user_message is None:                           # 没直接给 → 从 state 里捞
        ...
    text = _coerce_user_message_text(user_message)
    if not text.strip():
        return
    set_state_field(flow, "last_intent", None)         # ★每轮重新分类,不复用上轮意图
    resolved_intents = intents or (config.default_intents if config else None)
    resolved_llm = intent_llm or (config.intent_llm if config else None)
    if resolved_intents:
        if resolved_llm is None:
            raise ValueError("intent_llm is required when intents are provided")  # 要分类就得给 LLM
        receive_user_message(flow, text, outcomes=resolved_intents, llm=resolved_llm)
    else:
        receive_user_message(flow, text)               # 不分类,纯记录
classify_intent★Mixin 提供的方法:给 LLM 用户这句 + 候选意图列表 + 对话历史,让它挑一个。本质是一次受限的分类调用
last_intent → router★分类结果写进 state.last_intent。你的 @router 读它 return self.state.last_intent,就把对话分流到对应处理分支——意图无缝接回 D44
每轮重置 last_intent开轮先清空——防止"这轮没分类成功却用了上轮的旧意图"这种串味 bug。
要分类必须给 llm给了 intents 却没给 intent_llm → 直接报错。分类要 LLM,缺了没法干。fail-fast。
控制流:一轮对话怎么走 用户一句话 归一+读档 意图分类 → last_intent @router 分流 应答+存档 下一句用户消息 → 用同 session id 再 kickoff 一轮(读到刚存的档)
图注:每轮 = 归一化 → 读档 → 意图分类 → router 分流 → 应答存档;下一轮用同 session id 再来一遍。
L07

ConversationalConfig:类级默认,每轮可覆盖

对话的类级默认配置(flow/conversation.py:37):

# flow/conversation.py:37
@dataclass
class ConversationalConfig:
    """Optional class-level defaults for conversational flows."""
    default_intents: Sequence[str] | None = None       # 默认候选意图
    intent_llm: str | None = None                       # 默认分类用的 LLM
    interactive_prompt: str = "You: "                   # 交互模式提示符
    interactive_timeout: float | None = None
    exit_commands: Sequence[str] = field(default_factory=lambda: _EXIT_COMMANDS_DEFAULT)  # 退出词
    defer_trace_finalization: bool = True               # ★对话默认推迟 trace 收尾
default_intents / intent_llm在类上定义一次,每轮不用重复传。L06 里 config.default_intents 就是取这。
exit_commands交互式对话里输入 "exit"/"quit" 之类就结束——控制台聊天体验。
defer_trace_finalization=True★对话默认推迟trace 收尾(对比批处理默认 False)。因为对话是多轮的,不该每轮都把整个 trace 收尾——要等整个会话结束。这体现"对话有不同的生命周期语义"。
每轮可覆盖类级是默认,prepare_conversational_turn 的参数(intents/intent_llm)能逐轮覆盖——灵活。

👶 小白:意图分类和 D44 的 router,功能不重叠吗?

👨‍🏫 老师:不重叠,是接力。意图分类负责"把一句自然语言翻译成一个意图标签"("我想退货" → "refund"),这需要 LLM 理解语义。router 负责"拿着这个标签决定走哪个处理分支"("refund" → 退货流程),这是确定性的路由。分类是"听懂人话",router 是"照章办事"——前者喂给后者。

L08

取舍 + 今日小结

💡 设计取舍②:为什么"一轮一次 kickoff + 持久化",而不是"一个 kickoff 里 while 循环收多句"? 朴素做法是 kickoff 里写个 while True: msg = input() 一直收消息、常驻内存。但那样:① 进程崩了整段对话全丢;② 没法做无状态部署(Web 服务里每个请求是独立的,不能常驻一个 while);③ 用户走开半小时占着内存。CrewAI 选"每句一次 kickoff、状态落盘":每轮短命、无常驻,天然适配 Web/Serverless(HTTP 来一句请求 → kickoff 一轮 → 存档 → 返回),崩了读档还能续,用户走开也不占资源。代价是每轮要读写数据库(有 IO 开销)。用"持久化换常驻",把对话的连续性从"内存"挪到"存储"——这才是能上生产的对话架构。
⚠️ 边界:多轮对话必须传同一个 session id,否则每轮都是"新会话" 续聊的关键是每轮 kickoff 用同一个 session_id(即 state.id)——它经 normalize_kickoff_inputs 变成 inputs["id"],触发 D45 的读档。如果你每轮都不传 id(或传了不同的 id),每轮都会新建一个空 ChatState,机器人就"失忆"了——你问"那北京呢",它完全不知道你在接哪句。Web 场景要把 session_id 存在客户端(cookie/前端 state),每次请求带回来。另外:意图分类每轮都调一次 LLM,是额外的 token 成本,简单场景可以不配 intents、走别的方式分流。

🧠 今天你应该能回答

  • 对话式 Flow 引入了新引擎吗?它由哪三块旧积木组装?
  • 对话能力为什么做成 Mixin + 纯函数模块,而不是内建进 Flow?
  • ChatStateid 在对话里叫什么?多轮怎么共享 messages?
  • 续聊的 session_id 怎么复用 D45 的恢复机制?
  • append_message 的"优先 state、否则后备缓冲"解决了什么兼容问题?
  • 意图分类和 router 是什么关系?谁喂给谁?
  • 为什么"每句一次 kickoff"比"一个 kickoff 里 while 收消息"更适合生产?

✋ 10 分钟动手

P=lib/crewai/src/crewai/flow
sed -n '1,26p'     $P/flow.py               # 公开 Flow = Mixin + RuntimeFlow
sed -n '28,95p'    $P/conversation.py        # ChatState / Config / 归一化
sed -n '116,181p'  $P/conversation.py        # 消息读写 + 意图分类
sed -n '184,229p'  $P/conversation.py        # prepare_conversational_turn
# 看对话相关导出
python -c "import crewai.flow as f; print([x for x in f.__all__ if 'onvers' in x or x in ('ChatState',)])"
明日预告 · Day 48(阶段收官):七天把 Flow 拆到了底。最后一天做Flow vs Crew 选型取舍:什么场景用 Crew、什么场景用 Flow、怎么"Flow 编排 + Crew 干活"混合,用源码里两者的定位差异给出一张可操作的决策表。
← Day 46 表达式 Day 48 · Flow vs Crew 选型 →