Agent 批量图像生成流水线

手把手教程:构建 Agent 批量图像生成流水线——从产品描述到变体生成、自动评分、最佳选取,含并行处理、错误重试和成本追踪的完整 Python 代码。

结论先行 — 本教程构建一条完整的批量图像生成流水线:Agent 接收 100 个产品描述,每个生成 5 个图像变体(共 500 张),用视觉模型评分,选出最佳。完整 Python 代码基于 OpenAI SDK + SandBase,含并行执行、错误处理、重试逻辑和成本追踪。预计成本:$7.50–$20.00(视模型选择)。

我们要构建什么

一条面向电商的生产级图像生成流水线:

  1. 接收产品描述作为输入(名称、品类、核心卖点)
  2. 每个产品用不同提示词/角度生成 5 个图像变体
  3. 用视觉模型对每个变体评分(画质 + 相关性)
  4. 为每个产品选出最佳变体
  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 批量图像生成遵循一致的模式:大量生成、自动评分、留存最佳。关键工程决策:

  1. 模型选择 — Seedream Fast 在批量操作中提供最佳质量/成本比
  2. 并发度 — 10 并发是好的默认值;根据限流情况调整
  3. 错误处理 — 指数退避重试 + 持续失败进死信队列
  4. 成本控制 — 追踪每产品支出,设预算上限,用混合策略

完整流水线处理 100 个产品(500 张图),10 分钟内完成,成本 $9–$21(视画质需求)。通过提高并发或分批运行,可线性扩展到数千产品。