手头有几十上百条录音时,逐条跑太慢,但直接改成并发循环,中途一旦断开,最难回答的问题就是「哪条成了、哪条没成」。批量转写要先定的不是并发数字,而是四件事:一份能查状态的清单、一个不超过后端承受能力的排队节奏、一条「什么样的失败才值得重试」的边界,以及失败后把原始返回留档。这几件事定下来之前,脚本跑得越快越难收拾。
建议的顺序是:先扫目录生成带状态的清单,再用带并发上限的循环消费这份清单——成功改状态,可重试的失败按上限重试,不可重试的失败直接归档。并发上限和重试次数不要凭感觉写死,通常从保守值起步,看日志里的失败类型再调。清单只记录可核验的字段,跑完以清单对账为准,不靠翻文件夹判断。
把待处理文件整理成一份带状态的清单
清单是整套流程的唯一事实来源。字段建议固定为:文件路径、状态、开始时间、结束时间、备注,再加一个重试次数字段,方便后面判断是不是重试耗尽的失败。状态取值尽量少,用 pending、running、success、failed 四种就够了;如果跑完后还有 running 残留,基本说明进程中途被杀,那条需要补跑。
初始清单由一次目录扫描生成,不要手写。放在批量脚本同目录下执行即可:
import csv, pathlib
root = pathlib.Path("/data/audio")
exts = {".wav", ".mp3", ".m4a", ".flac"}
files = sorted(p for p in root.rglob("*") if p.suffix.lower() in exts)
with open("manifest.csv", "w", newline="", encoding="utf-8") as f:
w = csv.writer(f)
w.writerow(["path", "status", "started_at", "finished_at", "retry", "note"])
for p in files:
w.writerow([str(p), "pending", "", "", 0, ""])
print("待处理", len(files))
清单生成后先目测一下条数和实际文件数是否对得上,再开始跑。后续每一步都回写这份 CSV,不要在内存里记账。
给循环加上并发上限和间隔
并发控制的目的是把失败范围控制在小批次内,让失败原因容易区分,而不是把机器打满。通用骨架是「固定大小的线程池 + 每提交一条任务后短暂间隔」,两个参数都放在脚本顶部,方便改完再跑:
from concurrent.futures import ThreadPoolExecutor, as_completed
import time
MAX_WORKERS = 3 # 先保守起步,看日志再调
SUBMIT_GAP = 0.5 # 每次提交之间的间隔,按需调整
rows = [r for r in read_manifest("manifest.csv") if r["status"] != "success"]
with ThreadPoolExecutor(max_workers=MAX_WORKERS) as pool:
futures = {}
for row in rows:
futures[pool.submit(transcribe_one, row)] = row
time.sleep(SUBMIT_GAP)
for fut in as_completed(futures):
row = futures[fut]
try:
handle_success(row, fut.result())
except Exception as e:
handle_failure(row, e)
这里的 transcribe_one 换成自己的调用封装即可,输入是清单里的一行,输出是转写结果或抛异常。调整参数后要看的观察项:日志里超时和限流的占比有没有明显上升、清单里残留的 running 条数、单条耗时是否出现长尾、后端是否开始主动拒绝请求。如果拒绝变多,先降并发或加大间隔,而不是加重试。
定义什么样的失败才值得重试以及重试上限
重试只应该处理偶发问题。可重试的通常是连接超时、读写超时、连接被重置、429 限流、500/502/503/504 这类服务端临时错误。不可重试的通常是:音频文件为空或打不开、格式不在支持范围内、请求参数写错、鉴权失败、返回体里明确说明内容不合规。把后者也丢进重试队列,只会把真实问题掩盖成「重试了几次还是不行」。
import requests
RETRYABLE_EXC = (requests.Timeout, requests.ConnectionError)
RETRYABLE_CODE = {429, 500, 502, 503, 504}
def classify(exc=None, resp=None):
if exc is not None and isinstance(exc, RETRYABLE_EXC):
return "retry"
if resp is not None and resp.status_code in RETRYABLE_CODE:
return "retry"
return "fatal"
重试次数记录在清单的 retry 字段里,同时把每次尝试的时间写进运行日志,两个地方对得上才算数。达到自己设定的上限(例如三次)后,把状态改成 failed、备注写「重试耗尽」,然后交给下一步归档,不要无限重试。
失败条目单独归档并保留原始返回
归档的目标是事后能复现,而不是只看到一个错误码。建议按运行批次建目录,失败条目一条一个文件,另外留一份汇总:
runs/run-<时间戳>/
manifest.csv # 本次运行的清单快照
cases.jsonl # 每条失败的摘要,一行一条
raw/
<path-hash>.json # 单条失败的完整现场
每条失败至少保留:请求参数(文件路径、模型名、采样率等,密钥不要写入)、原始返回内容(状态码和响应体)、报错文本、发生时间、已重试次数。cases.jsonl 里的一行大概长这样:
{"path": "/data/audio/017.wav", "attempts": 3, "kind": "retry",
"error": "read timeout", "request": {"model": "qwen-audio", "path": "/data/audio/017.wav"},
"response": {"status": 500, "body": "..."}, "at": "..."}
归档和清单要互相指向:cases.jsonl 的 path 字段必须能在 manifest.csv 里查到,否则说明有失败被漏记。
跑完后用清单对账确认没有漏处理的文件
对账只做三件事:数状态、比总数、找残留。一条命令打印各状态计数和总条数:
python - <<'PY'
import csv, collections
rows = list(csv.DictReader(open("manifest.csv", encoding="utf-8")))
print(collections.Counter(r["status"] for r in rows))
print("总数", len(rows))
for r in rows:
if r["status"] in ("pending", "running"):
print("残留", r["path"])
PY
核对方式是:success 条数加上 failed 条数应等于总数,同时 cases.jsonl 的行数应等于 failed 条数。出现 pending 或 running 残留,通常是进程中断,先把它们改回 pending 再补跑;补跑时只挑 status 不等于 success 的条目生成新清单,跑完再对账一次。两次对账结果一致,这份结果才可以直接交付使用。