对话式 Flow:把"事件流程"套上"多轮聊天"的壳
前六天的 Flow 是"跑一次就完"的批处理。但很多产品要的是多轮对话:用户说一句、Flow 回一句、记住上下文、下一句接着聊。CrewAI 把这套能力做成一层可选的对话式扩展(_ConversationalMixin + conversation.py 工具函数),复用前面全部机制——对话就是"每收一句用户消息,就 kickoff 一轮,用持久化的 ChatState 串起多轮"。今天读 ChatState、一轮对话的归一化、消息读写、以及意图分类怎么接回 D44 的 router。
ChatState.messages)。用户每发一句,系统开一张新工单(kickoff 一轮),先判断"这句是想问天气还是想退货"(意图分类),再按意图派给对应处理流程(router)。对话 = 持久化状态 + 每轮一次 kickoff + 意图路由,全是前几天的积木。痛点:一次性的流程,怎么变成能记住上下文的对话?
ChatState 跨轮保留 messages;② 每收到一句用户消息,kickoff 一轮(用同一个 session id 恢复上轮状态);③ 轮内先做意图分类,把用户这句归到某个 intent,再用 D44 的 router 按 intent 分流。对话 = 持久化 + 每轮 kickoff + 意图路由的组装。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_intent、ask)由它提供,能覆盖/补充引擎行为。experimental它来自 crewai.experimental.conversational_mixin——标注"实验性"。核心引擎稳定,对话层还在演进,用 Mixin 隔离,改它不动引擎。conversation.py 是工具函数集不是一个大类,而是一堆纯函数(normalize_kickoff_inputs、append_message 等),Mixin 和引擎按需调用。函数式、易测。messages、意图分类这些硬塞进 Flow 基类,就是让 99% 不聊天的 Flow 背着一堆对话字段和逻辑。做成可叠加的 Mixin + 独立的纯函数模块,则"要对话才组合进来、不要就零负担",而且对话逻辑能单独测试、单独演进(标 experimental)。代价是多一层组合、读代码要跨文件。"可选能力做成可插拔层"是控制核心复杂度的关键手段——和 D45 持久化后端、D41 三层拆分一脉相承。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 是省事的默认。归一化:把各种"一句话"入参收敛成一轮 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,语义干净。消息读写:优先写 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,否则后备缓冲",读写路径一致。意图分类:把用户这句归到某个 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。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 是"照章办事"——前者喂给后者。
取舍 + 今日小结
while True: msg = input() 一直收消息、常驻内存。但那样:① 进程崩了整段对话全丢;② 没法做无状态部署(Web 服务里每个请求是独立的,不能常驻一个 while);③ 用户走开半小时占着内存。CrewAI 选"每句一次 kickoff、状态落盘":每轮短命、无常驻,天然适配 Web/Serverless(HTTP 来一句请求 → kickoff 一轮 → 存档 → 返回),崩了读档还能续,用户走开也不占资源。代价是每轮要读写数据库(有 IO 开销)。用"持久化换常驻",把对话的连续性从"内存"挪到"存储"——这才是能上生产的对话架构。session_id(即 state.id)——它经 normalize_kickoff_inputs 变成 inputs["id"],触发 D45 的读档。如果你每轮都不传 id(或传了不同的 id),每轮都会新建一个空 ChatState,机器人就"失忆"了——你问"那北京呢",它完全不知道你在接哪句。Web 场景要把 session_id 存在客户端(cookie/前端 state),每次请求带回来。另外:意图分类每轮都调一次 LLM,是额外的 token 成本,简单场景可以不配 intents、走别的方式分流。🧠 今天你应该能回答
- 对话式 Flow 引入了新引擎吗?它由哪三块旧积木组装?
- 对话能力为什么做成 Mixin + 纯函数模块,而不是内建进 Flow?
ChatState的id在对话里叫什么?多轮怎么共享 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',)])"