Instructor 架构深度解析:patch、Retry 重试层与响应分发机制的完整执行链路
Instructor 架构深度解析patch、Retry 重试层与响应分发机制的完整执行链路【免费下载链接】instructorstructured outputs for llms项目地址: https://gitcode.com/GitHub_Trending/in/instructor本篇技术指南基于 docs/architecture.md结合当前仓库的 v2 核心源码系统讲解 Instructor 库的内部架构与设计决策从patch()包装 provider 客户端开始到 tenacity 重试层、Mode分发器、Provider Handler 的请求整形与 reask 纠错再到最终 Pydantic 模型的解析与_raw_response挂载。读完本文你将掌握 Instructor 的最小同步/异步调用路径、流式Iterable与部分对象Partial模式如何接入分发器、并行工具模式的约束、Hook 埋点与重试语义以及如何为新的 provider 扩展注册一套完整的 handler。整体架构与高层执行流程Instructor 的核心思想并不复杂把请求整形 重试 响应解析这一整套管道透明地叠加在原生 LLM provider 客户端之上。开发者只需要传入一个response_modelPydantic 模型剩下的 schema 生成、工具参数注入、JSON 解析、校验失败重试全部由库内部完成。原文给出了一个非常清晰的时序图sequenceDiagram完整描述了从用户代码到最终模型实例的全过程U-I: chat.completions.create(response_model..., **kwargs) Note over I: patch() 用 cache/templating 和 retry 包装 create() I-R: retry_sync/async(funccreate, max_retries, strict, mode, hooks) loop attempts R-C: create(**prepared_kwargs) C--R: raw responseprovider 特有 R-D: process_response(_async)(response, response_model, mode, stream) alt Streaming/Partial D-M: Iterable/Partial.from_streaming_response(_async) D--R: Iterable/Partial 模型或条目列表 else Standard D-H: provider mode handlerformat/parse 选择 H--D: 调整后的 response_model/new_kwargs如需 D-M: response_model.from_response(...) M--D: 解析后的模型挂载 _raw_response D--R: model或适配后的简单类型 end R--I: parsed model end I--U: final model实例上附带 _raw_response Note over R,H: 遇到 validation/JSON 错误 → 走 reask 路径 R-H: handle_reask_kwargs(..., exception, failed_attempts) H--R: 用于下一次尝试的新 kwargs/messages在这个流程中四个参与者各自承担明确的职责patch()包装 provider 的create叠加 cache 查找/保存、templating、strict 模式、hooks 与 retry。Retry 层tenacity执行 provider 调用、发射 hooks、累计 usage、处理 validation/JSON 错误并触发 reask然后重新尝试。Dispatcherprocess_response按Mode选择正确的解析路径处理多模态消息转换并把_raw_response挂载到返回的模型上。Provider Handlers负责 provider/mode 特异的请求整形与 reask 准备。源码中的Dispatcher从 process_response 到 ModeRegistry值得说明的是当前仓库的 v2 架构已经将原文中Dispatcherprocess_response的职责进一步细分instructor/v2/core/registry.py中的ModeRegistry成为中央分发枢纽它把(Provider, Mode)二元组映射到一组 handler# instructor/v2/core/registry.py dataclass class ModeHandlers: request_handler: RequestHandler # 请求整形tools/schema 注入 reask_handler: ReaskHandler # 校验失败后的纠错 response_parser: ResponseParser # 响应解析 stream_extractor: StreamExtractor | None None stream_extractor_async: AsyncStreamExtractor | None None message_converter: MessageConverter | None None # 多模态消息转换 template_handler: TemplateHandler | None NoneModeRegistry支持惰性加载与动态注册handler 首次被查询时才 import 对应模块并用threading.Lock保护并发首次调用场景见 registry.py 的注释说明。而retry_sync_v2在发起调用前会先执行RegistryValidationMixin.validate_mode_registration(provider, mode)若该模式未注册会抛出带可用模式列表的RegistryError见 exceptions.py从源头上拦截模式与 provider 不匹配的误用。patch()入口处的统一包装patch()是整条管道的入口。在 v2 中instructor/v2/core/patch.py的patch_v2接受原始 provider 调用函数例如client.messages.create、Provider枚举与Mode通过is_async(func)判断后生成同步或异步包装器# instructor/v2/core/patch.py节选 def patch_v2(func, provider, mode, default_modelNone): RegistryValidationMixin.validate_mode_registration(provider, mode) func_is_async is_async(func) if func_is_async: return _create_async_wrapper(func, provider, mode, default_model) return _create_sync_wrapper(func, provider, mode, default_model)包装后的create在instructor/v2/core/client.py的Instructor/AsyncInstructor中对外暴露。以同步create为例它的职责链清晰可循通过reject_async_validators(response_model)拒绝在同步路径中混入异步 validatorhandle_kwargs将构造客户端时传入的默认 kwargs 与本次调用 kwargs 合并调用侧优先见 client.py合并客户端级 hooks 与单次调用 hooksself.hooks hooks最终委托给create_fn即真正被包装过的重试管道。# instructor/v2/core/client.pycreate 核心逻辑节选 kwargs self.handle_kwargs(kwargs) if token_budget is not None: kwargs[token_budget] token_budget combined_hooks self.hooks if hooks is not None: combined_hooks self.hooks hooks return self.create_fn( response_modelresponse_model, messagesmessages, max_retriesmax_retries, contextcontext, strictstrict, hookscombined_hooks, **kwargs, )客户端还通过__getattr__把未拦截的属性透传给底层原始 clientchat、messages、create除外因此client.chat.completions.create(...)这类 OpenAI 原生链式调用依然可以无缝工作见 client.py。Mode分发器的路由键Mode枚举定义了每种 provider 与解析策略的组合是 Dispatcher 的路由键。instructor/v2/core/mode.py中按 provider 分组维护了全套模式OpenAI 的TOOLS/TOOLS_STRICT/JSON/MD_JSON/JSON_SCHEMA/RESPONSES_TOOLS、Anthropic 的ANTHROPIC_TOOLS/ANTHROPIC_JSON/ANTHROPIC_PARALLEL_TOOLS、Vertex 与 Gemini 的VERTEXAI_TOOLS/GEMINI_TOOLS/GENAI_TOOLS等等并提供了三个分类辅助方法Mode.tool_modes()所有基于工具调用的模式Mode.json_modes()所有基于 JSON 输出的模式Mode.parallel_modes()并行工具模式PARALLEL_TOOLS、ANTHROPIC_PARALLEL_TOOLS、VERTEXAI_PARALLEL_TOOLS。仓库还维护了一张DEPRECATED_TO_CORE映射表将历史遗留模式如FUNCTIONS、TOOLS_STRICT、ANTHROPIC_TOOLS等归一化到核心模式并通过warn_deprecated_mode在会话内只警告一次避免日志刷屏见 mode.py。从源码结构看v2 的方向是provider 由客户端决定mode 只管解析策略因此 Provider 枚举与 Mode 枚举在注册表中是解耦的。最小调用路径同步与异步原文给出了两条最小代码路径这也是理解全库的最短切入点。同步路径import openai import instructor from pydantic import BaseModel class User(BaseModel): name: str age: int client instructor.from_provider(openai/gpt-5-nano) model client.create( modelgpt-4o-mini, messages[{role: user, content: {name: Ada, age: 37}}], response_modelUser, # 触发 schema/tool 接线与解析 max_retries3, # tenacity 支撑的校验重试 strictTrue, # 若 provider 支持则启用严格 JSON 解析 ) # 需要时访问原始 provider 响应 raw model._raw_response异步路径仅需换用async_clientTrue与awaitimport asyncio import openai import instructor from pydantic import BaseModel class User(BaseModel): name: str age: int async def main(): aclient instructor.from_provider(openai/gpt-5-nano, async_clientTrue) model await aclient.create( modelgpt-4o-mini, messages[{role: user, content: {\name\: \Ada\, \age\: 37}}], response_modelUser, max_retries3, strictTrue, ) print(model) asyncio.run(main())from_provider工厂位于 instructor/v2/auto_client.py通过openai/gpt-5-nano这样的provider/model字符串自动推导 Provider 与默认模型其签名支持async_client布尔开关据此返回Instructor或AsyncInstructor。Response._normalize_messages还允许直接传入字符串自动包装为单条 user 消息或在messages与input之间二选一二者同时传入会抛TypeError见 client.py。关于strict参数结合 openai/handlers.py 的实现可以看到prepare_request会从 kwargs 中弹出strict为 true 时在生成的 OpenAI schema 中写入strict: True随后将 schema 作为唯一 function tool 注入tools并把tool_choice固定为该 function。代码中还特意对generate_openai_schema的 lru_cache 返回值做了浅拷贝避免把strict字段永久污染到后续所有针对同一模型的调用——这是一个值得注意的缓存安全细节。流式、部分对象与并行工具流式 Iterablecreate_iterablecreate_iterable(response_modelModel)内部会强制streamTrue并把response_model包装为Iterable[response_model]最终由IterableBase.from_streaming_response(_async)负责把流式 chunk 组装成条目的生成器同步或异步生成器异步。对应实现见 client.pykwargs[stream] True response_model Iterable[response_model] # type: ignore return self.create_fn( messagesmessages, response_modelresponse_model, max_retriesmax_retries, contextcontext, strictstrict, hookscombined_hooks, **kwargs, )使用方式for item in client.create_iterable(messages..., response_modelMyModel): print(item)异步侧AsyncInstructor.create_iterable有额外一层判断当create收到的response_model本身就是Iterable[T]类型且当前 mode 不属于Mode.parallel_modes()时会自动转投到create_iterable见 client.py这让两种调用风格可以互相兼容。部分对象create_partialcreate_partial(response_modelModel)用于在流式过程中拿到逐字段填充的部分模型内部把模型包装为Partial[Model]并强制streamTrue每收到一个 chunk 就返回一次更新后的部分对象。实现同样位于 client.pyPartial类型定义在 instructor/v2/dsl/partial.py。for partial in client.create_partial(messages..., response_modelMyModel): # partial 中会包含已到达的字段 pass并行工具Mode.PARALLEL_TOOLS当需要一次请求中触发多个工具调用时使用Mode.PARALLEL_TOOLS并将response_model声明为模型列表如[PersonInfo, EventInfo]。注意并行工具模式不支持流式。from instructor.mode import Mode result client.create( modelgpt-4o, messages[{role: user, content: Extract person and event info.}], response_model[PersonInfo, EventInfo], modeMode.PARALLEL_TOOLS, )从源码看并行工具路径有两个关键点请求侧由instructor/v2/dsl/parallel.py的handle_parallel_model把多个模型展开为多个 tools 并设置tool_choiceauto见 handlers.py解析侧parse_response遍历response.choices[0].message.tool_calls按工具名查找type_registry用model_validate_json逐个解析并 yield见 handlers.py。之所以流式不支持从该代码路径可以推断并行工具需要拿到完整的tool_calls列表才能按工具名分发流式 chunk 无法保证工具调用的完整性。Hooks 与重试语义Hook 事件一览Hook 是观测与插桩整条调用链路的统一入口。HookName枚举定义在 instructor/v2/core/hooks.py共六个事件HookName触发时机completion:kwargs每次 provider 调用之前completion:response每次 provider 调用之后completion:error非校验类的完成错误completion:last_attempt重试序列即将停止时completion:usage有可用 usage 时携带累计用量快照parse:error校验/JSON 解析错误注册方式如下client.on内部委托给Hooks.on也支持传字符串形式的 hook 名from instructor.core.hooks import HookName client.on(HookName.COMPLETION_KWARGS, lambda **kw: print(KWARGS, kw)) client.on(HookName.PARSE_ERROR, lambda e: print(PARSE, e))Hooks类用defaultdict(list)按事件名维护处理器列表每个事件可挂多个 handler客户端级 hooks 与单次调用的hooks参数会通过运算符合并见 client.py。completion:error/completion:last_attempt的 handler 协议签名固定为(error, *, attempt_number, max_attempts, is_last_attempt)方便做精确的重试状态观测见 hooks.py。重试管道tenacity 与 reask重试逻辑的核心在instructor/v2/core/retry.py的retry_sync_v2异步对应retry_async_v2。关键机制停止条件max_retries为整数时构造stop_after_attempt(max(max_retries, 0) 1)若 kwargs 中带了数值型timeout还会追加stop_after_delay(timeout)与之取并集max_retries也可以直接传一个 tenacityRetrying实例做完全自定义见 retry.py。可重试异常仅ValidationError、json.JSONDecodeError、AsyncValidationError、ResponseParsingError四类会被重试_RETRYABLE_PARSE_ERRORS其余异常直接向上抛出并在completion:errorhook 中带出 attempt 元数据。每次尝试先 emitcompletion:kwargs→ 调用 provider → emitcompletion:response→ 累计 usageupdate_total_usage与 openaiCompletionUsage结构对齐Anthropic 模式则使用其自有 usage 初始化→ 用注册表中的response_parser解析。解析失败走 reask把FailedAttempt(attempt_number, exception, completion)记入failed_attemptsemitparse:error然后调用handlers.reask_handler(kwargs, response, exception)生成新的 kwargs/messages把错误反馈追加进对话再进入下一轮尝试。重试预算token_budget参数用于限制校验重试的累计 token 消耗超出预算时_budget_error会提前终止并抛出TokenBudgetError而不是继续盲目重试。最终若所有尝试失败抛出InstructorRetryException其中包含failed_attempts历次失败详情、last_completion最后一次原始响应、total_usage累计用量、messages用于复现的对话以及create_kwargs完整的调用参数便于复现与排查见 retry.py。以 OpenAI TOOLS 模式为例reask 与解析由 instructor/v2/providers/openai/handlers.py 中注册到Mode.TOOLS的 handler 类实现prepare_request注入 tools/tool_choicehandle_reask委托给reask_tools(kwargs, response, exception)parse_response则区分流式解析、finish_reason length的不完整输出抛IncompleteOutputException、并行工具生成器与标准工具调用解析四条子路径。多模态消息转换发生在哪里对于需要多模态能力的模式消息会经由processing.multimodal.convert_messagesv2 中为 instructor/v2/core/multimodal.py转换。特定 handler/mode 可以启用 Image/Audio/PDF 自动检测把字符串路径、URL、data URI 或字节流统一转成 provider 可消费的 payloadImage.autodetect会依次识别 base64 字符串、http(s)://、gs://、本地文件路径见 multimodal.pyMIME 类型白名单覆盖了常见的图片jpeg/png/gif/webp、音频wav/mp3/m4a/flac/opus 等与 PDFapplication/pdf远程内容通过instructor/v2/core/remote.py的fetch_remote_content拉取并附带MAX_IMAGE_BYTES/MAX_AUDIO_BYTES/MAX_PDF_BYTES大小上限约束。消息转换器message_converter作为ModeHandlers的可选成员注册因此在 registry 查询层面与请求/解析 handler 保持同一套生命周期。错误处理一览综合原文与源码错误处理的完整图景如下校验或 JSON 解码错误四类可重试异常触发 reask 路径Reask handlerhandle_reask/handle_reask_kwargs在 kwargs 中追加/调整带错误反馈的 message让下一次尝试能自我纠正IncompleteOutputException在响应因finish_reason length被截断时抛出handlers.py同样属于不可盲目重试的终止性信号在 retry 管道中被原样透传RegistryError在(Provider, Mode)组合未注册时于入口即失败并给出可用模式列表TokenBudgetError在token_budget被耗尽时提前终止重试所有重试用尽则抛InstructorRetryException携带failed_attempts、最后一次 completion、usage 汇总与可复现的create_kwargs。扩展性说明新增 Provider 的正确姿势原文最后给出了清晰的扩展约束源码印证了这套约定的落地方式新 provider 需要提供响应解析与 reask 处理的 utils并通过mode_registry.register(...)注册到(Provider, Mode)键上见 registry.py也可用register_mode_handler装饰器按 mode 声明 handler 类如 openai/handlers.py 中各 handler 类的用法。大部分 JSON/tool 模式是共享的例如DEPRECATED_TO_CORE把大量历史 provider 模式归一化到Mode.TOOLS/Mode.JSON/Mode.JSON_SCHEMA/Mode.MD_JSON/Mode.PARALLEL_TOOLS新 provider 应优先复用这些核心模式的 handler。provider 特有逻辑留在 provider utils 中中央 Dispatcherregistry retry只负责路由与编排不内嵌任何 provider 细节。Provider 规格HANDLER_SPECS/PROVIDER_SPECS集中在 instructor/v2/core/provider_specs.py新 provider 主要补齐这一层声明与对应 handlers 即可。小结Instructor 的架构可以概括为一条职责单一、层层委托的管道patch()负责入口包装tenacity 重试层负责调用 → 解析 → 失败则 reask的循环ModeRegistry负责按(Provider, Mode)分发到具体的 request/reask/response handler最后统一把_raw_response挂到解析出的 Pydantic 模型上。理解这条链路后无论是排查一次重试了却仍解析失败的问题、为内部调用接入completion:kwargs/parse:error观测还是为一个新 provider 注册 handler都能精准定位到对应的代码层次入口看 instructor/v2/core/patch.py 与 instructor/v2/core/client.py重试语义看 instructor/v2/core/retry.py路由分发看 instructor/v2/core/registry.py模式全集看 instructor/v2/core/mode.py。【免费下载链接】instructorstructured outputs for llms项目地址: https://gitcode.com/GitHub_Trending/in/instructor创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考