第 7 章 · 02 流式聚合与工具 schema 生成


第 7 章 · 02 流式聚合与工具 schema 生成

本节摘要:上一节看了 Model 层的"骨架",本节走"毛细血管"——四件配套机制。一是 generate_stream 流式输出:TransformersModel 用 TextIteratorStreamer + 后台线程把生成过程逐 token 吐出;二是 agglomerate_stream_deltas 流式 delta 聚合:把零碎 chunk 拼回完整 ChatMessage,工具调用的 arguments 是逐段追加的;三是 get_tool_json_schema:把第 5 章 Tool 的 inputs 自动转成 OpenAI function calling 格式;四是 TokenUsage 统计与 monitoring.py 的 AgentLogger(rich 彩色日志)、OpenTelemetry/Phoenix 遥测、utils.py 的 RateLimiter/Retrying 限流重试。读完本节,Model 层精读收官。

内容来源:原项目源码 src/smolagents/models.pymonitoring.pyutils.py,官方文档 tutorials/inspect_runs.md,精读并套用体系化模板。

⚠️ 注意:流式接口并非所有后端都实现——generate_stream 目前只在 TransformersModel、LiteLLMModel、InferenceClientModel、OpenAIModel 上有;VLLMModel、MLXModel、AmazonBedrockModel 等只有 generate。agent 侧 stream_outputs=True 的回退逻辑是:模型不支持流式就静默走非流式。

学习目标

阅读完本节,你应当能够:

  1. 讲清 TransformersModel 怎么用"后台线程 + streamer"把阻塞的 model.generate 变成生成器。
  2. 手写 agglomerate_stream_deltas 的核心逻辑,特别是工具调用按 index 分桶、arguments 逐段追加。
  3. 解释 get_tool_json_schema 如何把 Tool 的 inputs 规范化(any 转 string、anyOf 拆解、required 推导)。
  4. 说出 TokenUsage 从哪四个位置被采集,并配置一次 OpenTelemetry/Phoenix 遥测、解释 RateLimiter/Retrying 的指数退避公式。

一、generate_stream:把阻塞生成变成生成器

LLM 生成一个完整回答可能要几十秒,流式输出让用户看到"正在打字"。TransformersModel 的实现(models.py:1093-1132,节选):

1114 # Start generation in a separate thread 1115 thread = Thread(target=self.model.generate, kwargs={"streamer": self.streamer, **generation_kwargs}) 1116 thread.start() 1119 is_first_token = True 1120 count_generated_tokens = 0 1121 for new_text in self.streamer: 1122 count_generated_tokens += 1 1123 # Only include input tokens in the first yielded token 1124 input_tokens = count_prompt_tokens if is_first_token else 0 1125 is_first_token = False 1126 yield ChatMessageStreamDelta( 1127 content=new_text, 1128 tool_calls=None, 1129 token_usage=TokenUsage(input_tokens=input_tokens, output_tokens=1), 1130 ) 1131 count_prompt_tokens = 0 1132 thread.join()

这是经典的生产者-消费者模式:transformers 的 model.generate 是阻塞调用,只能塞进后台线程;TextIteratorStreamer(构造时创建于 models.py:959/971)是横跨两个线程的迭代器队列,主线程 for new_text in self.streamer 逐 token 消费,每 token 包一个 ChatMessageStreamDelta 向下 yield

两处细节值得咀嚼:第一,token 计数即流式副产品——每个 delta 记 output_tokens=1,数 delta 就数出了输出 token;输入 token 只挂在第一个 delta 上(1124 行),避免重复累加。第二,streamer__init__ 时就按 VLM/LLM 两条路径创建好,skip_prompt=True 保证不回吐 prompt。

云端模型无需线程——API 本身返回 SSE 事件流,generate_stream 直接翻译事件。以 InferenceClientModel 为例(models.py:1610-1632,LiteLLMModel/OpenAIModel 同构):

1610 for event in self.retryer(self.client.chat.completions.create, 1612 **completion_kwargs, stream=True, stream_options={"include_usage": True}, 1613 ): ... 1626 if choice.delta: 1627 yield ChatMessageStreamDelta( 1628 content=choice.delta.content, 1629 tool_calls=[ChatMessageToolCallStreamDelta( 1630 index=delta.index, id=delta.id, type=delta.type, 1631 function=delta.function, 1632 ) for delta in choice.delta.tool_calls] if choice.delta.tool_calls else None, 1633 )

注意 stream_options={"include_usage": True}——请求 API 在流末尾附带 usage 事件,token 计数就有了权威来源(替代本地数 delta)。

二、agglomerate_stream_deltas:把碎块拼回完整消息

