当你一次向 SenseNova-Vision 提交几千张图片时,最容易遇到的问题就是请求过快:一部分请求超时,另一部分被服务端限流。这个脚本的核心思路不是盲目降低速度,而是把并发控制拆成可观测的四个部分:输入任务的清单化、线程池与信号量、单请求超时与失败重试,最后用小批量样本来确认本环境的合理并发数。
适用场景:用 SenseNova-Vision 批量处理本地图片,且单次请求可能因网络或服务端限制而失败。操作动作:先定义输入清单,再通过线程池控制最大线程数,用信号量控制正在执行的请求数,给每个请求设置超时并记录失败路径,最后用 10、50、100 张图片的样本测试来确认当前网络和服务端容忍的并发边界。验证方式:观察 error.log 中的失败率与耗时变化。风险边界:并发上限由服务端配额和网络决定,不是脚本能单方面保证的值。
定义输入任务清单
程序需要先知道处理哪些图片,以及每张图片对应的结果标识。expected_key 用于校验收到的描述是否对应正确图片,也能作为输出文件名的命名依据。建议用 CSV 或 JSONL 保存任务清单,每行一个任务,避免把大量路径硬编码在脚本里。
# tasks.csv 示例
image_path,expected_key
/path/to/001.jpg,PLANET_001
/path/to/002.jpg,OCTAGON_002
如果字段中可能包含逗号或换行,建议改用 JSONL:
{"image_path": "/path/to/001.jpg", "expected_key": "PLANET_001"}
{"image_path": "/path/to/002.jpg", "expected_key": "OCTAGON_002"}
读取任务清单时,逐行解析成字典,作为后续并发任务的输入对象。
设计线程池大小与信号量
线程池决定了同时会有多少线程在抢任务,信号量则限制同时发出多少个网络请求。如果只控制线程数,一个线程内可能连续循环请求,仍然会打满服务端;如果只控制信号量,任务在池外的排队和回写又不够清晰。两者配合,线程池负责任务调度,信号量负责流量整形。
from concurrent.futures import ThreadPoolExecutor
import threading
MAX_WORKERS = 5 # 先给一个保守值,根据CPU核数和IO耗时调整
MAX_INFLIGHT = 3 # 同时在途的请求数,需要通过小样测试确定
semaphore = threading.Semaphore(MAX_INFLIGHT)
executor = ThreadPoolExecutor(max_workers=MAX_WORKERS)
MAX_INFLIGHT 是真正的并发上限,MAX_WORKERS 可以稍大一些,让线程在等待信号量时也能占用CPU做其他事。建议先保持 MAX_INFLIGHT 较小,再逐步加大。
实现单张图片处理与超时控制
每个任务必须在一个函数内完成请求、超时控制和结果捕获,否则异常会干扰线程池的后续任务。下面给出一个通用骨架,call_sensenova 需要替换为实际的 SenseNova-Vision SDK 或 HTTP 调用。
def process_one(task):
image_path = task["image_path"]
expected_key = task["expected_key"]
with semaphore:
try:
# 替换为实际的 SenseNova-Vision 请求调用
response = call_sensenova(image_path, timeout=30)
return {
"image_path": image_path,
"expected_key": expected_key,
"success": True,
"response": response,
}
except Exception as exc:
return {
"image_path": image_path,
"expected_key": expected_key,
"success": False,
"error": str(exc),
"attempt": 0,
}
timeout=30 是请求超时阈值,超过则抛出异常,避免一个请求卡住整个线程。返回结构固定为字典,方便后续判断成功和重试。
记录失败任务并实现重试
偶发的网络抖动会导致请求失败,直接丢弃会让整个批次不完整。建议给每个任务加上有限次数的重试,并在最终失败时把图片路径写入 error.log,供后续断点续跑。
import time
MAX_RETRIES = 3
def run_task(task):
for attempt in range(MAX_RETRIES):
result = process_one(task)
if result["success"]:
return result
if attempt < MAX_RETRIES - 1:
time.sleep(2 ** attempt) # 退避:1秒、2秒
# 最终失败,写入日志
with open("error.log", "a", encoding="utf-8") as f:
f.write(task["image_path"] + "\n")
return result
run_task 作为线程池的入口函数,内部调用 process_one。重试间隔用指数退避,避免重试风暴。如果连续失败,最终错误日志会记录所有未成功的图片路径。
用小规模样本验证稳定性
在全量处理前,先随机抽取 10、50、100 张图片,分别跑一遍上述流程,记录总耗时、失败数和失败率。这里的“失败”指最终重试后仍失败的任务数量。通过对比,找出当前网络和服务端配额下可接受的 MAX_INFLIGHT 值。
| 样本数 | MAX_INFLIGHT | 总耗时(秒) | 失败数 | 失败率 |
|---|---|---|---|---|
| 10 | 3 | — | — | — |
| 50 | 3 | — | — | — |
| 100 | 3 | — | — | — |
如果某个规模下失败率明显上升,说明并发已经超过本环境容忍值,此时应降低 MAX_INFLIGHT。如果各规模失败率都很低,可以尝试加大 MAX_INFLIGHT 并重新测试。