把高频、低难度的模型调用从"大模型全量跑"改成"按任务分层 + 缓存前缀 + 并发受控"的流水线,是压低推理账单比较直接的一条路。这篇教程带你从零搭出这样一条流水线:分类、抽取、改写三类任务走 Haiku 档位的模型,长链推理保留给上一级模型;同时给出可复用的并发控制、提示缓存写法和每千次调用的成本对照表。
这篇能做出什么
完成后你会得到三个可运行的文件:
1. tasks.py:把分类、抽取、改写封装成统一签名的异步函数,每个函数自带提示词、参数和 JSON 输出约束。
2. pipeline.py:带信号量限流、指数退避重试、失败升级(先跑小模型,不合格再升到大模型)的批处理入口。
3. cost.py:从每次响应的 usage 里采集输入、输出、缓存写入、缓存读取四类 token,汇总算出每个任务的单次成本和每千次成本。
最终你会拿到一张成本对照表,能回答"这 10 万条工单切到 Haiku 之后大概省多少"这个问题,而不是靠感觉。
前置条件清单
- 一个可用的 Anthropic API 密钥,放在环境变量里,不要写进代码。
- Python 3.9 以上,能创建虚拟环境。
- 一份 200~500 条的真实样本数据(JSONL 或者 CSV 都行),用来测 token 分布和准确率。
- 一个 100 条左右的人工标注小集合,用来判断"切到 Haiku 之后质量掉没掉"。
- 官方定价页和模型列表页各开一个标签页。本文不写具体单价和模型 ID,这些变化快,以官方文档当前版本为准。
第一步:装环境,跑通最小调用
```bash
python -m venv .venv
source .venv/bin/activate # Windows 用 .venv\Scripts\activate
pip install anthropic python-dotenv
```
密钥放到环境变量。模型 ID 不要硬编码在代码里,写进环境变量,将来换模型只改一处。
```bash
export ANTHROPIC_API_KEY="你的密钥"
export HAIKU_MODEL="从官方模型列表页复制的模型 ID"
export MAX_CONCURRENCY="8"
```
最小调用长这样:
```python
minimal_call.py
import os
from anthropic import Anthropic
client = Anthropic(api_key=os.environ["ANTHROPIC_API_KEY"])
resp = client.messages.create(
model=os.environ["HAIKU_MODEL"],
max_tokens=256,
temperature=0,
system="你是一个严谨的文本分类器,只输出标签,不要解释。",
messages=[{"role": "user", "content": "文本:这款耳机续航很长,但戴久了夹头。"}],
)
print(resp.content[0].text)
print(resp.usage)
```
resp.usage 是后面所有成本统计的数据源,先把它打印出来看一眼结构,心里有数。
第二步:做任务分层,决定谁走 Haiku
不要凭感觉分层,用两个维度判断:错误代价和调用量。
| 层级 | 典型任务 | 输入/输出特征 | 建议走向 |
|---|---|---|---|
| L1 分类打标 | 意图识别、情感判断、工单路由 | 输入短,输出极短,答案空间封闭 | Haiku 档位 |
| L2 结构化抽取 | 从工单抽产品、订单号、金额、城市 | 输入中等,输出结构化 JSON | Haiku 档位 + 结果校验 |
| L3 改写生成 | 摘要、润色、多语言改写、扩写 | 输出长,对文风敏感 | Haiku 先跑,抽样评估 |
| L4 多跳推理 | 长文档分析、复杂决策、代码修复 | 上下文长,需要多步推理 | 保留上一级模型 |
一条实用规则:能用关键词、正则、字典查出来的,就不要调模型。 例如"包含订单号格式的字符串"完全可以用正则抽,把这类判断前置到代码里,能省掉一整层调用。
第三步:三类任务的统一封装
关键约定:每个任务都返回 (结果, usage字典),这样成本统计不用回头改代码。
```python
tasks.py
import json
import os
from anthropic import AsyncAnthropic
client = AsyncAnthropic(api_key=os.environ["ANTHROPIC_API_KEY"])
MODEL = os.environ["HAIKU_MODEL"]
JSON_TAIL = "\n只输出合法 JSON,不要 Markdown 代码围栏,不要任何解释文字。"
CLASSIFY_SYSTEM = """你是客服工单分类器。
可选标签:退款、物流、质量、咨询、投诉、其他。
判断规则:
- 同时提到退款和质量问题时,优先归到"质量"。
- 信息不足无法判断时,输出"其他",不要猜。
输出格式:{"label": "<标签>"}""" + JSON_TAIL
EXTRACT_SYSTEM = """你是信息抽取器。字段定义:
- product: 产品名称,未提及填 null
- order_id: 订单号,未提及填 null
- city: 收货城市,未提及填 null
- amount: 金额数字(元),未提及填 null
- urgency: 高 / 中 / 低,出现"尽快""着急"等表述记为高
输出格式:{"product": ..., "order_id": ..., "city": ..., "amount": ..., "urgency": ...}""" + JSON_TAIL
```
每个任务用一个 _call 收敛样板代码:
```python
def _usage(resp):
u = resp.usage
return {
"in": getattr(u, "input_tokens", 0) or 0,
"out": getattr(u, "output_tokens", 0) or 0,
"cache_write": getattr(u, "cache_creation_input_tokens", 0) or 0,
"cache_read": getattr(u, "cache_read_input_tokens", 0) or 0,
}
async def _call(system_text, user_text, max_tokens, cache=False):
system = [{"type": "text", "text": system_text}]
if cache:
system[0]["cache_control"] = {"type": "ephemeral"}
resp = await client.messages.create(
model=MODEL,
max_tokens=max_tokens,
temperature=0,
system=system,
messages=[{"role": "user", "content": user_text}],
)
return resp.content[0].text, _usage(resp)
async def classify(text):
return await _call(CLASSIFY_SYSTEM, f"工单内容:\n{text}", max_tokens=64, cache=True)
async def extract(text):
return await _call(EXTRACT_SYSTEM, f"工单内容:\n{text}", max_tokens=256, cache=True)
async def rewrite(text, style="简洁"):
system = f"你是文本改写助手。要求:{style},保留全部事实信息,不新增内容。只输出改写后的正文。"
return await _call(system, text, max_tokens=1024)
```
分类和抽取的输出很短,max_tokens 给到 64 / 256 就够,写大了没有意义,写小了会被截断导致 JSON 解析失败。
第四步:提示缓存怎么用才不亏
提示缓存省钱的原理很朴素:同一段前缀被反复使用时,后续命中按更低的读取价计费。 但首次写入是按写入价计费的,通常比普通输入更贵。所以缓存不是"开了就省",而是"同一前缀重复用够多次才省"。
三条实践规则:
1. 把稳定内容全部放前面。 分类规则、标签定义、字段说明、少样本示例,这些每批任务都不变,适合缓存。
2. 动态内容一定放最后。 待处理文本、时间戳、随机 ID、用户 ID 都不要出现在缓存块里。一个变化的字符就会让前缀全废。
3. 缓存块要够长。 最小可缓存长度、缓存存活时间、写入与读取的具体倍率,以官方文档当前版本为准,别按记忆里的数字推算。
代码里体现为:cache_control 只加在 system 的第一块(也就是那段长规则)上,待处理文本在 messages 里。如果规则文本比较长,把它单独抽成一个文件加载,方便统一维护。
验证缓存是否真的命中,看两个字段:
```python
if usage["cache_read"] > 0:
print("命中缓存:", usage["cache_read"], "tokens")
else:
print("未命中,本次为写入或普通输入")
```
如果长期 cache_read 为 0,回到上面三条规则逐条排查。
第五步:并发控制,别一上来就 64
并发不是越大越快。请求在服务端排队,超过配额就是 429,重试反而拖慢整体。稳妥做法是信号量限流 + 指数退避 + 抖动。
```python
pipeline.py
import asyncio
import os
import random
from anthropic import APIConnectionError, APIStatusError, RateLimitError
import tasks
CONCURRENCY = int(os.environ.get("MAX_CONCURRENCY", "8"))
async def call_with_retry(fn, *args, retries=4, **kwargs):
delay = 1.0
for attempt in range(retries + 1):
try:
return await fn(*args, **kwargs)
except RateLimitError:
if attempt == retries:
raise
await asyncio.sleep(delay + random.random())
delay *= 2
except APIConnectionError:
if attempt == retries:
raise
await asyncio.sleep(delay + random.random())
delay *= 2
except APIStatusError as e:
5xx 可以重试,4xx 多半是参数问题,重试没意义
if e.status_code < 500 or attempt == retries:
raise
await asyncio.sleep(delay + random.random())
delay *= 2
async def run_batch(items, handler, concurrency=CONCURRENCY):
sem = asyncio.Semaphore(concurrency)
async def one(item):
async with sem:
return await call_with_retry(handler, item)
return await asyncio.gather(*(one(i) for i in items), return_exceptions=True)
if __name__ == "__main__":
import json
texts = [json.loads(l)["text"] for l in open("sample.jsonl", encoding="utf-8")]
results = asyncio.run(run_batch(texts[:100], tasks.classify, concurrency=8))
ok = sum(1 for r in results if not isinstance(r, Exception))
print(f"成功 {ok} / {len(results)}")
```
调并发的方法:从 4 开始,翻倍到 8、16,观察两件事——成功率有没有下降、总耗时是不是还在缩短。任一指标变差就退回去。如果业务有明确的服务端速率限制,把上限设在限制的七八成,留出余量。
离线的批量任务还有一条路:用官方的批处理接口(Batches)提交整批作业,通常有价格优惠,代价是延迟高、不适合实时交互。具体折扣和限制以官方文档当前版本为准。
第六步:每千次调用成本实测
先把单价填进配置。不要从这篇教程里抄单价,去官方定价页复制当前数值,单位统一成"元 / 百万 token"。
```python
cost.py
PRICES = {
"in": 0.0, # 普通输入单价
"out": 0.0, # 输出单价
"cache_write": 0.0, # 缓存写入单价
"cache_read": 0.0, # 缓存读取单价
}
def call_cost(u):
return (
u["in"] * PRICES["in"]
+ u["out"] * PRICES["out"]
+ u["cache_write"] * PRICES["cache_write"]
+ u["cache_read"] * PRICES["cache_read"]
) / 1_000_000
def per_1000(total_cost, n):
return total_cost / n * 1000 if n else 0.0
```
然后跑一次真实小批量,把 usage 按任务聚合:
```python
from collections import defaultdict
agg = defaultdict(lambda: {"n": 0, "in": 0, "out": 0, "cw": 0, "cr": 0})
def record(bucket, u):
a = agg[bucket]
a["n"] += 1
a["in"] += u["in"]
a["out"] += u["out"]
a["cw"] += u["cache_write"]
a["cr"] += u["cache_read"]
def report():
for name, a in agg.items():
avg = {k: a[k] / a["n"] for k in ("in", "out", "cw", "cr")}
c = call_cost({"in": avg["in"], "out": avg["out"],
"cache_write": avg["cw"], "cache_read": avg["cr"]})
print(f"{name}: 平均输入 {avg['in']:.0f} / 输出 {avg['out']:.0f} "
f"/ 缓存读 {avg['cr']:.0f} tokens,单次 {c:.6f} 元,"
f"每千次 {per_1000(c, 1):.4f} 元")
```
对照表按这个模板填:
| 任务 | 平均输入 token | 平均输出 token | 缓存读 token | 单次成本 | 每千次成本 | 上一级模型每千次 | 倍率 |
|---|---|---|---|---|---|---|---|
| L1 分类 | 实测填入 | 实测填入 | 实测填入 | 公式算出 | ×1000 | 按同口径算 | 相除 |
| L2 抽取 | |||||||
| L3 改写 |
算倍率时注意一件事:输入和输出的倍率通常不一样,上一级模型可能输入贵得少、输出贵得多。所以改写这类"输出很长"的任务,切到小模型的收益往往比分类更明显。
举个算术例子说明口径(数字是假设,方便看公式):某次分类平均输入 900 token,其中 800 走缓存读取、100 走普通输入,输出 20 token。单次成本 =(100 × 输入单价 + 800 × 缓存读单价 + 20 × 输出单价)÷ 1,000,000,每千次成本再乘 1000。把单价替换成官方页面上的当前数值即可。
还要单独记录缓存未命中的比例。第一批请求必然是写入,稳态才便宜;如果每批任务的前缀都在变,缓存基本白开。
第七步:校验与失败升级
小模型跑得快,但会有格式错误和判断错误。两类问题分开处理:
```python
import json
LABELS = {"退款", "物流", "质量", "咨询", "投诉", "其他"}
def parse_label(text):
try:
obj = json.loads(text)
except json.JSONDecodeError:
return None
label = obj.get("label")
return label if label in LABELS else None
```
- 格式错误(JSON 解析失败、字段缺失):直接重试一次,提示里补一句"上次输出不是合法 JSON,请只输出 JSON"。还失败就升级到上一级模型。
- 语义错误(标签不在白名单、抽取的金额明显不合理):升级到上一级模型重跑,并记一条样本到复盘集。
升级率是个关键指标。如果升级率低于 5%,说明这层任务下沉是划算的;如果高于 20%,要么提示词还要打磨,要么这个任务本来就不该走小模型。
常见坑与排错
缓存永远不命中。 优先查三处:system 里是不是混进了时间戳或请求 ID;前缀长度是不是低于最小可缓存长度;两批请求之间的间隔是不是超过了缓存存活时间。用 cache_read 是否为 0 直接判定。
一开并发就 429。 别把重试次数调高硬扛,先降并发。重试要带抖动,否则所有请求会在同一时刻再次撞上去。同时区分错误类型:4xx 里的参数错误重试一百次也不会成功。
成本算重复或者算漏。 开了缓存之后,input_tokens 和 cache_read_input_tokens 是分开的字段,直接相加会重复计价。按四个字段分别乘各自单价再求和。另外重试成功的调用是计费的,重试逻辑要能感知到,别只统计首次调用。
模型 ID 硬编码。 换来换去很痛苦。统一从环境变量读,代码里只出现变量名。
超时设置不够。 改写类任务输出长,默认超时可能不够用,客户端初始化时可以显式设大一点。长输出的交互场景用流式,离线批处理走批接口。
提示词里没说"不要解释"。 小模型有时会贴心地加一句"好的,以下是结果"。JSON 解析就会炸。system 结尾固定加一句输出约束,能省掉大量解析失败。
没有评测集就下结论。 只看成本不看质量,等于把问题推迟到线上。哪怕只有 100 条人工标注,也先跑一遍对比。
下一步建议
1. 把评测做成常态。 每次改提示词或换模型,都跑一遍那 100 条标注集,记录准确率和升级率两条曲线。
2. 把成本指标接进看板。 按任务、按天聚合每千次成本、缓存命中率、升级率,异常时能第一时间看到。
3. 给并发加上自适应。 遇到 429 就自动下调并发,稳定一段时间再逐步加上去,比固定值省心。
4. 继续往下分层。 分类里还可以再拆:命中关键词的直接用规则,长尾再调模型。省下来的往往不只是钱,还有延迟。
5. 探索批处理与结构化输出。 离线任务用批接口,需要严格字段的任务用工具调用或结构化输出约束,能进一步降低解析失败率。
把这几步跑完,你就有一条能说清"每一千次调用花多少钱、错在哪里、什么时候该升级"的流水线了,而不只是把模型换个名字。
