1. 项目概述:从零构建一个流式AI应用的数据骨架
最近在折腾一个叫Eino的AI应用项目,核心目标是想把大语言模型(LLM)的能力,以一种更流畅、更可控的方式集成到实际业务里。这听起来像是很多团队都在做的事,但真正动手时,你会发现,最头疼的往往不是模型调用本身,而是数据怎么“喂”给模型,以及模型“吐”出来的东西怎么被下游系统理解和处理。这就引出了我们项目的三个核心概念:Message、ToolCall和流式管道。
简单来说,Message定义了AI对话的“语言”,即用户输入、AI回复、系统指令这些信息单元的结构。ToolCall则是让AI从“聊天”升级到“行动”的关键,它定义了AI如何请求调用外部工具(比如查数据库、发邮件、执行计算)。而流式管道,就是把Message的组装、ToolCall的解析与执行、以及最终结果的流式返回,这一整套流程给串联起来的“高速公路”。这个项目的本质,就是在设计这条高速公路的交通规则和车辆规格。
很多开发者一开始会直接调用OpenAI或类似平台的API,把对话历史拼成一个数组就发过去了。这在Demo阶段没问题,但一旦要处理复杂逻辑、支持多种工具、或者需要稳定的流式输出,代码很快就会变成一团乱麻。Eino项目就是要解决这个痛点,通过一套清晰的数据建模和管道设计,让AI应用的开发变得像搭积木一样可控。接下来,我就拆开揉碎了,讲讲我们是怎么设计这套“积木”的。
2. 核心概念深度解析:Message, ToolCall与Schema
2.1 Message:不只是文本的对话单元
在Eino的设计里,Message不是一个简单的字符串。它是一个结构化的数据对象,这是与直接使用字符串数组最本质的区别。我们为Message定义了一个基础的Schema(模式),通常包含以下核心字段:
- role: 发送者角色。通常是
system(系统指令)、user(用户)、assistant(AI助手)。 - content: 消息内容。这里就有讲究了,它可以是字符串,也可以是一个复杂数组,用以支持多模态(如图片、文档片段)或特定格式的内容块。
- name(可选): 在多人对话或特定工具调用场景中,标识具体的参与者或工具名。
为什么需要这么复杂?举个例子,如果你想让AI记住用户的身份信息,并在回复时个性化称呼,你可能会在system消息里嵌入一段指令:“用户的名字是{name}”。但如果后续对话中名字变了,或者你需要动态注入信息,修改历史消息就很麻烦。更优雅的做法,是将user消息的content设计为结构化数据,比如{"text": "用户输入的问题", "metadata": {"user_name": "张三"}}。这样,在流式管道的某个环节,我们可以很容易地提取和注入这些元数据,而不污染核心的对话文本。
实操心得一:Content字段的设计是灵活性的关键。早期我们只用字符串,后来为了支持文件上传和复杂指令,改成了支持OpenAI风格的content数组格式(例如[{"type": "text", "text": "..."}, {"type": "image_url", "image_url": {"url": "..."}}])。这要求管道中的每个处理器都能理解这种格式。我们的经验是,在项目内部统一一种扩展性好的Content格式,并编写相应的编解码工具函数,能省去后期大量适配工作。
2.2 ToolCall:让AI“动手”的标准化指令
ToolCall是连接LLM“思考”与外部世界“行动”的桥梁。一个典型的ToolCall Schema包含:
- id: 本次调用的唯一标识符,用于在后续的
ToolCall结果(Tool Call Result)中进行匹配。 - type: 固定为
"function"(目前主流LLM工具调用都采用函数形式)。 - function: 具体函数信息。
name: 要调用的函数/工具名称。arguments: 调用参数,是一个JSON格式的字符串。
这里最大的坑在于arguments这个JSON字符串。LLM输出的arguments是一个字符串,你需要将其解析成真正的JSON对象才能调用工具。但LLM的生成并不总是稳定的,可能会输出格式错误、字段缺失或类型不对的JSON。比如,要求参数是整数,LLM可能生成带引号的数字字符串"5"。
实操心得二:必须为每个工具定义严格的JSON Schema,并在调用前进行校验和修复。我们使用jsonschema库来验证LLM生成的参数。更关键的一步是“软化”验证:不是一遇到错误就抛出异常导致流程中断,而是尝试自动修复常见问题,比如修剪多余的反斜杠、将字符串数字转为整数等。同时,将验证和修复过程记录下来,用于后续优化提示词(Prompt),教LLM生成更规范的参数。
2.3 Schema:一切契约的基石
上面反复提到的Schema,是这一切能运转起来的“宪法”。我们主要涉及两种Schema:
- Message Schema: 定义了整个对话历史数组的结构。确保无论是从前端接收、从数据库读取,还是发送给LLM,数据格式都是一致的。
- Tool Schema (Function Calling Schema): 描述每个工具(函数)的规格,包括函数名、描述、参数列表及其每个参数的详细定义(类型、描述、是否必填、枚举值等)。这个Schema会作为“工具清单”的一部分,在对话开始时提供给LLM,让LLM知道它能调用什么。
生成这些Schema是个技术活。对于Tool Schema,我们通常直接从后端的工具函数定义(比如Python的def函数)通过反射自动生成。这里会用到像pydantic这样的库来定义参数模型,然后用inspect模块或pydantic本身的能力来提取函数签名和类型注解,最终转换为OpenAI等LLM所需的格式。
避坑指南:注意Schema的版本兼容性。当你更新了一个工具的函数签名(比如增加了一个可选参数),对应的Tool Schema必须同步更新。我们建立了自动化流程:在CI/CD中,如果检测到工具函数定义变更,会强制重新生成并检查Schema文件,确保开发环境、测试环境和LLM认知中的工具定义保持一致,避免出现“AI以为能调用,但后端接口对不上”的运行时错误。
3. 流式管道设计:数据流动的引擎
有了标准化的“车辆”(Message)和“货物”(ToolCall),就需要设计高效的“公路网”(管道)。流式管道的核心思想是将AI对话的生成、工具调用、结果处理分解为一系列可插拔的步骤,并支持将每个步骤的中间结果实时地、一段一段地(流式)返回给客户端。
3.1 管道的基本结构
一个典型的Eino流式管道包含以下阶段,数据像流水一样依次经过:
- 输入预处理:接收原始用户输入,可能包含文件、指令标记等。将其标准化为内部的Message格式,并附加上下文(如用户ID、会话ID)。
- 对话历史管理:从存储中加载当前会话的历史Message,并根据策略(如Token数限制、关键信息摘要)进行裁剪或总结,组装成即将发送给LLM的上下文列表。
- LLM调用与流式解析:这是核心。我们将组装好的Message列表和Tool Schema清单发送给LLM(如GPT-4),并开启流式响应。我们不是等LLM全部生成完再处理,而是一边接收Token,一边实时解析。特别要解析其中是否包含了
tool_calls的起始标记和内容。 - 工具调用分派与执行:一旦在流中完整解析出一个ToolCall对象(通过检测到特定的结束标记或结构),立即暂停等待后续的文本生成,并异步或同步地执行该工具调用。执行需要解析
arguments,调用对应的后端函数,获取结果。 - 结果注入与继续生成:将工具执行的结果格式化为一个特殊的
Message(role为tool,包含对应的tool_call_id和结果content),插入到对话历史中。然后,将更新后的历史再次发送给LLM,让它基于工具执行结果继续生成后续内容。这个过程(生成->检测到ToolCall->执行->再生成)可能循环多次。 - 输出后处理与流式返回:将LLM生成的文本Token、ToolCall的解析事件、以及最终工具执行结果的摘要,通过Server-Sent Events (SSE) 或WebSocket实时推送给前端。同时,可能进行内容过滤、格式美化等后处理。
3.2 流式处理ToolCall的挑战
在流式模式下处理ToolCall是一大难点。LLM在生成tool_calls时,其输出在流中不是一次性完整出现的。它可能先输出{"id": "call_abc", "type": "function", "function": {"name": "get_weather", "arguments": "{",然后隔几个Token再输出"city": "北京",最后输出"}"}}`。
我们的解决方案是设计一个“流式解析器状态机”。这个解析器监听来自LLM的每一个Token或数据块:
- 初始状态为“等待文本”。
- 当检测到
"tool_calls"或类似起始标记时,进入“解析ToolCall对象”状态。 - 在解析状态下,它需要累积字符,直到能解析出一个完整的、语法正确的JSON对象片段(比如一个完整的
tool_calls数组项)。这里不能简单用字符串匹配,因为参数里的JSON字符串本身可能包含大括号。我们采用了一个轻量级的、容错的JSON分词器(tokenizer)来追踪括号匹配,从而判断一个JSON对象何时结束。 - 一旦解析出一个完整的ToolCall对象,立即触发工具调用流程,并将一个“工具调用开始”的事件推送给前端流,告知用户“AI正在调用XX工具”。
- 工具执行完成后,将结果注入,并让解析器状态回到“等待文本”,继续处理后续的LLM生成流。
实操心得三:流式解析的健壮性高于一切。必须对LLM输出的各种边界情况做处理:JSON片段不完整、编码转义错误、甚至LLM“胡言乱语”出非JSON内容。我们的解析器在无法确定得到一个完整对象时,会持续累积数据,并设置一个超时或缓冲区上限。如果累积了过多数据仍无法解析,则判定为LLM输出异常,本次ToolCall失效,转而向LLM发送一个错误提示,引导它重新生成或继续文本输出。这个错误处理逻辑本身也是管道可配置的一部分。
4. 核心实现细节与代码组织
4.1 数据模型定义(Pydantic实践)
我们使用Pydantic来严格定义所有的核心数据模型,这提供了运行时类型校验、自动文档生成和序列化/反序列化的便利。
from typing import Literal, Union, List, Dict, Any from pydantic import BaseModel, Field class TextContentBlock(BaseModel): type: Literal["text"] = "text" text: str class ImageContentBlock(BaseModel): type: Literal["image_url"] = "image_url" image_url: Dict[str, str] # 通常包含 `url` 字段 ContentBlock = Union[TextContentBlock, ImageContentBlock] class Message(BaseModel): role: Literal["system", "user", "assistant", "tool"] content: Union[str, List[ContentBlock]] # 支持字符串或复杂内容块 name: str | None = None tool_calls: List["ToolCall"] | None = None # 仅当 role=assistant 时可能有 tool_call_id: str | None = None # 仅当 role=tool 时必须有 class ToolCallFunction(BaseModel): name: str arguments: str # JSON字符串 class ToolCall(BaseModel): id: str type: Literal["function"] = "function" function: ToolCallFunction # 使向前引用生效 Message.model_rebuild() class ToolSchema(BaseModel): """对应OpenAI风格的函数定义""" type: Literal["function"] = "function" function: Dict[str, Any] # 包含name, description, parameters(JSON Schema) class ChatCompletionRequest(BaseModel): messages: List[Message] tools: List[ToolSchema] | None = None stream: bool = False # ... 其他LLM参数使用Pydantic后,任何不符合模型的数据在进入管道时就会被拦截,极大减少了后续环节的潜在错误。同时,.dict()和.json()方法让数据转换非常方便。
4.2 工具注册与Schema生成
我们建立一个中央注册表来管理所有可用的工具。
import inspect import json from typing import Callable, get_type_hints from pydantic import create_model, BaseModel class ToolRegistry: def __init__(self): self._tools: Dict[str, Callable] = {} self._schemas: Dict[str, Dict] = {} def register(self, func: Callable): """注册一个工具函数,并自动生成其Schema""" self._tools[func.__name__] = func self._schemas[func.__name__] = self._generate_schema(func) return func # 方便用作装饰器 def _generate_schema(self, func: Callable) -> Dict: # 1. 获取函数签名和类型注解 sig = inspect.signature(func) type_hints = get_type_hints(func) # 2. 为每个参数创建Pydantic模型字段 fields = {} for param_name, param in sig.parameters.items(): if param_name == 'self': continue param_type = type_hints.get(param_name, str) field_info = ... # 根据param的默认值等构造Field信息 fields[param_name] = (param_type, field_info) # 3. 动态创建参数模型 args_model = create_model(f"{func.__name__}Args", **fields) # 4. 生成OpenAI兼容的JSON Schema parameters_schema = args_model.schema() # 5. 组装完整工具Schema tool_schema = { "type": "function", "function": { "name": func.__name__, "description": func.__doc__ or "", "parameters": parameters_schema } } return tool_schema def get_schemas_for_llm(self) -> List[Dict]: """获取所有工具的Schema,用于发送给LLM""" return list(self._schemas.values()) async def execute(self, tool_call: ToolCall) -> Any: """执行一个ToolCall""" func = self._tools.get(tool_call.function.name) if not func: raise ValueError(f"Tool {tool_call.function.name} not found") # 解析参数 try: args_dict = json.loads(tool_call.function.arguments) except json.JSONDecodeError as e: # 尝试修复常见的JSON格式错误 args_dict = self._attempt_fix_json(tool_call.function.arguments) # 使用Pydantic模型验证参数(使用前面动态创建的模型) # ... 验证和转换参数 ... # 执行函数 result = await func(**args_dict) if inspect.iscoroutinefunction(func) else func(**args_dict) return result # 使用示例 registry = ToolRegistry() @registry.register def get_weather(city: str, date: str | None = None) -> str: """获取指定城市的天气信息。 Args: city: 城市名称,例如“北京”。 date: 查询日期,格式YYYY-MM-DD,默认为今天。 """ # ... 实现逻辑 ... return f"{city}的天气是..." # 获取所有Schema发送给LLM tools_for_llm = registry.get_schemas_for_llm()注意事项:动态创建模型可能带来轻微性能开销和序列化问题。在生产环境中,我们通常会在应用启动时一次性生成所有工具的Schema并缓存起来,而不是每次请求都动态生成。同时,确保工具函数的文档字符串(__doc__)清晰完整,因为LLM会依赖这个描述来决定是否以及如何调用该工具。
4.3 流式管道处理器实现
管道由一系列处理器(Processor)组成,每个处理器负责一个特定阶段。我们采用类似“中间件”或“责任链”的模式。
from abc import ABC, abstractmethod from typing import AsyncGenerator class PipelineContext(BaseModel): """管道上下文,携带数据流经整个管道""" session_id: str user_input: Message history: List[Message] llm_response_stream: AsyncGenerator | None = None tool_calls_to_execute: List[ToolCall] = [] final_output_stream: AsyncGenerator | None = None # ... 其他元数据和状态 class PipelineProcessor(ABC): @abstractmethod async def process(self, context: PipelineContext) -> PipelineContext: pass class LLMStreamingProcessor(PipelineProcessor): def __init__(self, llm_client, tool_registry: ToolRegistry): self.llm = llm_client self.tool_registry = tool_registry async def process(self, context: PipelineContext) -> PipelineContext: # 1. 准备LLM请求 request = ChatCompletionRequest( messages=context.history + [context.user_input], tools=self.tool_registry.get_schemas_for_llm(), stream=True ) # 2. 发起流式请求并创建解析器 raw_stream = await self.llm.chat.completions.create(**request.dict()) stream_parser = ToolCallStreamParser() # 前面提到的状态机解析器 # 3. 定义内部异步生成器,用于产出处理后的流事件 async def _processed_stream(): async for chunk in raw_stream: # 解析增量内容 delta = chunk.choices[0].delta text_delta = delta.content or "" tool_call_deltas = delta.tool_calls or [] # 将增量喂给解析器 parsed_events = stream_parser.feed(text_delta, tool_call_deltas) for event in parsed_events: if event.type == "text": # 产出文本Token yield {"type": "text", "data": event.data} elif event.type == "tool_call_start": # 产出工具调用开始事件 tool_call = event.data context.tool_calls_to_execute.append(tool_call) yield {"type": "tool_call", "data": {"status": "started", "id": tool_call.id, "name": tool_call.function.name}} elif event.type == "tool_call_ready": # 解析器判定一个完整的ToolCall已就绪 tool_call = event.data # 这里可以触发异步执行,但不阻塞流 asyncio.create_task(self._execute_and_requeue(tool_call, context)) elif event.type == "error": # 产出错误事件 yield {"type": "error", "data": event.data} context.llm_response_stream = _processed_stream() return context async def _execute_and_requeue(self, tool_call: ToolCall, context: PipelineContext): """执行工具,并将结果作为新消息加入历史,触发重新处理""" try: result = await self.tool_registry.execute(tool_call) tool_message = Message( role="tool", tool_call_id=tool_call.id, content=json.dumps(result, ensure_ascii=False) ) # 将工具结果消息加入历史,并可能触发新一轮的管道处理(例如通过一个消息队列) await self._requeue_for_next_round(context, tool_message) except Exception as e: # 处理执行错误,生成错误信息工具消息 error_message = Message(role="tool", tool_call_id=tool_call.id, content=f"Error: {str(e)}") await self._requeue_for_next_round(context, error_message)这个LLMStreamingProcessor是管道的核心,它连接了LLM的流式输出、ToolCall的流式解析、以及工具执行的异步触发。_processed_stream这个异步生成器是流式输出的源头,它产出的标准化事件(文本、工具调用开始、错误等)可以被后续的StreamOutputProcessor直接转发给客户端。
5. 常见问题、调试与性能优化
在实际开发和运维中,我们遇到了各种各样的问题,这里总结几个最有代表性的。
5.1 问题排查清单
| 问题现象 | 可能原因 | 排查步骤与解决方案 |
|---|---|---|
| LLM不调用工具 | 1. Tool Schema描述不清。 2. 提示词(System Message)未明确指示使用工具。 3. 对话历史过长,工具定义被截断。 | 1. 检查工具函数的description和参数描述是否清晰易懂。用简单任务测试。2. 在System Message中加入“请使用可用工具来回答问题”等指令。 3. 检查Token计数,确保包含工具定义的上下文未被裁剪。 |
| ToolCall参数解析失败 | 1. LLM生成的JSON格式错误(如缺少引号、尾逗号)。 2. 参数类型不匹配(如字符串传给了数字参数)。 3. 参数值为空或缺失必填参数。 | 1. 在流式解析器中加入更强大的容错JSON解析器,如json5或自定义修复逻辑。2. 在工具执行前,用Pydantic模型进行强校验和类型转换。 3. 在Tool Schema中明确标记 required字段,并在提示词中强调。 |
| 流式响应中断或卡住 | 1. 工具执行耗时过长,阻塞了流。 2. 网络问题或LLM服务端超时。 3. 管道中某个处理器抛出未处理异常。 | 1.关键:工具执行必须异步化,不阻塞文本流推送。使用asyncio.create_task。2. 设置合理的LLM调用超时和重试机制。 3. 在管道每个处理器外层添加全局异常捕获,将错误转化为流式错误事件输出,而不是崩溃。 |
| 上下文长度超限 | 1. 对话历史积累过多Token。 2. 工具执行结果(特别是大段文本)被追加后超限。 | 1. 实现对话历史总结/裁剪策略。例如,将较早的Message替换为一句摘要。 2. 对工具返回的大结果进行压缩或截断,只保留关键信息再注入历史。 |
| 多轮工具调用混乱 | 1. 同一轮对话中多个ToolCall的id匹配错误。2. 工具结果消息未正确关联 tool_call_id。 | 1. 确保ToolCall的id在单次LLM响应内唯一,并在整个会话中妥善管理。2. 严格遵循格式: role="tool"的消息必须包含tool_call_id,且内容是对应ToolCall的执行结果。 |
5.2 调试技巧
- 记录完整的管道日志:为每个会话(
session_id)记录管道每个阶段的输入输出。特别是记录发送给LLM的完整Message列表和Tool Schema,以及LLM返回的原始响应流。当工具调用不符合预期时,复查这些日志是最直接的。 - 使用“调试模式”:在开发环境,可以配置管道跳过实际的LLM调用和工具执行,使用预设的“剧本”来模拟整个流程。这能快速验证管道逻辑是否正确,特别是复杂的多轮工具调用场景。
- 可视化流事件:前端可以开发一个调试面板,实时显示接收到的SSE事件(文本块、工具调用开始/结束、错误等)。这能帮你直观看到流是否顺畅,ToolCall事件是否在正确的时间点被触发。
5.3 性能优化点
- Schema缓存:如之前所述,Tool Schema应在服务启动时生成并缓存,避免每次请求都进行反射和生成。
- LLM上下文管理:历史消息的裁剪和总结算法需要高效。可以计算每个Message的Token数并缓存,避免每次请求都重新计算。
- 工具执行并行化:如果一次LLM响应中包含了多个独立的ToolCall(比如同时查询天气和股票),应该并行执行这些工具,而不是串行,以降低整体延迟。
- 流式解析器优化:解析器状态机的实现要高效,避免在累积数据时进行复杂的字符串操作。可以考虑使用更底层的字节操作或特定优化的JSON流解析库。
- 连接与超时管理:对于长时间运行的流式对话,要处理好HTTP/WebSocket连接的超时、重连和状态恢复。确保在客户端意外断开后,服务器端能安全地清理资源。
构建Eino这样一个基于Message、ToolCall和流式管道的系统,是一个将离散技术点串联成稳定服务的过程。它要求你对LLM的工作原理、前后端数据交互、异步编程和错误处理都有深入的理解。这套架构的价值在于,它提供了一个清晰、可扩展的框架,让你能专注于业务工具的开发,而不必每次都重新发明轮子来处理AI交互的复杂性。当你的工具越来越多,交互逻辑越来越复杂时,一个健壮的数据流管道就是确保一切井然有序的关键。