ai-data-science-team 的 LangChain 对话迁移方案:Message-First 架构改造实战指南
【免费下载链接】ai-data-science-teamAn AI-powered data science team of agents to help you perform common data science tasks 10X faster.项目地址: https://gitcode.com/GitHub_Trending/ai/ai-data-science-team
本篇技术指南围绕仓库内规划文档 planning_docs/package_review_initial_release/langchain_conversation_plan.md 展开,系统讲解 ai-data-science-team 将全部 Agent 从"字符串指令"调用方式迁移至以HumanMessage/AIMessage为核心的消息优先(Message-First)结构的完整方案。该迁移使各 Agent 的状态统一以messages: Sequence[BaseMessage]为单一事实源,从而支撑 supervisor 调度与多智能体团队(Multi-Agent Teams)的对话式编排。读完本文,你将掌握:迁移的三大设计原则、按 Agent 逐一改造的 7 个标准步骤、试点 Agent 的源码级实现细节、向后兼容层的设计思路,以及回滚与风险控制方法。
一、迁移背景:为什么需要 Message-First 结构
ai-data-science-team 是一套由多个 LangGraph Agent 组成的数据科学工具集,早期各 Agent 的对外接口是invoke_agent(user_instructions=...)这类"字符串指令 + 业务字段"的调用约定。这种模式在单个 Agent 独立使用时没有问题,但一旦进入 supervisor/team 架构(例如 supervisor_ds_team.py 所实现的监督式数据科学团队),就会出现明显的结构性障碍:
- 对话上下文无法传递:字符串指令只表达"这一轮"的需求,无法携带多轮对话历史;
- 路由信息缺失:supervisor 需要根据完整的
HumanMessage/AIMessage序列判断下一个该派给哪个 worker,字符串接口无法提供统一、标准化的输入形态; - 状态难以聚合:各 Agent 自定义的返回键(如
data_wrangled、recommended_steps)形态各异,团队层无法用统一逻辑读取或评估。
因此,迁移计划定义的目标是:在不破坏现有 API 的前提下,将 agents 迁移到以HumanMessage/AIMessage为核心的消息优先结构,实现与 supervisor/team 架构的兼容。从源码看,这一目标已通过"新增消息入口 + 保留旧入口"的双轨方式落地。
二、三大设计原则
迁移计划明确了三条贯穿始终的原则,它们是后续所有改造决策的评判标准:
| 原则 | 含义 | 落地证据 |
|---|---|---|
| 保持向后兼容 | 保留invoke_agent(user_instructions=...)签名,新增基于消息的入口 | 各 Agent 类中invoke_agent与invoke_messages并存 |
| 单一事实源 | 状态(State)统一携带messages: Sequence[BaseMessage],所有下游逻辑从messages读取 | GraphState中messages: Annotated[Sequence[BaseMessage], operator.add] |
| 工具透明性 | 按需把工具调用消息(或摘要)纳入messages,由开关控制记录详略 | log_tool_calls=True参数控制工具调用日志输出 |
其中"单一事实源"是迁移的核心:一旦messages成为所有 Agent 共享的规范化状态键,supervisor 就可以用同一条 reducer 逻辑合并来自不同 worker 的对话片段,进而路由与评估完整对话。
三、迁移的 7 个标准步骤(按 Agent 逐一改造)
迁移计划给出了一套可复用的改造模板,每个 Agent 按以下 7 步执行:
步骤 1:输入规范化(Inputs)
Agent 的输入首选接受messages;当调用方仍传user_instructions字符串时,将其包装为[HumanMessage(content=...)],并在状态中统一规范化为messages。
以试点 Agent 为例,data_loader_tools_agent.py 中的prepare_messages节点正是这一原则的实现:
def prepare_messages(state: GraphState): print(format_agent_name(AGENT_NAME)) print(" * PREPARE MESSAGES") if state.get("messages"): return {} return {"messages": [("user", state.get("user_instructions"))]}该节点作为图的入口节点,确保无论调用方传入字符串还是消息列表,进入下游逻辑前messages必然存在。
步骤 2:状态 Schema 扩展(State Schema)
状态 Schema 需包含messages: Annotated[Sequence[BaseMessage], operator.add],同时保留各 Agent 特有的业务字段(data、errors、summaries等)。operator.add是 LangGraph 的累加式 reducer,用于把每轮节点返回的新消息追加到既有消息序列中。
在 data_loader_tools_agent.py 中,GraphState继承自 LangGraph 的AgentState并扩展了业务键:
class GraphState(AgentState): user_instructions: str messages: Annotated[Sequence[BaseMessage], operator.add] data_loader_artifacts: dict tool_calls: List[str]注意AgentState本身也定义了messages通道,这里通过显式声明强化了"单一事实源"的约定。
步骤 3:输出规范化(Outputs)
Agent 的最终回答以AIMessage形式追加进messages;同时保留原有结果键(如data_wrangled、recommended_steps)以维持兼容性。也就是说:messages负责"对话记录",业务键负责"机器可读产物",两者各司其职。
在试点 Agent 的post_process节点(data_loader_tools_agent.py)中可以看到这种双轨输出:
last_ai_message = AIMessage(getattr(last_ai, "content", ""), role=AGENT_NAME) # ... return { "messages": [last_ai_message], "internal_messages": internal_messages, "data_loader_artifacts": artifacts if artifacts else last_tool_artifact, "tool_calls": tool_calls, }此外,模板层的报告节点 node_func_report_agent_outputs 也遵循同一约定——将各业务键汇总为 JSON 后封装成AIMessage写入result_key="messages";模板层的 node_func_explain_agent_code 同样以AIMessage作为代码解释的输出载体。这说明"输出即消息"已成为整个仓库的通用惯例。
步骤 4:工具调用记录(Tool Calls)
对于 ReAct/工具型 Agent,可按需将工具调用消息纳入messages,并用一个开关参数(如log_tool_calls=True)控制转录详略,避免对话记录被工具日志淹没。
试点 Agent 在 post_process 中,先通过 utils/messages.py 的get_tool_call_names从消息中提取工具名,再依据log_tool_calls决定是否打印* Tool: <name>与捕获的 artifacts 键:
tool_calls = get_tool_call_names(internal_messages) if tool_calls and log_tool_calls: for name in tool_calls: # 尝试在上一消息中查找 artifact 路径作为提示 ... print(f" * Tool: {name}{path_hint}")同一开关模式也出现在 eda_tools_agent.py 与 mlflow_tools_agent.py 中,三处默认值均为True。
步骤 5:访问器更新(Accessors)
将 getter 方法改为从messages中读取"最后一条 assistant 消息",同时保留返回 artifacts/data 的辅助方法,确保旧调用方(期望字符串返回)不受影响。
试点 Agent 的 get_ai_message 演示了"倒序扫描最后一条 AI 消息"的标准写法:
msgs = self.response.get("messages", []) last_ai = None for msg in reversed(msgs): role = getattr(msg, "role", None) or getattr(msg, "type", None) if role in ("assistant", "ai"): last_ai = msg break if last_ai is None and msgs: last_ai = msgs[-1]与此配套,get_internal_messages 返回完整内部消息(支持markdown=True格式化输出),get_artifacts 则负责把工具产物按需转换为 DataFrame,形成"对话与产物分离"的访问模式。
步骤 6:入口节点保障(Entry Nodes)
每个 Agent 图的第一个节点必须确保messages已存在(字符串输入在此被包装),然后再进入下游逻辑。试点 Agent 的prepare_messages即承担此职责;而 data_cleaning_agent.py 的invoke_agent在入口处就完成了字符串到消息的转换:
self.response = self.invoke( { "messages": [("user", user_instructions)] if user_instructions else [], "user_instructions": user_instructions, "data_raw": data_raw.to_dict(), "max_retries": max_retries, "retry_count": retry_count, }, **kwargs, )步骤 7:Supervisor/团队接入(Supervisors/Teams)
当messages在所有子 Agent 间标准化后,supervisor 即可基于统一的消息序列完成路由与对话评估。这一点在 supervisor_ds_team.py 中得到充分体现:SupervisorDSState以messages为团队级共享对话通道,并使用自定义 reducer_supervisor_merge_messages进行合并;supervisor 的路由 prompt 通过MessagesPlaceholder(variable_name="messages")把完整对话注入路由链(supervisor_ds_team.py),再结合 OpenAI function-calling 或文本解析输出"下一个 worker"。
四、试点 Agent 源码深度解析:data_loader_tools_agent
迁移计划的推出顺序第一步是"先在单个 Agent 上试点",选中的是data_loader_tools_agent,并标记为 ✅ done(message-first、同步/异步消息入口、工具日志开关、目录/文件产物处理)。该 Agent 也是理解整个迁移模板的最佳范本。
4.1 四个消息入口
DataLoaderToolsAgent 提供四类调用方式:
| 入口 | 同步/异步 | 用途 |
|---|---|---|
invoke_agent(user_instructions=...) | 同步 | 兼容旧调用方,内部将字符串包装为消息 |
ainvoke_agent(user_instructions=...) | 异步 | 兼容旧调用方的异步版本 |
invoke_messages(messages) | 同步 | 面向 supervisor/teams 的首选入口,直接传入Sequence[BaseMessage] |
ainvoke_messages(messages) | 异步 | 异步版本,同样面向 supervisor/teams |
invoke_messages的实现(data_loader_tools_agent.py)把user_instructions置为None、直接透传messages,与旧入口在编译图上完全复用同一套节点逻辑——这正是"不破坏现有 API"原则的工程体现。
4.2 三节点线性图
该 Agent 的编译图为prepare_messages → react_agent → post_process的线性流程(data_loader_tools_agent.py):
- prepare_messages:入口保障,见步骤 1;
- react_agent:调用
create_react_agent构建的 ReAct 工具调用智能体,其内部使用 LangGraph 预置的AgentState作为状态 Schema,并注入system_hint指导工具选择(例如"用户问 LIST 文件时用搜索/列目录工具,而不是加载内容",见 data_loader_tools_agent.py); - post_process:从内部消息中提取最后一条 AI 回复、按工具名聚合 artifacts、汇总工具调用名,最终返回
messages+ 业务键。
该 Agent 挂载的 6 个工具(load_directory、load_file、list_directory_contents、list_directory_recursive、get_file_info、search_files_by_pattern,定义于 tools/data_loader.py)均在 data_loader_tools_agent.py 中注册,工具的 artifacts 会通过post_process汇总到data_loader_artifacts键。
五、向后兼容层:双轨并行的工程保障
迁移计划明确要求"保留invoke_agent签名与字符串返回的 getter",以应对两类既有依赖:
- 旧调用方期望字符串输入:所有 Agent 的
invoke_agent/ainvoke_agent继续接受user_instructions: str; - 旧调用方期望字符串返回:
get_ai_message(markdown=True)、get_internal_messages(markdown=True)等仍可输出 Markdown 格式化文本。
对数据清洗这类携带 DataFrame 的 Agent,消息入口还额外做了指令回填:当invoke_messages未显式传user_instructions时,data_cleaning_agent.py 通过 get_last_user_message_content 从消息列表中提取最近一条 Human 消息内容作为指令,保证下游"推荐清洗步骤"等节点无需改动即可工作。
此外,基类 BaseAgent 在invoke/ainvoke/stream/astream四个方法中统一对返回的messages执行remove_consecutive_duplicates去重(agent_templates.py),从框架层缓解多轮追加消息可能产生的重复记录问题——这是对"State Drift(状态漂移)"风险的又一道防线。
六、推出顺序与转换检查清单
6.1 分四步的灰度推出
迁移计划采用"由点及面"的灰度策略,降低整体回归风险:
- 试点单 Agent(
data_loader_tools_agent):✅ 已完成——message-first、同步/异步消息入口、工具日志开关、目录/文件产物处理全部落地; - 扩展到工具类 Agent(清洗、wrangling、可视化);
- 扩展到 SQL 与特征工程 Agent;
- 更新团队/supervisor 装配层,使其完全依赖
messages进行调度。
6.2 转换检查清单(当前状态)
计划文档记录的检查清单显示核心 Agent 已全部完成转换:
data_loader_tools_agent(试点)data_cleaning_agent(message-first 入口与 demo 已添加)data_wrangling_agent(同上)data_visualization_agent(同上)sql_database_agent(同上)feature_engineering_agent(同上)eda_tools_agent(ds_agents)h2o_ml_agent(message-first 入口、校验调整与 demo 已添加)mlflow_tools_agent(ml_agents)
源码搜索可以印证上述状态:invoke_messages/ainvoke_messages方法已出现在 data_cleaning_agent.py、data_visualization_agent.py、data_wrangling_agent.py、feature_engineering_agent.py、sql_database_agent.py、eda_tools_agent.py、h2o_ml_agent.py 与 mlflow_tools_agent.py 等文件中。
6.3 下一阶段:非核心模块
核心 Agent 转换完成后,迁移进入"非核心模块"阶段,计划包括:
multiagents 子包:
sql_data_analyst:以 message-first 方式包装子 Agent、加入系统提示(system hint)、暴露子图(subgraph)、补充 demo;pandas_data_analyst:同样实现 message-first 入口、系统提示、子图可见性并添加 demo;- supervised/team 变体(如有)应对齐 message-first 并暴露子图。
从源码看,pandas_data_analyst.py 与 sql_data_analyst.py 均已实现
invoke_messages/ainvoke_messages,且 pandas 分析器的状态归一化逻辑会保留system/human/assistant三类消息角色(pandas_data_analyst.py),说明该阶段工作已在持续推进;apps:更新 notebooks/demos 在合适场景改用
invoke_messages。例如 supervisor_ds_team.ipynb 已导入HumanMessage并调用team.invoke_agent,data_loader_tools_agent.ipynb 则大量演示了基于字符串指令的invoke_agent用法,两类调用方式在示例库中长期并存。
七、风险与缓解措施
迁移计划明确列出三大风险及对应缓解策略,这些约束也解释了为何仓库中保留大量"看似重复"的兼容代码:
7.1 转录膨胀(Transcript Bloat)
工具调用消息若全部写入messages,会显著膨胀 LLM 上下文。缓解方式:用log_tool_calls开关按需裁剪工具消息的详略。
团队层还有更激进的兜底——supervisor_ds_team.py 定义了TEAM_MAX_MESSAGES = 20与TEAM_MAX_MESSAGE_CHARS = 2000两个常量,其自定义 reducer_supervisor_merge_messages会:丢弃tool/function角色消息、剔除冗长的 JSON 型 "Agent Outputs" 报告、截断超长消息体、并仅保留最后 20 条消息(supervisor_ds_team.py),从源头防止多步骤工作流中的 token 与速率限制问题。
7.2 既有调用方期望字符串(Existing Callers Expect Strings)
缓解方式:完整保留invoke_agent签名与所有字符串返回型 getter。这正是第四节"双轨并行"设计的直接动机。
7.3 状态漂移(State Drift)
缓解方式:节点始终返回新的messages,避免对外部状态做原地修改。配合BaseAgent的去重逻辑(agent_templates.py),保证同一消息不会因 reducer 追加而在多次 invoke 之间反复累积。
八、总结
ai-data-science-team 的 LangChain 对话迁移方案是一份"工程上可复制"的 Message-First 改造模板:它以messages: Sequence[BaseMessage]为单一事实源,通过"新增消息入口 + 保留旧接口 + 工具日志开关 + 消息化访问器"四件套,在零破坏的前提下完成了 9 个核心 Agent 的对话化改造,并已向 multiagents 与 apps 层延伸。对于任何基于 LangGraph 构建多智能体系统的团队,这份方案在向后兼容策略、状态 Schema 设计、上下文裁剪防膨胀、灰度推出顺序四个维度上都具备直接的借鉴价值——当你想让多个独立 Agent "会聊天"并接受统一调度时,从invoke_messages开始,让对话历史成为唯一的权威数据源。
【免费下载链接】ai-data-science-teamAn AI-powered data science team of agents to help you perform common data science tasks 10X faster.项目地址: https://gitcode.com/GitHub_Trending/ai/ai-data-science-team
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考