流式 delta 是给 UI 看的,agent 的记忆与解析仍需要完整 ChatMessageagglomerate_stream_deltas(models.py:220-279)负责聚合。两个 dataclass 先认识一下:ChatMessageStreamDelta(models.py:213-217)含 content/tool_calls/token_usage 三字段;工具调用 delta ChatMessageToolCallStreamDelta 多一个 index 字段——这是 OpenAI 流式协议的关键:一次回复可能有多个工具调用,每个 chunk 只带某个调用的一小段(函数名前几个字符、arguments 的 JSON 片段),index 标明"这块属于第几个调用"。聚合逻辑(models.py:230-258,节选):

236 if stream_delta.tool_calls: 237 for tool_call_delta in stream_delta.tool_calls: 239 if tool_call_delta.index is not None: 240 if tool_call_delta.index not in accumulated_tool_calls: 241 accumulated_tool_calls[tool_call_delta.index] = ChatMessageToolCallStreamDelta( 242 id=tool_call_delta.id, type=tool_call_delta.type, 243 function=ChatMessageToolCallFunction(name="", arguments=""), 244 ) 247 tool_call = accumulated_tool_calls[tool_call_delta.index] ... 255 if tool_call_delta.function.arguments: 256 tool_call.function.arguments += tool_call_delta.arguments # 逐段追加

三件事并行累积:content 字符串直接 +=;token 计数求和;工具调用按 index 分桶——accumulated_tool_callsdict[int, delta],首见某 index 时建桶初始化空 name/arguments,之后 id/type 遇非空就覆盖、arguments 逐段拼接。函数名是"覆盖"而 arguments 是"追加",因为各家 API 对 name 有的整段给、有的逐字给。最后统一装配成 ChatMessage(models.py:260-279)。

调用方在哪?agents.py 的流式循环里,模型 delta 直接向下透传给 UI;真正做聚合的是 gradio_ui.py(导入了 agglomerate_stream_deltas)与记忆写入路径——边流式展示、边攒完整消息,两不误

💡 阶梯要点:流式架构的精髓是"delta 是过程,ChatMessage 是结果"。所有下游(agent 记忆、工具解析、token 统计)只消费 ChatMessage;delta 只给前端。index 分桶 + arguments 追加这两行,是接任何 OpenAI 兼容流式 API 都绕不开的样板代码。

三、get_tool_json_schema:工具 schema 自动生成

第 5 章讲过 Tool 的 inputs{"参数名": {"type": ..., "description": ...}} 字典;本节看它怎么变成 API 要的格式(models.py:288-329):

288 def get_tool_json_schema(tool: Tool) -> dict: 289 properties = deepcopy(tool.inputs) 290 required = [] 291 for key, value in properties.items(): 292 if value["type"] == "any": 293 value["type"] = "string" 294 if not ("nullable" in value and value["nullable"]): 295 required.append(key) 297 # parse anyOf ... (拆解 anyOf 联合类型,models.py:297-316) 318 return { 319 "type": "function", 320 "function": { 321 "name": tool.name, 322 "description": tool.description, 323 "parameters": { 324 "type": "object", 325 "properties": properties, 326 "required": required, 327 }, 328 }, 329 }

规范化三连:smolagents 私有的 "any" 类型降级为 string(292 行,JSON schema 不认识 any);required 列表从 nullable 反推(294 行,非可空即必填);anyOf(联合类型,如 str | None 经 @tool 装饰器生成)拆成单 type 或 type 数组并提取 enum。输出正是 OpenAI function calling 格式——因为这套格式已被 Anthropic/Gemini/DeepSeek 等广泛兼容,成了事实标准,_prepare_completion_kwargscompletion_kwargs["tools"] = [get_tool_json_schema(tool) for tool in tools_to_call_from](models.py:540)一次生成,全家通用。vLLM/MLX 这类本地后端则把同一份 schema 喂给 tokenizer 的 apply_chat_template(tools=...),由模型自带的聊天模板渲染进 prompt。

四、TokenUsage 与 monitoring.py:可观测性

TokenUsage(monitoring.py:36-54)极简:input_tokens + output_tokens,__post_init__ 算出 total_tokens。采集点有四处:API 响应的 response.usage(云端)、prompt/generated 张量长度(TransformersModel,models.py:1068/1089)、流式 usage 事件、本地数 delta。累计靠 Monitor(monitoring.py:81-117):

100 def update_metrics(self, step_log): 107 step_duration = step_log.timing.duration 108 self.step_durations.append(step_duration) ... 111 if step_log.token_usage is not None: 112 self.total_input_token_count += step_log.token_usage.input_tokens 113 self.total_output_token_count += step_log.token_usage.output_tokens

update_metrics 在第 4 章的 CallbackRegistry 里注册为 ActionStep 回调(agents.py:434),每步自动执行;run(return_full_result=True) 时汇总进 RunResult.token_usage——第 8 章 open_deep_research 就用它核算一次深度研究的成本。

