☰
LLM 批量处理实战:5000 条文本的并发、重试与成本核算(附代码)
2026/10/7 14:14:49 网站建设 项目流程

LLM 批量处理实战:5000 条文本的并发、重试与成本核算(附代码)

LLM 系列又一篇。结构化输出篇讲了怎么让模型返回数据——这篇讲量的另一面:5000 条数据怎么批处理。逐条调用慢、触发限流、System Prompt 重复发,这篇把完整链路走通:数据准备 → 小批量验证 → 并发放大 → 核验成本。

一、先看逐条调用会撞上的三堵墙

几千条文本要分类(评论情感、工单打标、文档摘要),最直觉的写法是循环逐条调 API。跑 200 条之后三堵墙就撞上了:

  1. 慢:每次都要等模型返回,串行 5000 条按分钟级延迟算要几个小时
  2. 限流:连续请求触发平台频控,被拒后等待重试,更慢
  3. token 浪费:每条都重发一遍 System Prompt,5000 条就是 5000 次重复(成本篇的账:Agent 输入是大头,批量场景同理)

解法不是硬扛,是换链路:JSONL 数据准备 → 先 20 条验证 → 并发放大 → 核验成本。四步全是通用工程,附完整代码。

二、数据准备:JSONL + custom_id

批量任务的数据格式用JSONL(每行一个 JSON),关键是每条带唯一custom_id——结果回填时靠它对上号(Function Calling 篇的 tool_use ID 同一个思路:靠 ID 对齐请求和结果):

