Agnes-2.5-Flash 批量生成商品描述任务的调度与并发控制

文章导读
几千条商品描述用脚本批量生成,核心不在生成本身,而在把单条调用、并发上限、失败重试和断点续跑拆成四个独立环节。直接 for 循环会因限流中断,手动粘贴又不可维护;可行的做法是先用 CSV 统一输入,串行跑通一条,再逐步放开并发,并把进度写到磁盘,保证中断后可重跑。
📋 目录
  1. 准备商品数据的输入文件格式
  2. 串行调用 Agnes-2.5-Flash 的底层函数
  3. 限制并发数避免触发限流
  4. 实现失败任务的重试与死信记录
  5. 记录进度并支持断点续跑
A A

几千条商品描述用脚本批量生成,核心不在生成本身,而在把单条调用、并发上限、失败重试和断点续跑拆成四个独立环节。直接 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 批量生成商品描述任务的调度与并发控制

串行调用 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。并发数不是固定值,需要结合环境确认。

Agnes-2.5-Flash 批量生成商品描述任务的调度与并发控制

实现失败任务的重试与死信记录

单条失败不应该中断整个批次。常见做法是重试 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 未写入的间隙,最坏会重复处理一条,但不会漏处理。需要结合环境确认你的业务是否容忍这种极小概率的重复。