AgentLogger(monitoring.py:130-273)基于 rich 做控制台美化:log_code 用 monokai 主题的 Panel + Syntax 渲染 agent 写的代码块,log_task 用黄色边框 Panel 展示新任务,log_rule 画分隔线,visualize_agent_tree 用 rich Tree 递归打印 agent 的工具表与 managed_agents 层级(第 8 章多智能体结构一图看清)。日志分 OFF/ERROR/INFO/DEBUG 四级,verbosity_level=2 时连发给模型的完整消息都打印。生产级可观测则走 OpenTelemetry:装好 smolagents[telemetry] 后三行插桩——

from phoenix.otel import register from openinference.instrumentation.smolagents import SmolagentsInstrumentor register() SmolagentsInstrumentor().instrument()

之后 agent 正常跑,每一步的 thought/code/observation/tool 调用全部变成 trace span 发往 Phoenix(或 MLflow、Langfuse——都兼容 OpenTelemetry 标准),在 Web 界面回放整次运行。遥测由 openinference 侧插桩,smolagents 核心代码零侵入。

五、RateLimiter 与 Retrying:限流与指数退避

ApiModel 底座上的两个守护者,都在 utils.py。先看限流(utils.py:497-525):

512 def __init__(self, requests_per_minute: float | None = None): 513 self._enabled = requests_per_minute is not None 514 self._interval = 60.0 / requests_per_minute if self._enabled else 0.0 515 self._last_call = 0.0 517 def throttle(self): 519 if not self._enabled: 520 return 521 now = time.time() 522 elapsed = now - self._last_call 523 if elapsed < self._interval: 524 time.sleep(self._interval - elapsed) 525 self._last_call = time.time()

纯客户端节流:两次调用间隔不足 60/每分钟请求数 秒就 sleep 补齐。requests_per_minute=None 时完全直通。

再看重试(utils.py:551-604,节选)。ApiModel 用四个模块常量装配它(models.py:38-41、1174-1183):最多 3 次、初始等 60 秒、指数底 2、带抖动,且只对限流类错误重试(is_rate_limit_error 在 models.py:1194-1202 检查异常文本含 "429"/"rate limit"/"too many requests"):

555 for attempt_number in range(1, self.max_attempts + 1): 556 try: 557 result = fn(*args, **kwargs) 558 return result 571 except BaseException as e: 573 should_retry = self.retry_predicate(e) if self.retry_predicate else False 575 if not should_retry or attempt_number >= self.max_attempts: 579 raise ... 593 delay *= self.exponential_base * (1 + self.jitter * random.random())

⚠️ 注意:退避公式 delay *= 2 * (1 + random()) 是指数增长加随机抖动——抖动防止多客户端同步重试形成"惊群"(公式参考 OpenAI 官方限流 cookbook)。重试只认限流错误,认证失败、参数错误、超时会立刻抛出——因为那些重试也没用,快速失败让上层 agent 自己走纠错路径(第 8 章 text_to_sql 的自我纠错正是利用这一点)。

本节要点回顾

  1. generate_stream:TransformersModel 用后台线程跑阻塞的 model.generate,主线程消费 TextIteratorStreamer;云端模型直接翻译 SSE 事件;stream_options={"include_usage": True} 让 token 计数有权威来源。
  2. delta 聚合:content 直接拼接、token 求和、工具调用按 index 分桶(name 覆盖、arguments 追加),agglomerate_stream_deltas 把碎块装配回完整 ChatMessage。
  3. 工具 schema:get_tool_json_schema 做三步规范化(any→string、nullable 反推 required、anyOf 拆解),输出 OpenAI function calling 格式,一套 schema 服务全部后端。
  4. TokenUsage 四处采集(API usage/张量长度/流式事件/数 delta),Monitor.update_metrics 按步累计并打印;return_full_result=True 汇总进 RunResult。
  5. AgentLogger 基于 rich:Panel/Syntax/Tree 彩色渲染代码、任务与多智能体树;生产可观测用 OpenTelemetry 插桩(Phoenix/MLflow/Langfuse),核心代码零侵入。
  6. RateLimiter 纯客户端节流(interval=60/rpm);Retrying 指数退避+抖动,最多 3 次、仅对 429 类错误重试,其他错误快速抛出。

下一节:进入第 8 章攀顶——managed_agents 层级式多智能体,看 manager 如何把子 agent 当工具直接写进 Python 代码里调用。


作者与出处
原作者: 灏天文库
整理: 灏天文库整理
本站整理收录,版权归原作者/开源协议所有;欢迎通过原文链接访问源仓库。
发布者: 作者: 灏天文库 转发
评论区 (0)
U