import json def to_jsonl(items: list, path: str, system: str, model: str): """items: [{"id": ..., "text": ...}],转成批量输入 JSONL""" with open(path, "w", encoding="utf-8") as f: for it in items: f.write(json.dumps({ "custom_id": it["id"], "body": { "model": model, "messages": [ {"role": "system", "content": system}, {"role": "user", "content": it["text"]}, ], "temperature": 0, # 分类任务低温(评测篇老话) "response_format": {"type": "json_object"}, # 结构化输出(第十七篇层 2) }, }, ensure_ascii=False) + "\n") SYSTEM = ('对评论做情感分类,只返回 JSON:' '{"sentiment": "正面|负面|中性", "confidence": 0-1}') to_jsonl(reviews, "batch_input.jsonl", SYSTEM, MODEL)

三、核心方法论:先 20 条验证,再放大

5000 条不要一上来就全跑。先抽 20 条跑通整条链路,确认三件事:

  1. 格式对:JSONL 被正确识别,没有解析错误
  2. 输出稳:每条都返回合法结构(结构化输出的校验兜底全过)
  3. 能对上号:custom_id和结果正确对应

这 20 条如果有标注,顺便算一遍准确率——模型不行的话,改 prompt 或换模型都在 20 条的规模上做,改完再放大。这一步是评测篇「先建评测集再调优」在批量场景的化身:验证是防御性成本,比 5000 条跑完发现全错便宜 250 倍。

四、并发放大:信号量 + 指数退避

没有平台批量接口时,自己写并发:信号量限并发 + 指数退避重试,40 行:

import asyncio, json from openai import AsyncOpenAI client = AsyncOpenAI(base_url=BASE_URL, api_key=KEY) def load_jsonl(path: str) -> list: return [json.loads(l) for l in open(path, encoding="utf-8") if l.strip()] async def process_one(task: dict, sem: asyncio.Semaphore, max_retries: int = 3) -> dict: async with sem: # 限并发,防触发频控 for attempt in range(max_retries): try: resp = await client.chat.completions.create(**task["body"]) content = resp.choices[0].message.content data = json.loads(content) # 结构化输出:解析校验 return {"custom_id": task["custom_id"], "ok": True, **data} except Exception as e: if attempt == max_retries - 1: return {"custom_id": task["custom_id"], "ok": False, "error": str(e)} await asyncio.sleep(2 ** attempt) # 指数退避:1s → 2s → 4s async def run_batch(path: str, concurrency: int = 10) -> list: tasks = load_jsonl(path) sem = asyncio.Semaphore(concurrency) # 并发度按平台限流调,起步 5~10 return await asyncio.gather(*(process_one(t, sem) for t in tasks)) if __name__ == "__main__": results = asyncio.run(run_batch("batch_input.jsonl")) ok = [r for r in results if r["ok"]] failed = [r for r in results if not r["ok"]] print(f"成功 {len(ok)},失败 {len(failed)}") json.dump(results, open("batch_result.json", "w", encoding="utf-8"), ensure_ascii=False, indent=1)

三个设计点:

  1. 信号量限并发:并发度从 5 起步,被限流就降——把频控当配置而不是敌人
  2. 指数退避重试:1s → 2s → 4s,被限流时等一等比硬重试快得多
  3. 失败单独落盘:ok: false的条目和结果放一起,补跑只跑失败的那几条(幂等:重跑不重复计费)

如果平台提供批量推理接口(上传 JSONL、后台并行、下载结果),优先用它——System Prompt 只写一次,并发和限流平台管,5000 条的成本只有几块钱。自己写并发的价值是灵活:中间步骤要自定义逻辑时才用。

五、结果核验:准确率和混淆矩阵

有标注的数据,核验用评测篇的方法:

from collections import Counter def evaluate(results: list, labels: dict) -> dict: ok = [r for r in results if r["ok"] and r["custom_id"] in labels] correct = sum(1 for r in ok if r["sentiment"] == labels[r["custom_id"]]) confusions = Counter((labels[r["custom_id"]], r["sentiment"]) for r in ok if r["sentiment"] != labels[r["custom_id"]]) print(f"准确率: {correct/len(ok):.1%} ({correct}/{len(ok)})") for (truth, pred), n in confusions.most_common(): print(f" {truth} → 误判为 {pred}: {n} 条") return {"accuracy": correct / len(ok), "confusions": dict(confusions)}

混淆矩阵比准确率值钱:「正面 → 误判为中性」集中在带转折句的评论(「音质不错,就是盒子偏大」),答案就清晰了——prompt 里加一句「转折句看主体情感倾向」,20 条测试集上验证,准确率再上一档。误判案例是 prompt 优化的原料,比盲目改一百遍 prompt 都管用。

六、成本核算(接成本篇)

# 成本篇的 task_cost 原样用:批量场景是 N 次独立调用,无上下文复用 cost = task_cost(rounds=5000, ctx_tokens=85, out_tokens=28, price_in=2, price_out=8) print(cost) # 5000 条分类:几块钱

批量场景省 token 的两刀:System Prompt 尽量短(每条都重发,5000 条 × 长 prompt = 大头)、输出 max_tokens 按需设(分类任务 100~200 足够,默认 4096 是浪费)。

七、踩坑提醒

  1. 别跳过小批量验证:5000 条跑完发现 prompt 有歧义,250 倍的钱和两小时白花——20 条先过三关
  2. 并发度不是越大越好:触发频控后重试反而更慢,从 5 起步逐步加
  3. 温度必须 0:批量分类要稳定,同一批数据重跑结果不一致,混淆矩阵就失真(评测篇 judge 铁律同理)
  4. 失败条目要补跑:指数退避后仍失败的(约 1%),单独补跑而不是全量重跑——幂等设计让补跑成本归零

总结

步骤一句话
数据准备JSONL + custom_id,靠 ID 对齐请求和结果
核心方法论先 20 条过三关,再放大——验证比返工便宜 250 倍
并发信号量限并发 + 指数退避,有批量接口优先用批量
核验准确率 + 混淆矩阵,误判案例是 prompt 优化的原料
成本prompt 短一截、max_tokens 按需,5000 条几块钱

结合系列:结构化输出管格式、成本篇管账、评测篇管核验、可观测性篇管失败日志——批量处理就是这四件套的合体应用。数据量大又不赶时间的任务(标注、分类、摘要),这套链路直接抄。觉得有用点个关注。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询