批量处理文本的标准做法就一句话:把任务写进 input.jsonl,用 ThreadPoolExecutor(max_workers=8) 并发发请求,每条结果即时追加进 output.jsonl,进程启动时先读一遍已完成的 id 并跳过——这样中途断了重跑只补没做完的部分。下面这份脚本包含断点续跑、指数退避重试、结构化输出兜底和成本核算,复制到本地改三行就能跑。
读一条、发一次请求、写一行结果,这个写法在 20 条数据上完全没问题,一到几万条就彻底崩掉。它不是慢一点的问题,是根本跑不完。
一次大模型请求从建连到拿完回答,正常在 1.5 到 3 秒之间。按每条 2 秒算,10 万条就是 100000 乘 2 等于 200000 秒,也就是 55.5 小时。这还没算网络抖动和失败重试,单条变成 3 秒,总时长直接冲到 83 小时,一个周末都收不了尾。要把这个数字压到 2 到 3 小时,只能靠并发,换更快的模型救不了。
串行脚本真正难受的地方在恢复能力。跑到第 8 万条时电脑休眠、SSH 掉线,进程一死,因为没记录哪些已完成,重启只能从第 1 条重新发,前面 8 万条的时间和钱全部作废。所以批量脚本第一个设计目标不是快,是可恢复:每条结果落盘,重跑自动跳过。
接口本身支持并发,串行发等于每个时刻只有一个请求在天上。同一个 10 万条的任务,8 并发把 55 小时压到 7 小时左右,代价只是多写十几行线程池代码。
小鱼API(xyuai.cc)走 OpenAI 兼容协议,/v1/chat/completions 和 /v1/models 都是标准路径,所以直接装官方 openai SDK,把 base_url 指向 https://xyuapi.top/v1 就行,其余代码一个字不用改。备用地址是 https://xyuai.cc/v1,Gemini 原生协议另走 /v1beta。
pip install --upgrade openai
import os
from openai import OpenAI
client = OpenAI(
api_key=os.environ["XYU_API_KEY"], # 令牌从平台后台拿,别写死在源码里
base_url=os.getenv("XYU_BASE_URL", "https://xyuapi.top/v1"),
timeout=60.0,
max_retries=0, # 重试逻辑自己写,方便记日志
)
每行一个 JSON 对象,最少两个字段:id 和 text。id 必须唯一且稳定,它就是断点续跑的钥匙;text 是待处理的原文。文件里不要有空行之外的装饰,也不要有逗号分隔的表头。
{"id": "doc-0001", "text": "这款耳机的降噪在中低频表现稳定,戴一小时耳朵不闷。"}
{"id": "doc-0002", "text": "物流比预计晚了两天,客服回复也慢,但包装没有破损。"}
{"id": "doc-0003", "text": "价格比同规格的便宜一些,做工一般,日常用足够了。"}
CSV 有个致命缺陷:文本里只要出现换行符,整行就被撕成两行,后面字段全部错位,而大模型返回的文本里换行几乎不可避免。jsonl 一行一条记录,内部换行会被转义,撑不破行结构,还能追加写、按行跳过坏行,这两点对断点续跑是刚需。
下面这份脚本是完整版,直接存成 batch_run.py 就能跑。它做了四件事:断点续跑、指数退避重试、结果即时落盘、进度与花费实时打印。
脚本启动后先扫 output.jsonl,把里面所有 id 收进一个集合,处理时跳过。这样无论跑到哪一步挂掉,重新执行同一条命令就自动接着跑。注意容错:进程被强杀时最后一行可能只写了一半,解析失败直接跳过,别让它把整个任务拦死。
429、超时、连接被重置这些都是暂时性故障,等一会儿重发大概率成功。退避就是让等待时间指数增长:1 秒、2 秒、4 秒、8 秒。还要加一点随机抖动,否则几十个线程会在同一秒一起重发,形成二次冲击,把限流撞得更狠。
线程池里多个线程同时写同一个文件会互相插队,写出半行、乱码、两条结果粘在一行。用一个 threading.Lock() 把「序列化加写一行加刷盘」包起来,写完立刻 os.fsync,断电也不会丢已完成记录。
Windows 的 PowerShell 串命令用分号,用两个与号会直接报语法错误。给每个批次单独的输出文件,跑完再合并。
python batch_run.py
$env:XYU_API_KEY="你的令牌"; $env:MAX_WORKERS="8"; python batch_run.py
# 输入文件太大先切分,每 2 万条一批,切完分别跑,互不影响
split -l 20000 input.jsonl part_
完整脚本:
#!/usr/bin/env python
# -*- coding: utf-8 -*-
"""批量文本处理:并发 + 断点续跑 + 指数退避重试 + 结构化输出兜底"""
import json
import os
import random
import re
import threading
import time
from concurrent.futures import ThreadPoolExecutor, as_completed
from openai import OpenAI
BASE_URL = os.getenv("XYU_BASE_URL", "https://xyuapi.top/v1")
API_KEY = os.getenv("XYU_API_KEY", "")
MODEL = os.getenv("XYU_MODEL", "gemini-2.5-pro")
IN_PATH, OUT_PATH, ERR_PATH = "input.jsonl", "output.jsonl", "error.jsonl"
MAX_WORKERS = int(os.getenv("MAX_WORKERS", "8"))
MAX_RETRY = 4
# 按次计费单价,用来实时估算花费
PRICE = {
"gemini-2.5-pro": 0.031,
"deepseek-v3.2-thinking": 0.049,
"deepseek-r1-thinking": 0.049,
"deepseek-v4-flash-thinking": 0.05,
"grok-4.1": 0.05,
"gemini-3-pro-preview": 0.05,
"grok-4.5": 0.09,
"grok-4.6": 0.09,
"claude-sonnet-4-5-thinking": 0.09,
"deepseek-v4-pro-thinking": 0.09,
"kimi-k2.5": 0.09,
"kimi-k2.6": 0.09,
"gemini-3.1-pro-preview": 0.09,
"claude-opus-4-5-thinking": 0.12,
"gpt-5.4-pro-thinking": 0.15,
"claude-sonnet-4-6-thinking": 0.2,
"claude-sonnet-4-7-thinking": 0.2,
"gpt-5.5": 0.2,
"gpt-5.5-pro-thinking": 0.2,
"claude-opus-4-6-thinking": 0.25,
"claude-opus-4-7-thinking": 0.25,
"gpt-5.3-pro": 0.3,
}
UNIT = PRICE.get(MODEL, 0.031)
client = OpenAI(api_key=API_KEY, base_url=BASE_URL, timeout=60.0, max_retries=0)
write_lock = threading.Lock()
def load_done(path=OUT_PATH):
"""扫一遍已完成 id,实现断点续跑"""
done = set()
if not os.path.exists(path):
return done
with open(path, "r", encoding="utf-8-sig") as f:
for line in f:
line = line.strip()
if not line:
continue
try:
done.add(json.loads(line)["id"])
except (json.JSONDecodeError, KeyError):
continue # 尾行被强杀截断属正常,跳过即可
return done
def load_tasks(path=IN_PATH):
tasks = []
with open(path, "r", encoding="utf-8-sig") as f: # utf-8-sig 兼容 BOM
for n, line in enumerate(f, 1):
line = line.strip()
if not line:
continue
try:
tasks.append(json.loads(line))
except json.JSONDecodeError as e:
print("[WARN] 第 %d 行不是合法 JSON,已跳过:%s" % (n, e))
return tasks
def append_line(path, obj):
"""加锁写一行并刷盘,避免多线程写乱文件"""
with write_lock:
with open(path, "a", encoding="utf-8") as f:
f.write(json.dumps(obj, ensure_ascii=False) + "\n")
f.flush()
os.fsync(f.fileno())
def extract_json(text):
"""response_format 失效时的兜底:抽出第一个花括号块"""
try:
return json.loads(text)
except json.JSONDecodeError:
pass
m = re.search(r"\{[\s\S]*\}", text or "")
if not m:
return None
try:
return json.loads(m.group(0))
except json.JSONDecodeError:
return None
PROMPT_TMPL = (
"你是文本结构化助手。请针对下面的文本输出 JSON,字段固定为:\n"
'{"summary": "不超过50字的中文摘要", "tags": ["标签1", "标签2"], '
'"sentiment": "正面/中性/负面"}\n'
"只输出 JSON,不要任何解释文字。\n\n文本:\n%s"
)
def call_one(item):
"""处理单条:失败按指数退避重试,最终失败返回 ok=False"""
err = None
for attempt in range(MAX_RETRY):
try:
resp = client.chat.completions.create(
model=MODEL,
messages=[{"role": "user", "content": PROMPT_TMPL % item["text"]}],
temperature=0.2,
max_tokens=800,
response_format={"type": "json_object"},
)
raw = resp.choices[0].message.content or ""
data = extract_json(raw)
if data is None:
raise ValueError("返回内容无法解析为 JSON:%s" % raw[:120])
return {"id": item["id"], "ok": True, "result": data}
except Exception as e:
err = e
wait = (2 ** attempt) + random.uniform(0, 1)
print("[RETRY] id=%s 第 %d 次失败 %s:%.1fs 后重试"
% (item["id"], attempt + 1, type(e).__name__, wait))
time.sleep(wait)
return {"id": item["id"], "ok": False, "error": repr(err)}
def main():
tasks, done = load_tasks(), load_done()
todo = [t for t in tasks if t["id"] not in done]
print("总计 %d 条,已完成 %d 条,本轮待处理 %d 条" % (len(tasks), len(done), len(todo)))
ok = fail = 0
t0 = time.time()
with ThreadPoolExecutor(max_workers=MAX_WORKERS) as pool:
futures = [pool.submit(call_one, t) for t in todo]
for i, fut in enumerate(as_completed(futures), 1):
res = fut.result()
if res["ok"]:
append_line(OUT_PATH, res)
ok += 1
else:
append_line(ERR_PATH, res)
fail += 1
if i % 50 == 0 or i == len(todo):
print("[进度] %d/%d 成功=%d 失败=%d 耗时=%.0fs 花费≈%.2f元"
% (i, len(todo), ok, fail, time.time() - t0, (ok + fail) * UNIT))
print("完成:成功 %d,失败 %d,耗时 %.1f 分钟,花费≈%.2f 元"
% (ok, fail, (time.time() - t0) / 60, (ok + fail) * UNIT))
if __name__ == "__main__":
main()
并发不是越大越快。上游对单位时间内的请求数有阈值,你同时发出去的请求超过阈值,接口就回 429 Too Many Requests,SDK 里抛的是 openai.RateLimitError。重试还会把压力叠上去:8 个线程各重试 3 次,瞬时并发就变成 24,限流更严重,总耗时反而比稳稳的 4 并发更长。
| max_workers | 10 万条预计耗时 | 撞 429 的概率 | 建议场景 |
|---|---|---|---|
| 1 | 约 55 小时 | 几乎没有 | 只跑几百条,不赶时间 |
| 4 | 约 14 小时 | 偶尔 | 长文本,单条经常超过 3 秒 |
| 8 | 约 7 小时 | 需配合限速 | 大多数批量任务,推荐起点 |
| 16 | 约 3.5 小时 | 明显变高 | 短文本,能承受较多重试 |
| 32 | 约 2 小时 | 大量 429,重试拖慢 | 不推荐,通常净收益为负 |
不要一上来就写 32。先拿 200 条做试跑,按 4 并发跑一轮,看日志里有没有 RETRY 行。没有就提到 8,再跑 200 条;还是没有重试就提到 16。一旦看到 429 或者 Connection reset by peer 频繁出现,立刻退回上一档并加令牌桶限速。这个爬坡十分钟能做完,比事后查限流便宜得多。
response_format={"type": "json_object"} 只是告诉模型「输出必须是合法 JSON」,它管不了字段名。字段叫什么、有几个、每个的类型,全部要靠提示词写死,并且给一个带假数据的完整示例,模型会照着填。
总会有几次返回带解释文字,比如「好的,以下是结果:{\"summary\":...}」。这时候不要整条判为失败,先 json.loads 试一次,失败就用 re.search(r"\{[\s\S]*\}", text) 把花括号块抽出来再解析。上面脚本里的 extract_json 就是这个兜底,能救回大部分脏输出。抽出来仍然失败,才把原始文本截前 120 个字符记进错误文件。
即使返回了合法 JSON,也可能缺 tags,或者把 sentiment 写成 情绪。写库或写 CSV 之前统一用 .get() 取值并给默认值,不要用中括号硬取,否则一条脏数据就能让写库步骤崩掉。要求更严时,可以在校验失败后补一次「严格按字段重发」的重试。
单条文本超过两千字之后,模型容易漏掉中间内容,输出质量明显下降,所以长文本必须先切块。硬按字符数切是最糟的做法,会把一句话从中截断,两块都读不通。正确做法是先按空行拆段落,再把段落累加到一个块里,累加超过上限才开新块。
import re
def chunk_by_paragraph(text, max_chars=1800):
"""按段落聚合切块,尽量不切断句子;单段超长再按中文句末符切"""
paragraphs = [p.strip() for p in re.split(r"\n\s*\n", text) if p.strip()]
chunks, buf = [], ""
for p in paragraphs:
if len(p) > max_chars: # 超长单段:按句末标点切
for sent in re.split(r"(?<=[。!?;])", p):
if not sent:
continue
if buf and len(buf) + len(sent) > max_chars:
chunks.append(buf)
buf = sent
else:
buf += sent
continue
if buf and len(buf) + len(p) + 1 > max_chars:
chunks.append(buf)
buf = p
else:
buf = (buf + "\n" + p) if buf else p
if buf:
chunks.append(buf)
return chunks
每块返回的摘要不要简单拼接,那样会变成五段废话。正确顺序是先用每块的 tags 做并集去重,再把各块摘要按原顺序拼起来交给模型做二次压缩,sentiment 走多数投票。多花一次请求,但结果读起来是一篇完整的东西。
切块必然切断上下文,第二块里出现「它」「该公司」时模型不知道指谁。解决办法是每块的提示词前面拼上上一块的最后 200 字,标成「上文仅供参考,不用处理」。按次计费下输入长度不影响价格,这点成本几乎为零。
令牌桶的核心是:桶里按固定速率补充令牌,每个请求先取一个,取不到就等。它比固定 time.sleep 更平滑,允许短时突发,又不会长时间超速。
import threading
import time
class TokenBucket:
"""rate = 每秒允许的请求数,burst = 允许的瞬时突发量"""
def __init__(self, rate, burst=None):
self.rate = float(rate)
self.capacity = float(burst or rate)
self.tokens = self.capacity
self.updated = time.monotonic()
self.lock = threading.Lock()
def acquire(self):
while True:
with self.lock:
now = time.monotonic()
self.tokens = min(self.capacity,
self.tokens + (now - self.updated) * self.rate)
self.updated = now
if self.tokens >= 1:
self.tokens -= 1
return
need = (1 - self.tokens) / self.rate
time.sleep(need)
LIMITER = TokenBucket(rate=6, burst=8) # 稳态 6 次/秒,最多瞬时 8 次
把它放进 call_one,位置在 client.chat.completions.create 之前一行:LIMITER.acquire()。全局一个实例,所有线程共用。
嫌上面的类麻烦,可以用提交节奏控制:主线程往线程池塞任务时每次睡 0.12 秒,等于把提交速率压在每秒 8 条左右,配合 8 个 worker 就够了。
with ThreadPoolExecutor(max_workers=8) as pool:
for t in todo:
pool.submit(call_one, t)
time.sleep(0.12) # 约 8 次/秒的提交节奏
看到 429 别只是加长重试等待,那治标不治本。正确的处理是动态降速:把令牌桶的 rate 乘 0.7,把 max_workers 往下调一档,然后重试。连续十次请求都没再出现 429,再把速率加回来。这套自适应逻辑二十行代码就能写完。
| 报错原文 | 含义 | 你该怎么做 |
|---|---|---|
429 Too Many Requests | 瞬时请求数超阈值 | 降 max_workers,令牌桶按 0.7 系数降速 |
openai.RateLimitError | SDK 封装的限流异常 | 同上,同时确认令牌额度是否已用尽 |
openai.APITimeoutError / Timeout | 请求超时 | 客户端设 timeout=60,重试两次 |
Connection reset by peer | 连接被中途断开 | 降并发,加抖动后重试 |
400 Bad Request | 参数或模型名不对 | 重试无用,打印原文人工核对 |
UnicodeDecodeError | 文件带 BOM 或有非 UTF-8 字节 | 用 encoding="utf-8-sig" 读 |
小鱼API 以按次计费为主:一次请求一个固定价,输入多长都不影响这一单的价格。所以批量场景最有效的降本手段是合并请求——把 10 条短文本塞进一次调用,让模型返回一个 JSON 数组,成本立刻变成原来的十分之一。实际合并时控制在 5 到 10 条、总输入 2000 到 3000 字以内,太多会撞上输出截断。
PROMPT = (
"下面是一个 JSON 数组,每项有 id 和 text。请对每项生成摘要,"
"严格返回同长度的 JSON 数组,每项形如 "
'{"id":"原样返回","summary":"不超过50字"}' + ",不要解释。\n" + batch_json
)
resp = client.chat.completions.create(
model="gemini-2.5-pro",
messages=[{"role": "user", "content": PROMPT}],
response_format={"type": "json_object"},
)
# 返回后按下标与输入 id 对齐,任何一个 id 对不上就整批降级为单条重跑
按次计费单价固定,成本就是「调用次数乘单价」,没有输入长度变量。下表全部按 30 天算。
| 日处理量 | 模型 | 单次价 | 每日调用次数 | 日成本 | 月成本 |
|---|---|---|---|---|---|
| 100 条 | gemini-2.5-pro | 0.031 元 | 100 | 3.1 元 | 93 元 |
| 1,000 条 | gemini-2.5-pro | 0.031 元 | 1000 | 31 元 | 930 元 |
| 1,000 条(10 条合 1) | gemini-2.5-pro | 0.031 元 | 100 | 3.1 元 | 93 元 |
| 10,000 条 | deepseek-v3.2-thinking | 0.049 元 | 10000 | 490 元 | 14,700 元 |
| 10,000 条(10 条合 1) | deepseek-v3.2-thinking | 0.049 元 | 1000 | 49 元 | 1,470 元 |
| 10,000 条 | claude-opus-4-5-thinking | 0.12 元 | 10000 | 1,200 元 | 36,000 元 |
假设你一次性清洗 10 万条评论,用 gemini-2.5-pro 逐条发,就是 100000 乘 0.031 等于 3100 元,这个数字已经超过很多小项目的整月预算。而它还是列表里单价最低的一档——换成 claude-opus-4-5-thinking 的 0.12 元一次,同样的量是 12000 元。所以批量场景的选型原则很清楚:优先挑 0.031 到 0.09 元这一档的模型,再叠加 10 条合 1 的请求合并,10 万条的实际花费能从 3100 元压到 310 元。
摘要、分类、打标签、情感判断这类结构化任务,gemini-2.5-pro 的 0.031 元一次足够用;需要多步推理、要求给出理由的任务,用 deepseek-v3.2-thinking 或 deepseek-r1-thinking 的 0.049 元档;文案改写、长文写作这类对表达质量敏感的任务才考虑 claude-sonnet-4-5-thinking 的 0.09 元档。原则是先跑 200 条样本,用低价那一档试,只有质量明确不够时才往上升。
在 Windows 上用 Excel 打开过再另存的文件,开头会多出三个字节的 BOM,Python 按 utf-8 读第一行就会抛 UnicodeDecodeError: 'utf-8' codec can't decode byte 0xef in position 0。解决办法统一用 encoding="utf-8-sig" 打开输入文件,它会自动吃掉 BOM;实在遇到混合编码的历史数据,再退一步加 errors="replace" 把坏字节替换掉,同时在日志里记下被替换的行号。
用 csv 模块写结果时不要手工拼字符串加逗号,字段里的逗号、换行、引号都要靠 csv.writer 的自动引用处理。用 pandas 读的时候显式指定 quoting=csv.QUOTE_ALL,或者干脆继续用 jsonl 输出——结果文件通常还要再进一次脚本,jsonl 更省事也更容易定位坏行。
输出到一半断掉、JSON 缺右花括号,绝大多数是 max_tokens 太小。中文一个字符通常占一到两个 token,写 800 的 max_tokens 只够四百来字,长摘要直接腰斩。按你的目标输出长度乘 2 再留三成余量;思考型模型还会额外消耗推理 token,这时候把 max_tokens 提到 2000 以上更保险。
10 万条跑完,如果把每条原文和完整返回都写进日志,日志会膨胀到几百兆,出问题时光是打开都费劲。日志里只留 id、耗时、状态、错误类型和重试次数这类短字段,原文和返回写进 jsonl 结果文件。再给日志按天分文件并定期清理,排查时按 id 定位一条,比翻整个大文件快得多。