数据截至 (上游 commit 3f15dc32871c)
横切层:@client 装饰器里的日志、缓存、成本与异常
30 秒导读: LiteLLM 的每一个公开函数(
completion/acompletion/embedding/responses…)都被同一个装饰器client包住。装饰器负责的事和"调哪家模型"毫无关系:生成调用 ID、装上日志回调、查预算、查缓存、算成本、发回调、把各家五花八门的错误收敛成同一组异常。本章讲的就是这层外壳,不讲壳里面的翻译和 HTTP(那是 02 和 03)。
1. 这是什么(零基础也能懂)
一句话定义: 横切层是一段被复用 177 次的函数包装代码,它在"真实调用"前后做所有与 provider 无关的事。
它解决什么问题。 LiteLLM 要支持 100+ 家 LLM API。如果不抽这一层,就会出现这种局面:
- OpenAI 的调用路径里写一遍"记日志、算钱、查缓存"
- Anthropic 的调用路径里再写一遍
- Bedrock 的调用路径里再写一遍……
于是同一个 bug 要修 100 次,同一个功能(比如新增一个 Langfuse 回调)要接 100 次。
它的解法: 把这些事全部提到函数外面,用一个装饰器统一做。真源码里就是一行:
# litellm/main.py:4899-4901,真实源码
@tracer.wrap()
@client
def completion(
整个包里一共 177 处 @client、分布在 23 个文件(grep -rn "^@client" litellm --include="*.py"):litellm/main.py 占 11 处(acompletion:388、completion:4901、embedding、text_completion、image_generation、transcription…),其余分布在 responses/main.py、rerank_api/main.py、batches/main.py、files/main.py、google_genai/main.py 等各能力子包的 main.py 里。
用起来什么样。 用户什么都不用做——他调的是普通函数,横切层是隐形的:
# 示意,非源码
import litellm
litellm.success_callback = ["langfuse"] # 装一个回调
litellm.cache = litellm.Cache(type="redis") # 开缓存
litellm.max_budget = 10.0 # 设预算上限
resp = litellm.completion(model="gpt-4o", messages=[{"role": "user", "content": "hi"}])
print(resp._hidden_params["response_cost"]) # 这次花了多少钱,装饰器算好塞进来的
上面三行全局配置,作用点全部在装饰器里,completion 函数体本身对它们一无所知。
建立直觉的一句话: 装饰器是一次调用的"检票口 + 记账台"——进场前查票(预算、缓存、重试上限),出场后记账(成本、日志、回调),中间那段旅程(真正调模型)它不碰。
2. 顶层全景(一次调用穿过装饰器的七道关)
怎么读这张图: 从上往下是时间顺序;左侧竖线是装饰器管的事,中间那格是被装饰的原函数(真正调模型的地方)。任何一步抛异常都跳到最右侧的失败路径。
用户调用 litellm.completion(...)
│
▼
┌──────────────────────── ──────────────┐
│ ① 早退判定 │ _is_async_request / _is_streaming_request
│ 异步入口的透传 / 流式请求不做后处理 │
├──────────────────────────────────────┤
│ ② function_setup │ 建 litellm_call_id、装 callbacks、
│ 建"记录仪"Logging 对象 │ 造 Logging 对象
├──────────────────────────────────────┤
│ ③ 预算闸 + 每请求重试上限闸 │ 超了直接 BudgetExceededError
├──────────────────────────────────────┤
│ ④ 查缓存 │ 命中 → 直接 return,下面全部跳过
└──────────────────────────────────────┘
│ 未命中
▼
╔══════════════════════════════════════╗
║ 原函数:翻译 + HTTP + 流式拼装 ║ ← 见 02 / 03 章
╚══════════════════════════════════════╝
│ 成功 │ 失败
▼ ▼
┌──────────────────────────────────────┐ ┌───────────────────────┐
│ ⑤ post-call rules + JSON schema 校验 │ │ ⑧ 异常映射 │
├──────────────────────────────────────┤ │ exception_type() │
│ ⑥ 写缓存 │ │ → 统一异常类 │
├───────────────────── ─────────────────┤ ├───────────────────────┤
│ ⑦ 算成本 + 发成功回调(不阻塞主路径) │ │ ⑨ failure_handler │
└──────────────────────────────────────┘ └───────────────────────┘
│ │
▼ ▼
返回 response raise 统一异常
(_hidden_params 带 response_cost) (Router 据此决定回退/冷却)
各关卡的落点:
| 关卡 | 干什么 | 主要符号 | 文件 |
|---|---|---|---|
| ① 早退判定 | 判断是异步入口的透传调用、还是流式请求 | _is_async_request / _is_streaming_request | litellm/utils.py:1978 / :2020 |
| ② setup | 生成调用 ID、归位回调、造 Logging | function_setup | litellm/utils.py:786 |
| ③ 预算闸 | 累计花费超 litellm.max_budget 就拒 | BudgetExceededError | litellm/exceptions.py:960 |
| ④ 查缓存 | 按请求参数哈希查两级缓存 | LLMCachingHandler._sync_get_cache | litellm/caching/caching_handler.py:284 |
| ⑤ 后置校验 | 跑用户自定义 rule + JSON schema 校验 | post_call_processing | litellm/utils.py:1255 |
| ⑥ 写缓存 | 结果回填缓存 | sync_set_cache / async_set_cache | litellm/caching/caching_handler.py:1013 / :946 |
| ⑦ 算成本 + 回调 | 算钱写进 _hidden_params,异步派发回调 | Logging.success_handler | litellm/litellm_core_utils/litellm_logging.py:2185 |
| ⑧ 异常映射 | 各家错误 → 同一组异常类 | exception_type | litellm/litellm_core_utils/exception_mapping_utils.py:2164 |
| ⑨ 失败日志 | 发失败回调(同步、不放线程) | Logging.failure_handler | litellm/litellm_core_utils/litellm_logging.py:3037 |
3. 一个装饰器、两个 wrapper
3.1 它要解决的小问题
litellm.completion 是同步函数,litellm.acompletion 是协程函数。同一份横切逻辑要同时服务这两种,但 两者写法完全不同(一个用线程池派发回调,一个用 asyncio.create_task)。
3.2 思路
client 里定义两个闭包,最后按被装饰函数是不是协程挑一个返回:
# litellm/utils.py:1969-1975,真实源码(略去注释)
is_coroutine = get_coroutine_checker().is_async_callable(original_function)
if is_coroutine:
return wrapper_async
else:
return wrapper
即:wrapper(utils.py:1347) 给同步函数,wrapper_async(utils.py:1651) 给协程函数。
3.3 关键细节:为什么同步 wrapper 里还要判"这是不是异步请求"
因为 acompletion 内部其实又调了同步的 completion:
用户 → acompletion() [被 wrapper_async 包住]
│ completion_kwargs["acompletion"] = True (main.py:593)
▼
loop.run_in_executor(None, partial(completion, ...)) (main.py:631-635)
│
▼
completion() [被 wrapper 包住 —— 这里会二次横切!]
如果不管,一次异步调用会被记两次日志、查两次缓存、算两次钱。所以同步 wrapper 的第一行逻辑就是这个短路:
# litellm/utils.py:1348-1351,真实源码
# DO NOT MOVE THIS. It always needs to run first
call_type = original_function.__name__
if _is_async_request(kwargs):
...
_is_async_request(utils.py:1978) 只是检查 kwargs 里有没有 acompletion / aembedding / atranscription 等一堆 a* 标志位为 True。命中就走精简分支:只查一下 num_retries_per_request,然后直接调原函数并返回(utils.py:1353-1374),setup、预算、缓存、成本、回调全部跳过——因为外层的 wrapper_async 已经做过了。
3.4 关键细节:流式请求为什么要提前 return
流式响应返回的是 CustomStreamWrapper,此刻内容还没生成完,算不了成本、也没法跑内容校验。所以两个 wrapper 都在真实调用之后立刻判定流式并早退:
# litellm/utils.py:1496-1504(同步)/ :1763-1766(异步),真实源码结构
if _is_streaming_request(kwargs=kwargs, call_type=call_type):
if "complete_response" in kwargs and kwargs["complete_response"] is True:
chunks = []
...
_is_streaming_request(utils.py:2020) 的判定只有两条:kwargs 里 stream=True,或 call_type 属于 _STREAMING_CALL_TYPES(utils.py:2010,Google GenAI 的 generate_content_stream 系列)。
早退的代价: 流式调用的成功日志和成本核算不在装饰器里发生,而是延后到流结束时由 CustomStreamWrapper 触发(utils.py:1815-1816 的注释明确写了 "streaming requests return early (before this point) via CustomStreamWrapper")。除非用户传了 complete_response=True,那样装饰器会把所有 chunk 收齐再用 stream_chunk_builder 拼成完整响应返回。
4. function_setup:给这次调用装上记录仪
4.1 它要解决的小问题
后面所有横切动作(日志、成本、回调、缓存命中标记)都需要一个共享的"这次调用的上下文对象"。function_setup(utils.py:786) 就是造这个对象的地方,一次调用只跑一次。
4.2 它做的三件事
第一件:给这次调用发身份证。
# litellm/utils.py:1383-1384,真实源码
if "litellm_call_id" not in kwargs:
kwargs["litellm_call_id"] = str(uuid.uuid4())
注意这行在 wrapper 里、在 function_setup 之前执行。这个 litellm_call_id 会一路带到 _hidden_params(llm_response_utils/response_metadata.py:41)和所有回调的 payload 里,是串起分布式追踪的那根线。
第二件:把回调归位。 用户可以在三个地方声明回调,function_setup 负责把它们摊平并按"同步/异步"分家:
| 声明位置 | 变量 | 处理位置 |
|---|---|---|
| 全局 | litellm.success_callback / litellm.failure_callback / litellm.callbacks | utils.py:818-905 |
| 单次请求 | kwargs["callbacks"] | utils.py:818 get_dynamic_callbacks |
| 单次请求 | kwargs["success_callback"] / kwargs["failure_callback"] | utils.py:883-905 |
分家的判据是 coroutine_checker.is_async_callable(callback)(utils.py:855)——协程回调进 dynamic_async_success_callbacks,普通函数留在 dynamic_success_callbacks。这样后面派发时不用再判一次。
字符串形式的回调名("langfuse"、"lago"、"s3"…)由 _init_custom_logger_compatible_class(utils.py:824) 实例化成真正的 logger 类。全局回调列表的增删则统一走 litellm.logging_callback_manager(litellm/litellm_core_utils/logging_callback_manager.py:20 LoggingCallbackManager),它做的关键一件事是去重:_add_custom_logger_to_list(:289) 按 _get_custom_logger_key(:314) 判重,避免同一个 logger 被注册两次导致每条日志发两遍。
第三件:造 Logging 对象。
# litellm/utils.py:1083-1098,真实源码(节选)
logging_obj = get_litellm_logging_class()( # Victim for object pool
model=model,
messages=messages,
stream=stream,
litellm_call_id=kwargs["litellm_call_id"],
...
dynamic_success_callbacks=dynamic_success_callbacks,
...
)
这个 Logging 类定义在 litellm/litellm_core_utils/litellm_logging.py:375,构造函数在 :391。造好后立刻 update_environment_variables(:669) 把 litellm_params / metadata 灌进去,然后返回给 wrapper,wrapper 再写回 kwargs["litellm_logging_obj"](utils.py:1399),让下游的翻译层和 HTTP 层也能拿到同一个对象去打 pre_call / post_call。
一个容易忽略的入口: 若 kwargs 里已经带了 litellm_logging_obj(比如 Router 或 Proxy 在上层已经造好了),wrapper 会直接复用、跳过 function_setup(utils.py:1389-1390)。这就是"一次逻辑调用只有一条日志"的实现方式。
5. 三道闸:预算、每请求重试上限、缓存
三道闸都在真实调用之前,任何一道拦下就不会产生 API 花费。
5.1 预算闸
# litellm/utils.py:1409-1414,真实源码
if litellm.max_budget:
if litellm._current_cost > litellm.max_budget:
raise BudgetExceededError(
current_cost=litellm._current_cost,
max_budget=litellm.max_budget,
)
异步 wrapper 里有一份完全相同的(utils.py:1694-1699)。
_current_cost 是进程内的一个全局浮点数(litellm/__init__.py:404),累加发生在成功日志的准备函数里(litellm_logging.py:2051 litellm._current_cost += litellm.completion_cost(...))。
这里有个真实的局限(读代码才看得出):累加那段的判定条件是 litellm.max_budget 且 self.stream is False 且 isinstance(result, dict) 且 "content" in result(litellm_logging.py:2042-2047),而正常 completion 返回的是 ModelResponse 对象、不是 dict。也就是说 SDK 层这个 max_budget 闸在常规路径上几乎不会真正累加。要做真正的预算控制,应该用 Proxy 的 key/team 预算(见 06-proxy-gateway),而不是这个进程内变量。
BudgetExceededError(exceptions.py:960) 本身有个设计细节:它不继承 RateLimitError(怕破坏用户已有的 except BudgetExceededError:),但把 status_code = 429、category、rate_limit_type 三个字段照着限流错误的样子填上(exceptions.py:972-980),这样下游按"限流类错误"消费的回调仍能正确归类。
5.2 每请求重试上限闸
# litellm/utils.py:1417-1424,真实源码
if litellm.num_retries_per_request is not None:
previous_models = (kwargs.get("metadata") or {}).get("previous_models", None)
if previous_models is not None:
if litellm.num_retries_per_request <= len(previous_models):
raise Exception("Max retries per request hit!")
判据是 metadata["previous_models"] 的长度——这个列表由 Router 在每次重试/回退时追加。所以这道闸的本质是给 Router 的重试链条设一个全局硬顶,防止"回退到下一个模型 → 又失败 → 又回退"无限展开。注意它只在同步 wrapper 里出现(两处::1353 的异步透传分支和 :1417 的主路径),异步 wrapper 里没有对应的前置检查。
5.3 缓存闸
同步侧的入口:
# litellm/utils.py:1447-1455,真实源码
caching_handler_response: "CachingHandlerResponse" = _llm_caching_handler._sync_get_cache(
model=model or "",
original_function=original_function,
logging_obj=logging_obj,
start_time=start_time,
call_type=call_type,
kwargs=kwargs,
args=args,
)
命中就立刻 return caching_handler_response.cached_result(utils.py:1457-1459),后面的真实调用、成本、写缓存全部不发生。
进入这个分支的门槛写在 utils.py:1429-1445 的那个长条件里,读法是三段:
- 开关:
kwargs["caching"] is True,或者用户设了litellm.cache且没显式传caching=False; - 旁路:
kwargs["cache"]["no-cache"]不为True; - 排除: call_type 不是
aembedding/acompletion/atranscription等一串a*异步标志——这些交给异步 wrapper 处理。
异步侧简单得多:wrapper_async 无条件调 _async_get_cache(utils.py:1706,实现在 caching_handler.py:145),由 handler 内部决定要不要真查。它还多一条同步侧没有的返回路径:embedding 的部分命中——一批 input 里命中一部分,没命中的那部分照常发请求,最后由 _combine_cached_embedding_response_with_api_result(caching_handler.py:588) 把两半拼起来(utils.py:1854-1857)。
6. 缓存:两级存储 + 一张后端清单
6.1 键是怎么算出来的
Cache.get_cache_key(litellm/caching/caching.py:320) 的做法是:把 kwargs 里属于"LLM API 参数"的字段(ModelParamHelper._get_all_llm_api_params())按 "参数名: 值" 拼成一个长字符串,再哈希(_get_hashed_cache_key)。
含义很直接:模型、messages、temperature 任何一个不同,就是不同的 key。这是精确匹配,不是"意思差不多"。想要"意思差不多也算命中",得上 §6.3 的语义缓存。
6.2 两级:内存 + Redis
DualCache(litellm/caching/dual_cache.py:51) 是"本地内存 + Redis"的组合。读路径:
get_cache(key)
│
▼
┌───────────────┐ 命中 ┌──────────┐
│ InMemoryCache │ ───────► │ 返回值 │
└───────┬───────┘ └──────────┘
│ 未命中
▼
┌───────────────┐ 命中 ┌──────────────────────┐
│ RedisCache │ ── ─────► │ 回填 InMemoryCache │──► 返回值
└───────┬───────┘ └──────────────────────┘
│ 未命中
▼
None
对应源码 dual_cache.py:153-182(同步 get_cache)与 :217-247(async_get_cache),两者结构完全一致,回填那行是 self.in_memory_cache.set_cache(key, redis_result, **self._backfill_kwargs(kwargs))(:175)。
内存那一级是有界的:InMemoryCache(in_memory_cache.py:28) 默认 max_size_in_memory=200、default_ttl=600 秒(:31-40),超出时 evict_cache(:101) 按最早过期时间淘汰。
6.3 后端清单
| 后端 | 类 | 文件 | 匹配方式 | 特点 |
|---|---|---|---|---|
| 内存 | InMemoryCache | caching/in_memory_cache.py:28 | 精确 | 默认 200 条 / TTL 600s,进程内 |
| Redis | RedisCache | caching/redis_cache.py:267 | 精确 | 跨进程共享,Proxy 多副本的默认选择 |
| Redis 集群 | RedisClusterCache | caching/redis_cluster_cache.py | 精确 | 需 redis_startup_nodes |
| 磁盘 | DiskCache | caching/disk_cache.py:14 | 精确 | 基于 diskcache 库,默认目录 .litellm_cache |
| S3 | S3Cache | caching/s3_cache.py:23 | 精确 | boto3 直写对象;不支持批量写,embedding 走单条路径(caching_handler.py:995) |
| GCS / Azure Blob | GCSCache / AzureBlobCache | caching/gcs_cache.py、caching/azure_blob_cache.py | 精确 | 同上,云对象存储 |
| Redis 语义 | RedisSemanticCache | caching/redis_semantic_cache.py:38 | 向量相似 | 必须传 similarity_threshold,内部转成余弦距离阈值 1 - threshold(:99) |
| Valkey 语义 | ValkeySemanticCache | caching/valkey_semantic_cache.py | 向量相似 | Redis 语义缓存的 Valkey 版 |
| Qdrant 语义 | QdrantSemanticCache | caching/qdrant_semantic_cache.py:37 | 向量相似 | 支持 binary / scalar / product 三种量化(:100-118),默认二值量化 |
选型由 Cache.__init__(caching/caching.py:56) 的 type 参数决定;哪些 call_type 允许缓存由 supported_call_types(caching.py:69-82) 控制,默认包含 completion / embedding / transcription / rerank / responses 五类的同步与异步版本。
6.4 语义缓存的读法
精确缓存问"这条请求我见过吗";语义缓存问"我见过意思差不多的请求吗"。实现上是把 prompt 过一遍 embedding 模型(redis_semantic_cache.py:330 _get_embedding),在向量库里做近邻检索,把返回的余弦距离换算回相似度再和阈值比:
# litellm/caching/redis_semantic_cache.py:471-476,真实源码
vector_distance = float(cache_hit["vector_distance"])
# Convert vector distance back to similarity score
# For cosine distance: 0 = most similar, 2 = least similar
similarity = 1 - vector_distance
代价也很直白:每次查缓存都要多打一次 embedding API。语义缓存的 key 里还会剔掉一批参数(_SEMANTIC_CACHE_SCOPE_EXCLUDED_PARAMS,caching.py:288),否则连 temperature 改一下都会分桶,语义匹配就失去意义了。
6.5 写回
- 同步:
sync_set_cache(caching_handler.py:1013) 直接litellm.cache.add_cache(...),阻塞。 - 异步:
async_set_cache(caching_handler.py:946) 一律asyncio.create_task(...)(:997、:1003、:1011)——发射后不管,写缓存的耗时不进用户的延迟。
两者都先过 _should_store_result_in_cache(caching_handler.py:1038) 这道判断。
7. 成本:从 usage 到 _hidden_params["response_cost"]
7.1 它要解决的小问题
用户想知道"这次调用花了多少钱"。但每家 provider 的计价单位不一样:OpenAI 按 token,Vertex 部分模型按字符,Replicate 按秒,Rerank 按 query 数,转录按音频时长。横切层要把这些统一成一个美元数字。
7.2 一条计算链
响应对象 result
│
▼
Logging._response_cost_calculator litellm_logging.py:1489
│ ├─ cache_hit → 直接返回 0.0
│ └─ _hidden_params 里已有 response_cost → 直接复用
▼
response_cost_calculator cost_calculator.py:1715
│
▼
completion_cost cost_calculator.py:1112
│ ├─ _select_model_name_for_cost_calc :729 决定"按哪个模型名查价"
│ └─ _get_usage_object :869 把各家 usage 统一成 Usage
▼
cost_per_token cost_calculator.py:300
│ 查 litellm.model_cost[模型名] 里的单价
▼
(prompt_cost, completion_cost) → 相加 → float 美元
_response_cost_calculator 的前两个短路很值得注意:
# litellm/litellm_core_utils/litellm_logging.py:1524-1536,真实源码(节选)
if cache_hit is True:
return 0.0
...
if isinstance(result, BaseModel) and hasattr(result, "_hidden_params"):
hidden_params = getattr(result, "_hidden_params", {})
if ("response_cost" in hidden_params and hidden_params["response_cost"] is not None):
return hidden_params["response_cost"] # use cost if already calculated
第一条保证缓存命中不重复计费;第二条保证一次响应只算一次钱(Router 场景下同一个响应可能被多层看到)。
7.3 模型名的选择比想象中麻烦
_select_model_name_for_cost_calc(cost_calculator.py:729) 的 docstring 直接写了优先级:
- 用户开了自定义定价 → 用传进来的 model 名
- 设了
base_model(Azure 常见:部署名 ≠ 模型名)→ 用 base_model - 响应对象里的
model字段 - 兜底用调用时传的 model
第 2 条是 Azure 用户最容易踩的坑:Azure 的 deployment 名可以随便起(比如 my-gpt4-prod),价格表里根本没有这个名字,必须靠 base_model 告诉 LiteLLM"它其实是 gpt-4"。
7.4 价格数据从哪来
| 来源 | 位置 | 何时用 |
|---|---|---|
| 远程价格表 | model_cost_map_url(litellm/__init__.py:412,默认指向 GitHub raw 上的 model_prices_and_context_window.json) | 默认,import litellm 时拉一次 |
| 仓库内价格表 | model_prices_and_context_window.json(仓库根) | 远程表的源文件 |
| 打包备份 | litellm/model_prices_and_context_window_backup.json | 拉取失败、或校验不过时兜底 |
| 用户自定义 | litellm.register_model(...)(utils.py:2831) | 私有模型 / 议价后的单价 |
加载逻辑在 get_model_cost_map(litellm/litellm_core_utils/get_model_cost_map.py:426)。它有两个防呆设计:
- 环境变量
LITELLM_LOCAL_MODEL_COST_MAP=true→ 完全离线,只读本地备份(:277-282)。适合不想让库在 import 时联网的生产环境。 - 缩水校验:远程表拉回来后要跟本地备份的条目数比(
GetModelCostMap.validate_model_cost_map,常量MODEL_COST_MAP_MAX_SHRINK_RATIO/MODEL_COST_MAP_MIN_MODEL_COUNT)。远程文件如果被截断成半个,会被判定为不可信并回落本地备份——避免"上游一次坏发布让全世界的成本统计归零"。
最终产物 是 litellm.model_cost 这个全局 dict(litellm/__init__.py:530)。
7.5 算完的钱去哪了
答案是 _hidden_params。链路是:wrapper 在返回前调 update_response_metadata(utils.py:1563-1571),它构造 ResponseMetadata 并调 set_hidden_params:
# litellm/litellm_core_utils/llm_response_utils/response_metadata.py:40-48,真实源码(节选)
new_params = {
"litellm_call_id": getattr(logging_obj, "litellm_call_id", None),
"api_base": get_api_base(model=model or "", optional_params=kwargs),
"model_id": model_id,
"response_cost": logging_obj._response_cost_calculator(
result=self.result, litellm_model_name=model, router_model_id=model_id
),
...
}
所以用户拿到的 resp._hidden_params["response_cost"] 就是这里写进去的。同一个数字还会走另一条路进日志:_process_hidden_params_and_response_cost(litellm_logging.py:1903) 把它写进 model_call_details["response_cost"],再打包进 standard_logging_object— —这就是所有回调(Langfuse、Prometheus、Proxy 的账单表)看到的那份成本。
细项拆分(输入/输出/缓存读/缓存写/推理 token 分别多少钱)由 Logging.set_cost_breakdown(litellm_logging.py:1416) 存成 CostBreakdown。
8. 日志与回调:四个钩子 + 一个后台队列
8.1 四个钩子分别在什么时候打
| 钩子 | 时机 | 谁来调 | 位置 |
|---|---|---|---|
pre_call | 请求即将发出(已知 api_base / headers / body) | 翻译层与 HTTP 层 | litellm_logging.py:1064 |
post_call | 原始响应刚回来(还没解析成对象) | HTTP 层 | litellm_logging.py:1269 |
success_handler | 成功且拿到结构化响应 | 装饰器(线程池) | litellm_logging.py:2185 |
async_success_handler | 同上,异步路径 | 后台 logging worker | litellm_logging.py:2619 |
failure_handler | 抛异常 | 装饰器(当前线程) | litellm_logging.py:3037 |
async_failure_handler | 同上,异步路径 | 装饰器(await) | litellm_logging.py:3231 |
pre_call 里做的一件重要的事是把请求还原成 curl 命令存进 metadata(litellm_logging.py:1090-1097,_get_request_curl_command),Langfuse 之类的平台可以直接展示"这条请求原样长什么样"——但 headers 会先过 _get_masked_headers(:1261)。
8.2 成功回调怎么做到不拖慢主路径
同步路径用线程池,并且显式复制 contextvars,否则 OpenTelemetry 的 span 上下文会在跨线程时丢失:
# litellm/utils.py:1552-1560,真实源码
ctx = contextvars.copy_context()
executor = getattr(sys.modules[__name__], "executor")
executor.submit(
ctx.run,
logging_obj.success_handler,
result,
start_time,
end_time,
)
异步路径进后台队列。_client_async_logging_helper(utils.py:1136) 把 async_success_handler 这个协程扔给全局 worker:
# litellm/utils.py:1152-1156,真实源码
from litellm.litellm_core_utils.logging_worker import GLOBAL_LOGGING_WORKER
GLOBAL_LOGGING_WORKER.ensure_initialized_and_enqueue(
async_coroutine=logging_obj.async_success_handler(result=result, start_time=start_time, end_time=end_time)
)
LoggingWorker(litellm/litellm_core_utils/logging_worker.py:35) 是个有界队列 + 信号量并发控制的消费者:
主路径 ──enqueue()──► asyncio.Queue (上限 50_000)
(从不阻塞) │
▼
_worker_loop ── 先 acquire 信号量(并发 100) 再取任务
│
▼
_process_log_task (单任务超时 20s)
常量在 litellm/constants.py:453-455:并发 100、队列 50000、单协程 20 秒超时。队列满了不阻塞主路径,而是走 _handle_queue_full(:187) 主动丢弃/退避重试——类注释里把这个取舍写得很清楚:这是 best-effort,换来 "+200 RPS"(logging_worker.py:37-42)。
进程退出时 atexit 注册的 _flush_on_exit(:448) 会尽量把队列排空,减少"最后几条日志丢了"的情况。
8.3 失败回调为什么偏偏不放线程
两个 wrapper 的异常分支里,failure_handler 都带着同一句大写注释:
# litellm/utils.py:1643-1647,真实源码
logging_obj.failure_handler(
e, traceback_exception, start_time, end_time
) # DO NOT MAKE THREADED - router retry fallback relies on this!
原因是 Router 的重试和回退依赖失败已经被完整记录(冷却计数要用到它)。如果丢进线程池异步执行,Router 可能在日志落下来之前就已经切到下一个 deployment 了,冷却统计会失真。
8.4 敏感信息不落日志
两层机制,作用对象不同:
| 机制 | 遮什么 | 入口 | 开关 |
|---|---|---|---|
| 消息脱敏 | messages / prompt / 响应内容 | redact_message_input_output_from_logging(redact_messages.py:347) | litellm.turn_off_message_logging,或请求头 |
| 字段掩码 | 参数里像密钥的字段(值变成 sk-1***********ef12) | SensitiveDataMasker(sensitive_data_masker.py:9) | 默认按 key 名匹配 |
消息脱敏的开关判定有明确优先级,写在 should_redact_message_logging(redact_messages.py:295) 的 docstring 里:动态参数 > 请求头(litellm-disable-message-redaction / x-litellm-enable-message-redaction)> 全局 litellm.turn_off_message_logging。
字段掩码靠 key 名里的敏感词表(password / secret / key / token / credentials / certificate …,sensitive_data_masker.py:19-34)。有个很小但很实用的补丁在 non_sensitive_overrides(:39):input_cost_per_token 里含 token 却是价格字段,靠 cost 这个反向词把它排除,否则价格会被打成星号。
9. 异常统一:Router 能回退的前提
9.1 它要解决的小问题
同样是"上下文超长",各家给的信号完全不同:
- OpenAI:
"This model's maximum context length is ..." - Anthropic:
"prompt is too long"/"prompt: length" - llama.cpp:
"exceeds the available context size" - Gemini:
"exceeds the maximum number of tokens allowed" - Cerebras:
"Current length is X while limit is Y"
这些文案被集中列在一个共用判定里(ExceptionCheckers.is_error_str_context_window_exceeded,exception_mapping_utils.py:74-103,目前 9 条已知子串 + 1 条 Cerebras 模式)。
上层(Router、Proxy、用户代码)不可能为每家写一套 if。所以横切层要把这些收敛成同一组异常类。
9.2 异常体系:全部继承 openai.*
litellm/exceptions.py 的做法很干脆:所有异常都继承 openai SDK 的对应异常。这样已经在用 openai SDK 的代码,except openai.RateLimitError: 原样能接住 LiteLLM 抛的错。
| LiteLLM 异常 | 继承自 | 定义 | 典型触发 |
|---|---|---|---|
AuthenticationError | openai.AuthenticationError | exceptions.py:129 | 401 ,key 错 |
NotFoundError | openai.NotFoundError | :173 | 404,模型不存在 |
BadRequestError | openai.BadRequestError | :216 | 400 |
Timeout | openai.APITimeoutError | :330 | 408 / 读超时 |
RateLimitError | openai.RateLimitError | :413 | 429 |
InternalServerError | openai.InternalServerError | :729 | 500 |
ServiceUnavailableError | openai.APIStatusError | :633 | 503 |
BadGatewayError | openai.APIStatusError | :681 | 502 |
ContextWindowExceededError | BadRequestError | :504 | 上下文超长(LiteLLM 特有细分) |
ContentPolicyViolationError | BadRequestError | :588 | 内容被安全策略拦(LiteLLM 特有细分) |
BudgetExceededError | Exception | :960 | 预算闸拦下(不走 openai 体系) |
MidStreamFallbackError | ServiceUnavailableError | :1073 | 流已经开始吐字后才失败,需要中途换模型 |
加粗那四个是 LiteLLM 自己加的。前两个之所以要从 BadRequestError 里再细分出来,就是为了 §9.4 的 Router 回退。
MidStreamFallbackError 值得单说:它带一个 generated_content 字段(exceptions.py:1084)保存"