尧图建网站 尧图建网站 YAOTU WEB BUILD 免费咨询
ARTICLE DETAIL

资讯详情

深耕网站建设与建站编程的一线实战洞察。

构建生产级AI工具调用循环:超时、重试、熔断、验证与监控五件套

构建生产级AI工具调用循环:超时、重试、熔断、验证与监控五件套 1. 从“玩具”到“生产级”ToolUse循环的可靠性鸿沟在AI应用开发的圈子里ToolUse工具调用功能已经不是什么新鲜事了。无论是OpenAI的Function Calling还是各大模型平台提供的类似能力开发者们都能轻松地让大模型学会“使用工具”——调用一个API、查询一次数据库、发送一封邮件。写个Demo跑通流程看着模型精准地返回一个结构化的JSON成就感满满。这感觉就像第一次成功点亮了LED灯宣告着“Hello World”的诞生。但当你真的要把这个“会调用工具的AI”塞进一个7x24小时运行的生产环境去处理真实的用户请求、对接复杂的业务系统时那感觉就完全不同了。你会发现之前Demo里一切顺滑的假设在现实面前脆弱得不堪一击。网络会波动API会超时模型会“抽风”返回一个无法解析的胡言乱语下游服务返回的数据格式可能和文档描述的天差地别。这时你才明白一个能“跑起来”的ToolUse循环和一个能在生产环境“扛得住”的ToolUse循环中间隔着一道巨大的可靠性鸿沟。这道鸿沟不是靠更复杂的Prompt工程或者换一个更强大的模型就能填平的。它考验的是工程化思维是对“错误”的系统性处理能力。一个生产级的ToolUse循环其核心价值不在于它能多么“聪明”地调用工具而在于当一切都不按预期发展时它能否“体面”地失败、清晰地报告、并尽可能地自动恢复或降级。今天我们就来聊聊填平这道鸿沟所必需的“五件套”超时控制、重试机制、熔断降级、结构化输出验证以及全面的日志与监控。这五样东西缺一不可它们共同构成了生产级AI应用稳健运行的基石。2. 第一件套超时控制——为不确定性设定边界超时控制是可靠性的第一道防线。它的核心思想很简单为任何可能阻塞或长时间无响应的操作设定一个明确的时间上限。在ToolUse循环中超时主要发生在两个环节向大模型服务发起请求的等待时间以及模型调用外部工具API的执行时间。2.1 为什么超时如此关键没有超时控制的系统就像一个没有守门员的足球场。一个慢速或无响应的下游服务无论是模型API还是你的业务工具会迅速耗尽你应用服务器的线程或连接池资源导致整个服务“雪崩”影响所有用户。超时的作用就是及时切断这种有害的等待释放资源并给上游调用方一个明确的失败信号使其有机会采取备用方案如重试或降级。2.2 分层超时策略的设计一个粗糙的全局超时是远远不够的。我们需要一个分层的、精细化的超时策略。第一层HTTP客户端超时。这是最底层的网络超时。当你使用requests、httpx或aiohttp等库调用工具API时必须设置连接超时connect_timeout和读取超时read_timeout。连接超时防止在建立TCP连接阶段卡住读取超时防止服务器响应过慢。对于内部微服务连接超时可以设得短一些如2-5秒读取超时根据接口的SLA服务等级协议来定。对于调用第三方不可控API读取超时需要更保守比如10-30秒。import httpx async def call_external_tool(url: str, payload: dict): # 分层超时设置示例 timeout httpx.Timeout(connect5.0, read30.0, write5.0, pool10.0) async with httpx.AsyncClient(timeouttimeout) as client: try: response await client.post(url, jsonpayload) response.raise_for_status() return response.json() except httpx.TimeoutException: # 记录日志触发重试或降级 raise ToolExecutionTimeoutError(f调用工具 {url} 超时) except httpx.HTTPStatusError as e: # 处理HTTP错误如4xx, 5xx raise ToolExecutionError(f工具返回HTTP错误: {e.response.status_code})第二层工具调用整体超时。即使HTTP请求没超时工具本身的业务逻辑也可能很慢。我们需要为单个工具的整个执行过程包括可能的预处理、后处理设置一个总超时。这可以通过asyncio.wait_for或给任务添加超时装饰器来实现。import asyncio from functools import wraps def tool_timeout(seconds: int): def decorator(func): wraps(func) async def wrapper(*args, **kwargs): try: return await asyncio.wait_for(func(*args, **kwargs), timeoutseconds) except asyncio.TimeoutError: # 记录哪个工具、哪些参数导致了超时 raise ToolExecutionTimeoutError(f工具 {func.__name__} 执行超过 {seconds} 秒) return wrapper return decorator tool_timeout(45) # 为该工具设置45秒总超时 async def complex_data_processing_tool(query: str): # 模拟一个耗时操作 await asyncio.sleep(50) # 这将触发超时 return {result: processed}第三层单轮ToolUse循环超时。一轮完整的ToolUse可能包含用户输入 - 模型思考并决定调用工具 - 执行工具 - 将工具结果返回给模型 - 模型生成最终回答。这个过程可能因为模型“思考”过久或工具链过长而超时。你需要为整个循环设置一个合理的总时间预算比如60秒。一旦超时立即中断并向用户返回一个友好的提示如“处理超时请简化您的问题或稍后再试”。实操心得超时时间的设置不是拍脑袋决定的。需要结合监控数据P95/P99响应时间和业务容忍度来动态调整。对于关键路径上的工具超时设得太短会导致不必要的失败设得太长又会拖累整体体验。一个实用的技巧是采用“阶梯式超时”首次调用用一个较短的超时如P95时间如果触发重试再适当放宽。3. 第二件套智能重试机制——应对瞬时故障超时帮我们快速失败但很多故障是瞬时的、可恢复的如网络抖动、下游服务短暂过载。盲目重试会加重下游负担不重试则会降低成功率。因此我们需要一个“智能”的重试机制。3.1 重试策略的核心要素一个完整的重试策略需要考虑以下几点重试条件When to Retry不是所有错误都值得重试。通常只对“幂等”的操作且错误类型为网络超时、连接拒绝、5xx服务器错误等进行重试。对于4xx客户端错误如参数错误、权限不足重试是无效的。重试间隔Retry Delay立即重试No Delay往往会把已经压力山大的下游服务直接击垮。常用的策略有固定间隔Fixed Delay每次等待相同时间简单但可能低效。指数退避Exponential Backoff等待时间随重试次数指数级增加如1s, 2s, 4s, 8s...。这是应对下游过载最有效的策略之一能给服务恢复留出时间。随机抖动Jitter在退避时间上增加一个随机值避免多个客户端在故障恢复后同时发起重试造成“惊群效应”。重试上限Max Retries必须设置一个最大重试次数如3次避免因持续失败陷入无限重试循环。3.2 实现一个带指数退避和抖动的重试装饰器我们可以利用tenacity或backoff这类成熟的库也可以自己实现一个简易版本以更清晰地理解其原理。import asyncio import random from typing import Callable, Any class RetryableError(Exception): 标记为可重试的错误 pass class NonRetryableError(Exception): 标记为不可重试的错误 pass def async_retry( max_retries: int 3, initial_delay: float 1.0, exponential_base: float 2.0, jitter: bool True, ): 异步重试装饰器支持指数退避和随机抖动。 def decorator(func: Callable): async def wrapper(*args, **kwargs): last_exception None for attempt in range(max_retries 1): # 1 包含第一次尝试 try: if attempt 0: # 不是第一次尝试需要等待 delay initial_delay * (exponential_base ** (attempt - 1)) if jitter: # 增加最多25%的随机抖动 delay * random.uniform(0.75, 1.25) print(f第{attempt}次重试失败等待{delay:.2f}秒后重试...) await asyncio.sleep(delay) return await func(*args, **kwargs) except NonRetryableError as e: # 明确不可重试的错误直接抛出 raise e except (ToolExecutionTimeoutError, ConnectionError, TimeoutError) as e: # 将这些视为可重试错误 last_exception e if attempt max_retries: print(f已达到最大重试次数{max_retries}放弃。) raise last_exception continue except Exception as e: # 其他未知错误默认不重试直接抛出 raise e # 理论上不会走到这里 raise last_exception return wrapper return decorator # 使用示例 async_retry(max_retries3, initial_delay1.0, exponential_base2.0, jitterTrue) async def call_unstable_api(): # 模拟一个不稳定的API前两次调用失败第三次成功 if not hasattr(call_unstable_api, _call_count): call_unstable_api._call_count 0 call_unstable_api._call_count 1 if call_unstable_api._call_count 3: raise ToolExecutionTimeoutError(模拟API超时) return {data: success on retry} # 在ToolUse循环中调用 async def tool_use_loop(): try: result await call_unstable_api() # 将结果交给模型继续处理... except Exception as e: # 处理最终失败 print(f工具调用最终失败: {e})踩坑实录我曾在一个项目中将重试机制用在了非幂等的写操作上比如创建一个订单结果因为网络问题导致重试下游服务收到了两次相同的请求生成了两个重复订单。这是一个严重的生产事故。黄金法则重试只用于读操作或明确幂等的写操作。对于非幂等操作要么使用唯一请求ID让下游服务做去重要么就避免重试。4. 第三件套熔断与降级——防止故障扩散当某个工具服务持续失败达到重试上限后依然失败很可能意味着该服务已经不可用或严重过载。如果继续让请求发往这个“瘫痪”的服务不仅会浪费资源还会导致调用方线程池被占满引发连锁故障这就是所谓的“雪崩效应”。熔断器Circuit Breaker模式就是为了防止这种情况。4.1 熔断器的工作原理熔断器有三种状态关闭Closed请求正常通过并统计失败率。打开Open当失败率超过阈值熔断器“跳闸”进入打开状态。此时所有对该服务的请求会立即失败快速失败不再真正发起调用。半开Half-Open经过一个设定的重置时间后熔断器进入半开状态允许少量试探请求通过。如果这些请求成功则认为服务已恢复熔断器关闭如果失败则继续保持打开状态。4.2 实现一个简易的熔断器虽然可以使用pybreaker这样的库但理解其原理很重要。下面是一个概念性的简化实现import time from enum import Enum from dataclasses import dataclass from typing import Optional class CircuitState(Enum): CLOSED CLOSED OPEN OPEN HALF_OPEN HALF_OPEN dataclass class CircuitBreaker: failure_threshold: int 5 # 连续失败多少次后熔断 reset_timeout: float 60.0 # 熔断后多久进入半开状态秒 half_open_success_threshold: int 2 # 半开状态下成功多少次后关闭 def __post_init__(self): self.state CircuitState.CLOSED self.failure_count 0 self.last_failure_time: Optional[float] None self.half_open_success_count 0 def record_success(self): if self.state CircuitState.HALF_OPEN: self.half_open_success_count 1 if self.half_open_success_count self.half_open_success_threshold: self._close() elif self.state CircuitState.CLOSED: self.failure_count 0 # 成功则重置连续失败计数 def record_failure(self): self.failure_count 1 self.last_failure_time time.time() if self.state CircuitState.CLOSED and self.failure_count self.failure_threshold: self._open() elif self.state CircuitState.HALF_OPEN: # 半开状态下试探请求也失败重新打开 self._open() def _open(self): print(熔断器状态变为 OPEN) self.state CircuitState.OPEN self.half_open_success_count 0 def _close(self): print(熔断器状态变为 CLOSED) self.state CircuitState.CLOSED self.failure_count 0 self.half_open_success_count 0 def allow_request(self) - bool: if self.state CircuitState.OPEN: # 检查是否到了该进入半开状态的时间 if time.time() - self.last_failure_time self.reset_timeout: print(熔断器进入 HALF_OPEN 状态进行试探) self.state CircuitState.HALF_OPEN return True # 允许一次试探请求 return False # 仍在熔断期拒绝请求 return True # 关闭或半开状态允许请求 # 使用熔断器包装工具调用 class ToolExecutor: def __init__(self, tool_name: str): self.tool_name tool_name self.circuit CircuitBreaker(failure_threshold3, reset_timeout30) async def execute(self, *args, **kwargs): if not self.circuit.allow_request(): raise CircuitBreakerOpenError(f工具 {self.tool_name} 熔断中请求被快速失败) try: result await self._call_actual_tool(*args, **kwargs) self.circuit.record_success() return result except (ToolExecutionError, ToolExecutionTimeoutError) as e: self.circuit.record_failure() raise e async def _call_actual_tool(self, *args, **kwargs): # 这里是实际的工具调用逻辑 await asyncio.sleep(0.1) # 模拟随机失败 if random.random() 0.7: # 70%失败率模拟服务异常 raise ToolExecutionTimeoutError(模拟工具调用失败) return {status: ok}4.3 服务降级Fallback当熔断器打开或者工具调用最终失败时我们不应该只是给用户抛出一个冷冰冰的错误。服务降级就是我们的“B计划”。降级策略可以有很多种返回缓存数据对于查询类工具可以返回上一次成功的缓存结果需标注可能过时。返回简化/静态结果提供一个功能简化版的结果或者一个友好的提示页面。切换备用服务如果有同质的备用服务或接口可以切换过去。在ToolUse场景中降级可以表现为当获取实时天气的API失败时模型可以转而回答“目前无法获取实时天气但根据历史数据这个季节通常...”当计算工具失败时模型可以回答“计算服务暂时不可用您可以尝试使用公式 XXXX 进行估算”。在代码中降级逻辑通常和熔断、重试结合在一起。async def call_tool_with_fallback(tool_func, fallback_func, *args, **kwargs): 带有降级策略的工具调用 try: # 这里可以集成重试逻辑 return await tool_func(*args, **kwargs) except (ToolExecutionError, CircuitBreakerOpenError) as e: print(f主工具调用失败: {e}, 尝试降级方案。) # 执行降级逻辑 return await fallback_func(*args, **kwargs) async def get_weather_fallback(city: str): 天气API的降级方案返回静态提示 return { source: fallback, message: f实时天气数据暂时不可用。{city}的典型气候信息可在相关气象网站查询。 }5. 第四件套结构化输出验证——守住数据的最后一道门大模型并不总是可靠的“程序员”。即使你定义了再清晰的Function Calling Schema模型也可能返回格式错误、字段缺失、甚至类型不匹配的JSON。更糟糕的是下游工具API返回的数据也可能不遵守契约。直接把这些“脏数据”传递给下一个环节或返回给用户轻则导致功能异常重则可能引发安全漏洞如代码注入。因此对输入用户输入需谨慎处理和输出进行严格的验证是必须的。5.1 模型输出验证超越JSON解析首先是对模型返回的“工具调用请求”进行验证。这不仅仅是json.loads()那么简单。基础语法验证捕获JSON解码错误。模式Schema验证验证返回的JSON对象是否完全符合你定义的函数调用结构。包括必需字段name工具名、arguments是否存在。字段类型arguments是否是一个对象其中的参数类型是否符合预期字符串、数字、布尔值、数组等。参数范围对于枚举型参数值是否在允许的列表内。业务逻辑验证比如开始日期不能晚于结束日期。推荐使用Pydantic库来进行强大的数据验证和序列化。它能将JSON数据自动转换为类型安全的Python对象并在转换过程中完成所有验证。from pydantic import BaseModel, Field, validator from typing import Literal, Optional from datetime import date # 1. 定义工具参数模型 class WeatherQueryArgs(BaseModel): location: str Field(..., description城市名称如北京) unit: Literal[celsius, fahrenheit] Field(celsius, description温度单位) class CalculatorArgs(BaseModel): operation: Literal[add, subtract, multiply, divide] a: float b: float validator(b) def check_division_by_zero(cls, v, values): if values.get(operation) divide and v 0: raise ValueError(除数b不能为0) return v # 2. 定义模型返回的工具调用结构 class ToolCallRequest(BaseModel): name: str Field(..., description要调用的工具名称) arguments: dict Field(..., description工具调用参数) def parse_arguments(self, tool_model_map: dict): 根据工具名将arguments字典解析为具体的参数模型 if self.name not in tool_model_map: raise ValidationError(f未知的工具名: {self.name}) ModelClass tool_model_map[self.name] try: # Pydantic会自动进行类型转换和验证 return ModelClass(**self.arguments) except Exception as e: # 记录详细的验证错误信息便于调试和优化Prompt raise ValidationError(f工具{self.name}参数验证失败: {e}) # 3. 在ToolUse循环中使用 tool_model_map { get_weather: WeatherQueryArgs, calculator: CalculatorArgs, } def validate_and_parse_model_output(raw_text: str) - ToolCallRequest: 解析并验证模型返回的文本提取工具调用请求。 # 第一步尝试从文本中提取JSON模型可能返回带Markdown或说明的文本 import json import re # 简单的JSON提取实际中可能需要更复杂的解析 json_match re.search(r\{.*\}, raw_text, re.DOTALL) if not json_match: raise ValidationError(无法从模型输出中提取JSON结构) try: data json.loads(json_match.group()) except json.JSONDecodeError as e: raise ValidationError(f模型返回的JSON格式无效: {e}) # 第二步使用Pydantic模型验证 try: tool_call ToolCallRequest(**data) # 第三步进一步验证参数是否符合具体工具的要求 parsed_args tool_call.parse_arguments(tool_model_map) # 可以将解析后的参数对象附加到tool_call上方便后续使用 tool_call.parsed_arguments parsed_args return tool_call except Exception as e: # 这里可以记录下raw_text和data用于分析模型为何输出异常格式优化Prompt或fine-tuning print(f验证失败原始输出: {raw_text[:200]}...) raise ValidationError(f工具调用请求验证失败: {e})5.2 工具返回结果验证同样对工具执行后返回的结果也要进行验证确保其格式和内容在模型能够安全、正确理解的范围内。例如一个返回股票价格的工具应该确保结果是数字而不是一串错误信息。这同样可以用Pydantic模型来约束。class WeatherToolResult(BaseModel): location: str temperature: float unit: str condition: str humidity: Optional[int] None # 可以添加更多字段和验证器 # 在工具函数内部对返回给模型的数据进行验证 async def get_weather_tool(location: str, unit: str) - dict: # ... 调用真实API ... raw_api_response await call_weather_api(location, unit) # 验证和清洗API响应 try: validated_result WeatherToolResult( locationlocation, temperatureraw_api_response.get(temp), unitunit, conditionraw_api_response.get(weather, [{}])[0].get(description, unknown), humidityraw_api_response.get(humidity) ) # 返回给模型的是经过验证的字典 return validated_result.dict() except Exception as e: # API返回了意外格式记录日志并返回一个安全的错误信息结构 logger.error(f天气API返回格式异常: {e}, 原始响应: {raw_api_response}) # 返回一个结构化的错误信息让模型知道工具调用出了问题 return { error: True, message: 获取天气数据时遇到问题请稍后再试。, location: location }经验之谈验证失败时的处理策略很重要。不要简单地让整个对话崩溃。对于模型输出验证失败可以设计一个“修复循环”将验证错误信息作为系统提示的一部分要求模型重新生成正确的格式。对于工具结果验证失败则应该返回一个结构化的错误信息让模型有能力向用户解释“某个服务暂时不可用”而不是胡言乱语。这提升了系统的健壮性和用户体验。6. 第五件套可观测性——为系统装上眼睛和耳朵前面四件套都是在“做事”而可观测性Observability是让你“看清”系统到底在怎么做事。没有完善的日志、指标和追踪一个黑盒般的ToolUse循环在生产环境就是一场噩梦。你无法知道瓶颈在哪无法快速定位故障更谈不上优化。6.1 结构化日志Structured Logging告别print语句和杂乱无章的文本日志。结构化日志将日志信息以键值对通常是JSON的形式输出便于后续的收集、筛选和分析。需要记录的关键信息请求/会话ID唯一标识一次用户对话或请求用于串联所有相关日志。用户输入脱敏后记录。模型请求与响应记录发送给模型的Prompt可采样或哈希处理以避免数据过大和返回的完整响应。这对于调试模型“诡异”行为至关重要。工具调用详情工具名称、输入参数、开始时间、结束时间、耗时、成功/失败状态、错误信息、返回结果可摘要或哈希。验证错误记录具体的验证失败信息。重试与熔断事件记录重试次数、退避时间、熔断器状态变化。import json import logging import uuid from contextvars import ContextVar from datetime import datetime # 使用ContextVar来传递请求上下文如request_id request_id_ctx: ContextVar[str] ContextVar(request_id, default) class StructuredLogger: def __init__(self, name): self.logger logging.getLogger(name) def _log(self, level: str, event: str, **extra): log_entry { timestamp: datetime.utcnow().isoformat() Z, level: level.upper(), request_id: request_id_ctx.get(), event: event, **extra } # 输出为JSON字符串方便ELK等系统收集 self.logger.log(getattr(logging, level.upper()), json.dumps(log_entry, ensure_asciiFalse)) def info(self, event: str, **extra): self._log(INFO, event, **extra) def error(self, event: str, **extra): self._log(ERROR, event, **extra) def warning(self, event: str, **extra): self._log(WARNING, event, **extra) # 在ToolUse循环入口处设置request_id async def handle_user_query(user_input: str): request_id str(uuid.uuid4()) token request_id_ctx.set(request_id) logger StructuredLogger(__name__) logger.info(tool_use_loop_started, user_input_previewuser_input[:50]) try: # ... 循环逻辑 ... logger.info(model_called, tool_callvalidated_tool_call.dict()) # ... 调用工具 ... logger.info(tool_executed, tool_nametool_name, duration_msduration, successTrue) except ValidationError as e: logger.error(validation_failed, errorstr(e), raw_outputraw_text[:200]) except ToolExecutionError as e: logger.error(tool_execution_failed, tool_nametool_name, errorstr(e)) finally: request_id_ctx.reset(token) logger.info(tool_use_loop_finished)6.2 关键指标Metrics监控日志用于事后分析指标用于实时告警和性能评估。你需要监控吞吐量与延迟tooluse_requests_total总请求数。tooluse_request_duration_seconds请求处理耗时分布Histogram。tooluse_tokens_total消耗的总Token数分输入/输出。成功率与错误率tooluse_success_total成功完成最终给出有效回答的请求数。tooluse_failure_total失败的请求数按失败类型分类validation_error, tool_error, model_error, timeout。tooluse_tool_call_total工具调用总数按工具名分类。tooluse_tool_call_failure_total工具调用失败数按工具名和错误类型分类。可靠性相关指标tooluse_retries_total重试总次数。circuit_breaker_state熔断器状态Gauge 0Closed, 1Open, 2Half-Open。tooluse_fallback_invoked_total降级策略触发次数。使用像Prometheus这样的监控系统来暴露这些指标并配置Grafana仪表盘进行可视化。当错误率飙升或P99延迟超过阈值时能第一时间收到告警。6.3 分布式追踪Distributed Tracing在一个复杂的ToolUse循环中一次用户请求可能触发多次模型调用和多个工具调用这些调用可能分布在不同的服务上。分布式追踪如OpenTelemetry能帮你还原出一次请求的完整生命周期视图清晰地看到时间消耗在哪个环节是模型生成慢还是某个工具API慢是性能调优和根因分析的利器。为你的AI服务、模型网关、各个工具服务都集成OpenTelemetry SDK为每次调用生成唯一的Trace ID并传播下去。你就能在Jaeger或Zipkin这样的界面上看到一幅完整的“调用链火焰图”。7. 五件套的整合一个生产级ToolUse循环的代码骨架最后让我们把以上所有组件整合到一个简化的、但体现了生产级思维的ToolUse循环骨架中。请注意这是一个概念性示例真实环境需要更完善的依赖注入、配置管理、和错误处理。import asyncio import random from typing import Dict, Any, Optional # 假设已导入必要的库pydantic, httpx, logging, opentelemetry等 # ---------- 1. 定义核心异常与模型 ---------- class ToolUseError(Exception): ToolUse循环基础异常 pass class ValidationError(ToolUseError): pass class ToolExecutionError(ToolUseError): pass class ToolExecutionTimeoutError(ToolExecutionError): pass class CircuitBreakerOpenError(ToolExecutionError): pass # (此处省略之前定义的 Pydantic 模型: ToolCallRequest, WeatherQueryArgs等) # (此处省略之前定义的 CircuitBreaker, StructuredLogger, async_retry 装饰器) # ---------- 2. 核心工具执行器集成重试、熔断、降级 ---------- class ProductionToolExecutor: def __init__(self, tool_name: str, circuit_breaker: CircuitBreaker, fallback_funcNone): self.tool_name tool_name self.circuit circuit_breaker self.fallback fallback_func self.logger StructuredLogger(ftool.{tool_name}) async_retry(max_retries2, initial_delay1.0) async def execute(self, parsed_args: BaseModel) - Dict[str, Any]: 执行工具集成重试 # 检查熔断器 if not self.circuit.allow_request(): self.logger.warning(circuit_breaker_open, toolself.tool_name) raise CircuitBreakerOpenError(f工具 {self.tool_name} 已熔断) start_time asyncio.get_event_loop().time() try: # 实际调用工具这里用模拟函数代替 result await self._call_tool_internal(parsed_args) duration (asyncio.get_event_loop().time() - start_time) * 1000 self.circuit.record_success() self.logger.info(tool_execution_success, toolself.tool_name, duration_msround(duration, 2)) return result except (ToolExecutionError, ToolExecutionTimeoutError) as e: duration (asyncio.get_event_loop().time() - start_time) * 1000 self.circuit.record_failure() self.logger.error(tool_execution_failed, toolself.tool_name, errorstr(e), duration_msround(duration, 2)) # 如果定义了降级函数则执行降级 if self.fallback is not None: self.logger.info(fallback_triggered, toolself.tool_name) try: return await self.fallback(parsed_args) except Exception as fallback_e: self.logger.error(fallback_failed, toolself.tool_name, errorstr(fallback_e)) # 否则重新抛出异常 raise async def _call_tool_internal(self, args: BaseModel) - Dict[str, Any]: 模拟实际工具调用包含超时控制 # 这里应使用带有超时设置的HTTP客户端 await asyncio.sleep(random.uniform(0.1, 1.5)) # 模拟网络延迟 # 模拟随机失败 if random.random() 0.2: raise ToolExecutionTimeoutError(内部工具调用超时) # 模拟成功返回 return {result: fTool {self.tool_name} executed with args: {args.dict()}} # ---------- 3. 生产级ToolUse循环主逻辑 ---------- class ProductionToolUseLoop: def __init__(self, llm_client, available_tools: Dict[str, ProductionToolExecutor]): self.llm llm_client self.tools available_tools self.logger StructuredLogger(tool_use_loop) self.tool_model_map { # 工具名到参数模型的映射 get_weather: WeatherQueryArgs, calculator: CalculatorArgs, } async def run(self, user_input: str, session_id: str) - str: 运行一轮完整的ToolUse循环 self.logger.info(loop_started, session_idsession_id, user_input_previewuser_input[:100]) loop_start asyncio.get_event_loop().time() try: # Step 1: 调用大模型获取工具调用请求 model_response await self._call_llm_with_timeout(user_input) self.logger.info(model_response_received, response_previewmodel_response[:200]) # Step 2: 验证并解析模型输出 tool_call_request validate_and_parse_model_output(model_response) self.logger.info(tool_call_validated, tool_nametool_call_request.name) # Step 3: 获取对应的工具执行器并执行 tool_name tool_call_request.name if tool_name not in self.tools: raise ValidationError(f请求调用了未注册的工具: {tool_name}) tool_executor self.tools[tool_name] parsed_args tool_call_request.parse_arguments(self.tool_model_map) tool_result await tool_executor.execute(parsed_args) # Step 4: 将工具结果再次喂给模型生成最终回答 final_prompt self._construct_final_prompt(user_input, tool_call_request, tool_result) final_response await self._call_llm_with_timeout(final_prompt) total_duration (asyncio.get_event_loop().time() - loop_start) * 1000 self.logger.info(loop_succeeded, session_idsession_id, total_duration_msround(total_duration, 2), final_response_previewfinal_response[:100]) return final_response except ValidationError as e: self.logger.error(loop_failed_validation, session_idsession_id, errorstr(e)) # 可以在这里尝试让模型修复格式或者直接返回友好错误 return 抱歉我暂时无法理解您的请求请尝试换一种方式提问。 except (ToolExecutionError, CircuitBreakerOpenError) as e: self.logger.error(loop_failed_tool, session_idsession_id, errorstr(e)) return 抱歉处理您请求所需的某个服务暂时不可用请稍后再试。 except asyncio.TimeoutError: self.logger.error(loop_timeout, session_idsession_id) return 处理请求超时您的问题可能过于复杂请尝试简化或稍后重试。 except Exception as e: self.logger.error(loop_failed_unexpected, session_idsession_id, errorstr(e), exc_infoTrue) # 未知异常返回通用错误避免泄露内部信息 return 系统处理时遇到意外错误请联系管理员。 async def _call_llm_with_timeout(self, prompt: str, timeout: int 30) - str: 调用LLM并施加超时控制 try: # 假设llm.call是异步的 return await asyncio.wait_for(self.llm.call(prompt), timeouttimeout) except asyncio.TimeoutError: self.logger.error(llm_call_timeout) raise def _construct_final_prompt(self, user_input, tool_call, tool_result): # 构造包含工具调用和结果的Prompt给模型 return f 用户问题{user_input} 我决定调用工具 {tool_call.name}参数是{tool_call.arguments}。 工具返回的结果是{tool_result}。 请根据以上信息生成对用户的最终回答。 # ---------- 4. 初始化与运行示例 ---------- async def main(): # 初始化熔断器 weather_circuit CircuitBreaker(failure_threshold2, reset_timeout20) calc_circuit CircuitBreaker(failure_threshold2, reset_timeout20) # 初始化工具执行器 tools { get_weather: ProductionToolExecutor( get_weather, weather_circuit, fallback_funclambda args: {result: 降级天气服务暂不可用。} ), calculator: ProductionToolExecutor(calculator, calc_circuit), } # 模拟一个LLM客户端实际应替换为OpenAI/Azure等SDK class MockLLM: async def call(self, prompt): await asyncio.sleep(0.5) # 模拟模型返回一个工具调用请求 if 天气 in prompt: return {name: get_weather, arguments: {location: 上海, unit: celsius}} elif 计算 in prompt: return {name: calculator, arguments: {operation: add, a: 5, b: 3}} else: return {name: get_weather, arguments: {location: 北京}} # 缺省参数会触发验证错误 llm_client MockLLM() loop_processor ProductionToolUseLoop(llm_client, tools) # 模拟处理用户请求 test_inputs [上海天气怎么样, 帮我计算一下5加3, 无效请求测试] for query in test_inputs: print(f\n 处理用户输入: {query} ) try: result await loop_processor.run(query, session_idtest_session_123) print(f最终回答: {result}) except Exception as e: print(f循环处理失败: {e}) if __name__ __main__: asyncio.run(main())这个骨架代码展示了如何将超时、重试、熔断、降级、验证和日志观测性编织在一起构建出一个具备基本生产韧性的ToolUse循环。在实际项目中你需要根据具体的框架如LangChain、Semantic Kernel或自研架构和基础设施Kubernetes、云服务进行适配和扩展。记住可靠性不是一项功能而是一个贯穿设计、编码、测试和运维全过程的系统工程。把这“五件套”内化为你的开发习惯你的AI应用才能真正地从实验室走向生产线。
返回列表