几千条商品描述用脚本批量生成,核心不在生成本身,而在把单条调用、并发上限、失败重试和断点续跑拆成四个独立环节。直接 for 循环会因限流中断,手动粘贴又不可维护;可行的做法是先用 CSV 统一输入,串行跑通一条,再逐步放开并发,并把进度写到磁盘,保证中断后可重跑。
适用场景:商品数量在千条量级、使用 Agnes-2.5-Flash API 生成描述、可接受脚本化处理的项目。操作动作:统一 CSV 输入、串行验证单条生成、再以 5 个线程并发、失败重试 3 次并记录死信。验证方式:先打印前三条记录,再观察并发日志中的请求间隔,最后中断重跑对比输出行数。风险边界:并发数需按实际账号限流调整,重试只解决临时错误,持续 4xx/5xx 应优先检查密钥和配额。
准备商品数据的输入文件格式
输入文件建议使用 CSV,字段固定为 id、title、keywords,其中 id 用于断点续跑,title 和 keywords 是生成描述的主要输入。示例文件 product_list.csv:
id,title,keywords
1001,无线蓝牙耳机,降噪 长续航 运动
1002,便携榨汁杯,USB充电 随行杯 水果
读取时用 pandas 或标准库 csv 均可。pandas 适合后续做字段清洗和批量处理:
import pandas as pd
df = pd.read_csv("product_list.csv", dtype={"id": str})
print(df.head(3))
print(df.columns.tolist()) # 确认列名一致
若不想引入 pandas,用 csv 模块同样直接:
import csv
with open("product_list.csv", encoding="utf-8") as f:
rows = list(csv.DictReader(f))
print(rows[:3])
验证方式:运行后能打印出前三条记录的 id 和 title,且没有 KeyError,说明字段读取成功。这一步不建议跳过,输入格式错了后面再排查耗时更多。
串行调用 Agnes-2.5-Flash 的底层函数
先写一个独立的生成函数,输入标题和关键词,返回描述文本。不要急着加并发,先把单条链路跑通。若使用 OpenAI 兼容接口,骨架如下:
from openai import OpenAI
client = OpenAI(
api_key="你的密钥",
base_url="你的接口地址" # 按实际服务商配置
)
def generate_description(title: str, keywords: str) -> str:
prompt = f"为商品生成一段中文描述。\n商品标题:{title}\n关键词:{keywords}\n要求:描述自然、突出卖点,不超过100字。"
resp = client.chat.completions.create(
model="agnes-2.5-flash",
messages=[{"role": "user", "content": prompt}],
temperature=0.7
)
return resp.choices[0].message.content.strip()
如果服务商只提供 HTTP 接口,用 requests 也可以;重点是把返回文本提取出来,并对空返回做异常处理。验证方式是运行一次:
print(generate_description("无线蓝牙耳机", "降噪 长续航 运动"))
返回结果非空、内容与商品相关,说明函数可用。这一步不要跳过,串行失败时不要进入并发阶段。
限制并发数避免触发限流
并发数不是越大越好。Agnes-2.5-Flash 的限流阈值取决于账号套餐和接口服务商,通常建议从较小的并发开始,比如 5 个线程。用 ThreadPoolExecutor 控制并发,并在每个任务里打印当前活跃线程数,方便观察压力:
from concurrent.futures import ThreadPoolExecutor, as_completed
import threading
import time
def safe_generate(row):
active = threading.active_count()
print(f"处理 id={row['id']} 活跃线程数={active}", flush=True)
desc = generate_description(row["title"], row["keywords"])
return row["id"], desc
with ThreadPoolExecutor(max_workers=5) as executor:
futures = {executor.submit(safe_generate, row): row for row in rows}
for f in as_completed(futures):
row = futures[f]
try:
id_, desc = f.result()
print(f"id={id_} 生成成功,长度={len(desc)}")
except Exception as e:
print(f"id={row['id']} 失败: {e}")
验证方式:观察日志中“活跃线程数=6”左右(主线程加 5 个工作线程),且两次任务之间的间隔保持稳定。如果出现大量超时或 429 错误,就把 max_workers 调低到 3 或 2。并发数不是固定值,需要结合环境确认。
实现失败任务的重试与死信记录
单条失败不应该中断整个批次。常见做法是重试 3 次,重试间隔按 2 的指数增长,即 2 秒、4 秒、8 秒,并加上小幅随机抖动,避免多个线程同时重试造成瞬间请求叠加。若 3 次都失败,就把该行写入 error.csv,继续处理其他商品:
import random
import time
def generate_with_retry(row, retries=3):
last_exc = None
for attempt in range(retries):
try:
desc = generate_description(row["title"], row["keywords"])
return row["id"], desc
except Exception as e:
last_exc = e
sleep_time = 2 ** attempt + random.random()
print(f"id={row['id']} 第{attempt+1}次失败,{sleep_time:.1f}s后重试")
time.sleep(sleep_time)
print(f"id={row['id']} 重试耗尽,写入 error.csv")
with open("error.csv", "a", encoding="utf-8") as ef:
ef.write(f"{row['id']},{row['title']},{row['keywords']},{last_exc}\n")
raise last_exc # 或者返回 None,由外部决定是否继续
验证方式:临时把 API 密钥改成一个错误值,运行脚本,应能看到每个任务重试 3 次后依次写入 error.csv,且脚本不会崩在中途。恢复密钥后,error.csv 里的记录就是需要手工或二次处理的内容。
记录进度并支持断点续跑
几千条任务可能跑几十分钟甚至更久,中途断网或停电都会导致前功尽弃。建议在生成成功后将商品 id 追加写入 done.txt,每次启动脚本时先读取已完成 id,跳过这些行。这样重跑时只处理尚未成功的商品:
import os
def load_done_ids():
if not os.path.exists("done.txt"):
return set()
with open("done.txt", encoding="utf-8") as f:
return {line.strip() for line in f if line.strip()}
done_ids = load_done_ids()
remaining = [r for r in rows if r["id"] not in done_ids]
print(f"剩余待处理:{len(remaining)} 条")
for row in remaining:
try:
id_, desc = generate_with_retry(row)
# 保存描述到输出文件
with open("output.csv", "a", encoding="utf-8") as out:
out.write(f"{id_},{desc}\n")
with open("done.txt", "a", encoding="utf-8") as done:
done.write(f"{id_}\n")
except Exception:
# 已写入 error.csv,跳过后续处理
continue
验证方式:先运行脚本处理一部分商品,然后强制中断(Ctrl+C 或杀进程)。重新运行脚本,应看到“剩余待处理”数量比上次少,且 output.csv 中没有重复行。如果中断发生在 write 刚完成但 done.txt 未写入的间隙,最坏会重复处理一条,但不会漏处理。需要结合环境确认你的业务是否容忍这种极小概率的重复。