跳到主内容
快讯直播
AI智模界
教程

Claude Haiku 低成本流水线:分层、并发与成本实测

把高频、低难度的模型调用从"大模型全量跑"改成"按任务分层 + 缓存前缀 + 并发受控"的流水线,是压低推理账单比较直接的一条路。这篇教程带你从零搭出这样一条流水线:分类、抽取、改写三类任务走 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 结构化抽取从工单抽产品、订单号、金额、城市输入中等,输出结构化 JSONHaiku 档位 + 结果校验
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. 探索批处理与结构化输出。 离线任务用批接口,需要严格字段的任务用工具调用或结构化输出约束,能进一步降低解析失败率。

把这几步跑完,你就有一条能说清"每一千次调用花多少钱、错在哪里、什么时候该升级"的流水线了,而不只是把模型换个名字。

AI 生成本文由 AI 基于公开信息自动生成,仅供参考。