Agent 批量图像生成流水线
手把手教程:构建 Agent 批量图像生成流水线——从产品描述到变体生成、自动评分、最佳选取,含并行处理、错误重试和成本追踪的完整 Python 代码。
结论先行 — 本教程构建一条完整的批量图像生成流水线:Agent 接收 100 个产品描述,每个生成 5 个图像变体(共 500 张),用视觉模型评分,选出最佳。完整 Python 代码基于 OpenAI SDK + SandBase,含并行执行、错误处理、重试逻辑和成本追踪。预计成本:$7.50–$20.00(视模型选择)。
我们要构建什么
一条面向电商的生产级图像生成流水线:
- 接收产品描述作为输入(名称、品类、核心卖点)
- 每个产品用不同提示词/角度生成 5 个图像变体
- 用视觉模型对每个变体评分(画质 + 相关性)
- 为每个产品选出最佳变体
- 追踪成本、处理失败、全程日志
这个模式适用于任何批量生成场景——营销活动、社交内容日历、产品目录刷新、A/B 测试素材。
模型选择指南参见最佳 AI 图像生成 API。视频领域的类似模式参见视频生成 Agent 广告创意教程。
架构概览
┌───────────────────────────────────────────────────┐
│ 批量流水线 │
├───────────────────────────────────────────────────┤
│ 输入:100 个产品描述 │
│ ↓ │
│ 提示词生成器(每产品 5 条 = 500 条提示词) │
│ ↓ │
│ 并行图像生成(10 并发批次) │
│ ↓ │
│ 质量评分器(视觉模型评估) │
│ ↓ │
│ 选择器(每产品选最佳) │
│ ↓ │
│ 输出:100 张最佳图 + 元数据 │
└───────────────────────────────────────────────────┘
前置依赖
pip install openai asyncio aiofiles tenacity pydantic
步骤一:定义数据模型
from pydantic import BaseModel
from typing import Optional
from datetime import datetime
class Product(BaseModel):
id: str
name: str
category: str
features: list[str]
style_preference: Optional[str] = None
class ImageVariant(BaseModel):
product_id: str
variant_index: int
prompt: str
image_url: Optional[str] = None
score: Optional[float] = None
error: Optional[str] = None
latency_ms: int = 0
cost: float = 0.0
class PipelineResult(BaseModel):
product_id: str
best_variant: Optional[ImageVariant] = None
all_variants: list[ImageVariant] = []
total_cost: float = 0.0
步骤二:提示词生成
每个产品获得 5 条不同角度的提示词以最大化多样性:
class PromptGenerator:
"""为每个产品生成多样化提示词。"""
ANGLES = [
("主图", "纯白背景,摄影棚灯光,产品居中,电商主图风格,超清细节"),
("场景", "产品使用场景,自然环境,生活方式摄影,暖光,真实感"),
("特写", "微距特写,展示材质纹理和做工细节,浅景深,棚拍"),
("陈列", "产品置于精心布置的桌面,搭配互补道具,平铺或货架,编辑风格"),
("氛围", "戏剧性灯光,深色背景加轮廓光,高端质感,高对比,奢华展示"),
]
def generate_prompts(self, product: Product) -> list[str]:
"""为一个产品生成 5 条多样化提示词。"""
prompts = []
features_text = "、".join(product.features[:3])
for angle_name, angle_style in self.ANGLES:
prompt = (
f"{product.name},{features_text},"
f"{angle_style},"
f"专业产品摄影,8K 画质"
)
if product.style_preference:
prompt += f",{product.style_preference}"
prompts.append(prompt)
return prompts
步骤三:并行图像生成 + 错误处理
import asyncio
import time
import aiohttp
from tenacity import retry, stop_after_attempt, wait_exponential
class BatchImageGenerator:
"""并行生成图像,含限流和错误处理。"""
def __init__(
self,
api_key: str,
model: str = "bytedance/seedream/5.0/pro/fast",
max_concurrent: int = 10,
cost_per_image: float = 0.015
):
self.api_key = api_key
self.base_url = "https://api.sandbase.ai/v1"
self.model = model
self.semaphore = asyncio.Semaphore(max_concurrent)
self.cost_per_image = cost_per_image
self.total_cost = 0.0
self.success_count = 0
self.failure_count = 0
@retry(stop=stop_after_attempt(3), wait=wait_exponential(min=2, max=30))
async def _generate_single(self, prompt: str) -> dict:
"""单张图像生成(提交 + 轮询),带重试。"""
async with self.semaphore:
start = time.time()
headers = {
"Authorization": f"Bearer {self.api_key}",
"Content-Type": "application/json",
}
async with aiohttp.ClientSession() as session:
# 提交生成任务
async with session.post(
f"{self.base_url}/run",
headers=headers,
json={"model": self.model, "prompt": prompt},
) as resp:
submit = await resp.json()
task_id = submit["id"]
# 轮询等待完成
poll_headers = {"Authorization": f"Bearer {self.api_key}"}
while True:
async with session.get(
f"{self.base_url}/generations/{task_id}",
headers=poll_headers,
) as resp:
result = await resp.json()
if result["status"] in ("completed", "failed", "timeout"):
break
await asyncio.sleep(2)
latency = int((time.time() - start) * 1000)
return {"url": result["outputs"][0]["url"], "latency_ms": latency}
async def generate_variant(
self, product_id: str, variant_index: int, prompt: str
) -> ImageVariant:
"""生成一个变体,优雅处理错误。"""
variant = ImageVariant(
product_id=product_id,
variant_index=variant_index,
prompt=prompt
)
try:
result = await self._generate_single(prompt)
variant.image_url = result["url"]
variant.latency_ms = result["latency_ms"]
variant.cost = self.cost_per_image
self.total_cost += self.cost_per_image
self.success_count += 1
except Exception as e:
variant.error = str(e)
self.failure_count += 1
return variant
async def generate_batch(
self, products: list[Product], prompts_per_product: dict[str, list[str]]
) -> list[PipelineResult]:
"""并行生成所有产品的所有变体。"""
all_tasks = []
for product in products:
prompts = prompts_per_product[product.id]
for i, prompt in enumerate(prompts):
all_tasks.append((product.id,
self.generate_variant(product.id, i, prompt)))
tasks = [task for _, task in all_tasks]
product_ids = [pid for pid, _ in all_tasks]
variants = await asyncio.gather(*tasks)
results_by_product: dict[str, list[ImageVariant]] = {}
for pid, variant in zip(product_ids, variants):
results_by_product.setdefault(pid, []).append(variant)
return [
PipelineResult(
product_id=p.id,
all_variants=results_by_product.get(p.id, []),
total_cost=sum(v.cost for v in results_by_product.get(p.id, []))
)
for p in products
]
def report(self):
total = self.success_count + self.failure_count
print(f"生成完成: {self.success_count}/{total} 张")
print(f"失败: {self.failure_count}/{total}")
print(f"总成本: ${self.total_cost:.2f}")
print(f"成功率: {self.success_count/max(total,1)*100:.1f}%")
步骤四:视觉模型评分
from openai import AsyncOpenAI
class ImageScorer:
"""用视觉模型为图像评分(画质 + 相关性)。"""
def __init__(self, api_key: str):
self.client = AsyncOpenAI(
base_url="https://api.sandbase.ai/v1",
api_key=api_key
)
async def score_variant(self, variant: ImageVariant, product: Product) -> float:
"""为一个变体评分 0-1。"""
if not variant.image_url:
return 0.0
scoring_prompt = f"""对这张产品图打分(0-100)。
产品:{product.name}
品类:{product.category}
应可见的特征:{'、'.join(product.features)}
评分标准:
- 画面质量和锐度(25 分)
- 产品可见度和焦点(25 分)
- 专业构图(25 分)
- 与产品描述的相关性(25 分)
只返回 0-100 之间的数字。"""
try:
response = await self.client.chat.completions.create(
model="openai/gpt-4o-mini",
messages=[{
"role": "user",
"content": [
{"type": "text", "text": scoring_prompt},
{"type": "image_url", "image_url": {"url": variant.image_url}}
]
}],
max_tokens=10
)
score = float(response.choices[0].message.content.strip()) / 100.0
return min(max(score, 0.0), 1.0)
except Exception:
return 0.5 # 错误时给默认分
async def score_all(
self, results: list[PipelineResult], products: dict[str, Product]
) -> list[PipelineResult]:
"""评分所有变体,每产品选最佳。"""
tasks = []
for result in results:
product = products[result.product_id]
for variant in result.all_variants:
if variant.image_url:
tasks.append(self._score_and_assign(variant, product))
await asyncio.gather(*tasks)
for result in results:
scored = [v for v in result.all_variants if v.score is not None]
if scored:
result.best_variant = max(scored, key=lambda v: v.score)
return results
async def _score_and_assign(self, variant: ImageVariant, product: Product):
variant.score = await self.score_variant(variant, product)
步骤五:完整流水线
import json
async def run_batch_pipeline(
products: list[Product],
api_key: str,
model: str = "bytedance/seedream/5.0/pro/fast",
max_concurrent: int = 10,
):
"""运行完整批量图像生成流水线。"""
print(f"启动批量流水线: {len(products)} 产品 × 5 变体 = "
f"{len(products) * 5} 张图")
print(f"模型: {model} | 并发: {max_concurrent}")
print("=" * 60)
# 阶段一:生成提示词
prompt_gen = PromptGenerator()
prompts_per_product = {
p.id: prompt_gen.generate_prompts(p) for p in products
}
print(f"阶段一: 生成 {sum(len(v) for v in prompts_per_product.values())} 条提示词")
# 阶段二:生成图像
generator = BatchImageGenerator(
api_key=api_key, model=model, max_concurrent=max_concurrent
)
results = await generator.generate_batch(products, prompts_per_product)
print(f"阶段二: 图像生成完成")
generator.report()
# 阶段三:评分选优
scorer = ImageScorer(api_key=api_key)
products_dict = {p.id: p for p in products}
results = await scorer.score_all(results, products_dict)
selected_count = sum(1 for r in results if r.best_variant)
total_cost = sum(r.total_cost for r in results)
print(f"\n阶段三: 评分完成")
print(f"已选图产品: {selected_count}/{len(products)}")
print(f"流水线总成本: ${total_cost:.2f}")
# 保存结果
output = [{
"product_id": r.product_id,
"best_image": r.best_variant.image_url if r.best_variant else None,
"best_score": r.best_variant.score if r.best_variant else None,
"cost": r.total_cost
} for r in results]
with open("pipeline_results.json", "w") as f:
json.dump(output, f, indent=2, ensure_ascii=False)
return results
# 使用示例
if __name__ == "__main__":
products = [
Product(
id=f"prod_{i:03d}",
name=f"高端无线耳机 型号{i}",
category="数码",
features=["主动降噪", "30h 续航", "IPX5 防水"],
style_preference="现代极简"
)
for i in range(100)
]
asyncio.run(run_batch_pipeline(
products=products,
api_key="your-sandbase-api-key",
model="bytedance/seedream/5.0/pro/fast",
max_concurrent=10
))
成本估算
按模型选择(100 产品 × 5 变体 = 500 张)
| 模型 | 单张 | 500 张 | + 评分(500 次) | 总计 |
|---|---|---|---|---|
| Nano Banana Lite | $0.008 | $4.00 | ~$1.50 | $5.50 |
| Nano Banana 2 Lite | $0.01 | $5.00 | ~$1.50 | $6.50 |
| Seedream Fast | $0.015 | $7.50 | ~$1.50 | $9.00 |
| Qwen-Image-3 | $0.03 | $15.00 | ~$1.50 | $16.50 |
| Seedream Pro | $0.04 | $20.00 | ~$1.50 | $21.50 |
时间估算(串行 vs 并行)
| 模型 | 串行 500 张 | 并行 10 并发 |
|---|---|---|
| Nano Banana Lite | ~17 分钟 | ~2 分钟 |
| Seedream Fast | ~25 分钟 | ~3 分钟 |
| Qwen-Image-3 | ~58 分钟 | ~6 分钟 |
| Seedream Pro | ~83 分钟 | ~9 分钟 |
10 并发下,整条 500 张图流水线可在 2–9 分钟内完成。
进阶:混合模型策略
探索用 Fast,定稿用 Pro:
async def hybrid_pipeline(products: list[Product], api_key: str):
"""两阶段流水线:Fast 探索 + Pro 精修。"""
# 阶段 A: Fast 生成全部 500 张 ($7.50)
fast_results = await run_batch_pipeline(
products=products, api_key=api_key,
model="bytedance/seedream/5.0/pro/fast"
)
# 阶段 B: Pro 重新生成最佳提示词 ($4.00,100 张)
pro_gen = BatchImageGenerator(
api_key=api_key, model="bytedance/seedream/5.0/pro", max_concurrent=5
)
for result in fast_results:
if result.best_variant:
pro_variant = await pro_gen.generate_variant(
result.product_id, 99, result.best_variant.prompt
)
if pro_variant.image_url:
result.best_variant = pro_variant
# 总成本: $7.50 (Fast) + $4.00 (Pro 定稿) = $11.50
# 全 Pro: $21.50 → 节省 46%,最终输出仍为 Pro 画质
return fast_results
监控与可观测性
class PipelineMetrics:
"""追踪流水线性能指标。"""
def __init__(self):
self.start_time = time.time()
self.images_generated = 0
self.images_failed = 0
self.total_cost = 0.0
self.latencies: list[int] = []
def record_success(self, latency_ms: int, cost: float):
self.images_generated += 1
self.total_cost += cost
self.latencies.append(latency_ms)
def summary(self) -> dict:
elapsed = time.time() - self.start_time
return {
"耗时秒": round(elapsed, 1),
"已生成": self.images_generated,
"已失败": self.images_failed,
"成功率": f"{self.images_generated / max(self.images_generated + self.images_failed, 1) * 100:.1f}%",
"总成本": f"${self.total_cost:.2f}",
"平均延迟ms": round(sum(self.latencies) / max(len(self.latencies), 1)),
"P95延迟ms": sorted(self.latencies)[int(len(self.latencies) * 0.95)] if self.latencies else 0,
"图/秒": round(self.images_generated / max(elapsed, 1), 2),
}
总结
Agent 批量图像生成遵循一致的模式:大量生成、自动评分、留存最佳。关键工程决策:
- 模型选择 — Seedream Fast 在批量操作中提供最佳质量/成本比
- 并发度 — 10 并发是好的默认值;根据限流情况调整
- 错误处理 — 指数退避重试 + 持续失败进死信队列
- 成本控制 — 追踪每产品支出,设预算上限,用混合策略
完整流水线处理 100 个产品(500 张图),10 分钟内完成,成本 $9–$21(视画质需求)。通过提高并发或分批运行,可线性扩展到数千产品。


