一、写在前面:这几个坑你踩过吗?

在做 AI 生成类应用时,相信很多人都遇到过这几个问题:

  • 图像模型是同步返回,视频模型是异步的,两套处理逻辑完全不同,代码越写越乱
  • 视频生成完成后,拿到的是一个有效期只有几小时的临时 URL,稍不注意文件就丢了
  • 供应商升级接口字段后,整套轮询逻辑要跟着重写
  • 想同时支持多个图像模型(比如 Flux 和 DALL·E),却不得不维护两套独立的适配代码

这些问题的根本原因,是每个模型供应商都有自己的一套 API 风格。本文介绍一种统一任务层的接入方式,把图像和视频生成的调用逻辑收敛到一套标准协议里,从根本上减少维护成本。

二、环境准备

2.1 核心思路

本文使用的方案,是通过一个统一任务协议(Task API)来屏蔽底层模型的差异。这种设计模式并不局限于某一个具体的 SDK,核心是:

  • 用同一套请求结构提交图像 / 视频任务
  • 用同一套轮询或回调机制获取结果
  • 用统一的文件托管方案管理生成物

具体选型可以根据项目情况决定,本文以 Crun.ai API 为例演示。

2.2 安装依赖

pip install requests python-dotenv

2.3 配置环境变量

# .env
AI_API_KEY=your_api_key_here
AI_BASE_URL=https://api.crun.ai/v1

加载配置:

import os
from dotenv import load_dotenv

load_dotenv()

API_KEY = os.getenv("AI_API_KEY")
BASE_URL = os.getenv("AI_BASE_URL")

HEADERS = {
    "Authorization": f"Bearer {API_KEY}",
    "Content-Type": "application/json"
}

三、图像生成:同步调用完整示例

图像生成任务通常在几秒内完成,采用同步返回。

3.1 提交图像生成任务

import requests

def generate_image(prompt: str, model: str = "flux", size: str = "1024x1024") -> dict:
    """
    提交图像生成任务
    :param prompt: 图像描述文本
    :param model: 使用的模型,如 flux / dall-e-3
    :param size: 图像尺寸
    :return: 任务结果字典
    """
    url = f"{BASE_URL}/tasks/image"

    payload = {
        "model": model,
        "prompt": prompt,
        "size": size,
        "response_format": "url"
    }

    response = requests.post(url, headers=HEADERS, json=payload, timeout=30)
    response.raise_for_status()

    return response.json()

3.2 解析返回结果

result = generate_image("一只在雪地里奔跑的柴犬,摄影风格,高清")

if result.get("status") == "succeeded":
    image_url = result["output"]["url"]
    print(f"图像生成成功,URL:{image_url}")
else:
    print(f"生成失败:{result.get('error')}")

返回结构示例:

{
  "task_id": "img_abc123",
  "status": "succeeded",
  "output": {
    "url": "https://storage.example.com/outputs/img_abc123.png",
    "width": 1024,
    "height": 1024
  },
  "usage": {
    "model": "flux",
    "cost": 0.004
  }
}

⚠️ 注意response_format 若选 b64_json,返回的是 Base64 字符串,需要 decode 后再存储。生产环境建议直接使用 url 格式,由服务商统一托管,避免自行管理临时文件。

四、视频生成:异步任务完整处理流程

视频生成耗时较长(通常 30 秒~3 分钟),必须走异步流程。以下提供主动轮询Webhook 回调两种方案。

4.1 方案一:轮询模式(适合小规模/测试环境)

import time

def generate_video_with_polling(
    prompt: str,
    model: str = "kling",
    duration: int = 5,
    poll_interval: int = 5,
    max_wait: int = 300
) -> dict:
    """
    提交视频生成任务并轮询等待结果
    :param prompt: 视频描述文本
    :param model: 视频模型,如 kling / luma
    :param duration: 视频时长(秒)
    :param poll_interval: 轮询间隔(秒)
    :param max_wait: 最大等待时间(秒)
    :return: 最终任务结果
    """
    # Step 1: 提交任务
    submit_url = f"{BASE_URL}/tasks/video"
    payload = {
        "model": model,
        "prompt": prompt,
        "duration": duration,
        "aspect_ratio": "16:9"
    }

    response = requests.post(submit_url, headers=HEADERS, json=payload, timeout=30)
    response.raise_for_status()
    task = response.json()
    task_id = task.get("task_id")
    print(f"任务已提交,task_id:{task_id}")

    # Step 2: 轮询任务状态
    status_url = f"{BASE_URL}/tasks/{task_id}"
    elapsed = 0

    while elapsed < max_wait:
        time.sleep(poll_interval)
        elapsed += poll_interval

        status_response = requests.get(status_url, headers=HEADERS, timeout=10)
        status_response.raise_for_status()
        status_data = status_response.json()

        current_status = status_data.get("status")
        print(f"[{elapsed}s] 当前状态:{current_status}")

        if current_status == "succeeded":
            return status_data
        elif current_status == "failed":
            raise RuntimeError(f"任务失败:{status_data.get('error')}")

    raise TimeoutError(f"任务超时,已等待 {max_wait} 秒")

调用示例:

result = generate_video_with_polling(
    prompt="夕阳下,一匹白马在草原上慢跑,电影感镜头",
    model="kling",
    duration=5
)
print(f"视频生成完成:{result['output']['url']}")

⚠️ 避坑:轮询间隔建议不低于 5 秒,避免触发频率限制。务必设置 max_wait 上限,防止死循环挂起进程。

4.2 方案二:Webhook 回调(推荐生产环境)

轮询会持续占用服务器资源,生产环境更推荐 Webhook 模式:任务完成后由服务端主动通知,业务侧不轮询。

提交任务时附带 Webhook 地址:

def generate_video_with_webhook(
    prompt: str,
    model: str = "kling",
    webhook_url: str = "https://your-server.com/api/webhook/callback"
) -> str:
    submit_url = f"{BASE_URL}/tasks/video"
    payload = {
        "model": model,
        "prompt": prompt,
        "duration": 5,
        "webhook": {
            "url": webhook_url,
            "events": ["succeeded", "failed"]
        }
    }

    response = requests.post(submit_url, headers=HEADERS, json=payload, timeout=30)
    response.raise_for_status()
    return response.json()["task_id"]

服务端接收回调(Flask 示例):

from flask import Flask, request, jsonify

app = Flask(__name__)

@app.route("/api/webhook/callback", methods=["POST"])
def handle_callback():
    data = request.get_json()
    task_id = data.get("task_id")
    status = data.get("status")

    if status == "succeeded":
        video_url = data["output"]["url"]
        save_to_db(task_id, video_url)      # 写库更新业务状态
        print(f"任务 {task_id} 完成:{video_url}")

    elif status == "failed":
        error_msg = data.get("error", "未知错误")
        mark_task_failed(task_id, error_msg)

    # 必须在 5 秒内返回 200,否则会触发重试
    return jsonify({"received": True}), 200

⚠️ 避坑:Webhook 接口必须在 5 秒内返回 200,建议先响应后异步处理业务逻辑(可结合消息队列)。本地调试时可用 ngrok 临时暴露端口。

五、文件管理:解决临时 URL 过期问题

直接对接原生模型 API 时,返回的文件 URL 通常是临时地址,有效期从几小时到 24 小时不等。对用户来说,URL 过期等同于文件丢失。

常见的两种处理策略:

策略一:使用服务商统一托管(推荐) 部分统一任务层 API 会把生成结果存储在自己的对象存储中,返回长期有效的 URL,无需自行转存。这是最省事的方式。

策略二:主动下载转存到自有 OSS

import requests

def download_and_save(source_url: str, save_path: str) -> None:
    """
    将生成的文件下载并保存到本地或转存
    :param source_url: 来源文件 URL
    :param save_path: 本地保存路径
    """
    response = requests.get(source_url, stream=True, timeout=60)
    response.raise_for_status()

    with open(save_path, "wb") as f:
        for chunk in response.iter_content(chunk_size=8192):
            f.write(chunk)

    print(f"文件已保存:{save_path}")

# 使用示例
download_and_save(
    source_url="https://storage.example.com/outputs/video_xyz.mp4",
    save_path="./outputs/result.mp4"
)

建议:无论采用哪种策略,都应该在任务完成回调触发后立即处理文件,不要依赖"稍后再存"。生产环境里"稍后"往往等于"丢失"。

六、生产级重试策略封装

网络抖动、模型偶发超时,在生产环境不可避免。下面封装一个带指数退避的重试装饰器,适用于所有 AI API 调用:

import functools
import time
import random
import requests

def retry_with_backoff(max_retries: int = 3, base_delay: float = 1.0, max_delay: float = 30.0):
    """
    指数退避重试装饰器
    """
    def decorator(func):
        @functools.wraps(func)
        def wrapper(*args, **kwargs):
            last_exception = None
            for attempt in range(max_retries + 1):
                try:
                    return func(*args, **kwargs)
                except requests.exceptions.HTTPError as e:
                    # 4xx 错误(参数问题)不重试
                    if e.response.status_code < 500:
                        raise
                    last_exception = e
                except (requests.exceptions.ConnectionError,
                        requests.exceptions.Timeout) as e:
                    last_exception = e

                if attempt < max_retries:
                    # 指数退避 + 随机抖动,避免惊群效应
                    delay = min(
                        base_delay * (2 ** attempt) + random.uniform(0, 1),
                        max_delay
                    )
                    print(f"第 {attempt + 1} 次重试,等待 {delay:.1f}s...")
                    time.sleep(delay)

            raise last_exception
        return wrapper
    return decorator


@retry_with_backoff(max_retries=3)
def generate_image_safe(prompt: str, model: str = "flux") -> dict:
    return generate_image(prompt, model)

七、常见报错与解决方案

报错信息 原因 解决方法
401 Unauthorized API Key 无效或格式错误 检查 Authorization Header,确认 Key 有效
422 Unprocessable Entity 参数格式不符合要求 检查 payload 字段,注意 duration 通常只支持固定值
429 Too Many Requests 请求频率超限 降低并发数,重试时增加等待时间
task status: failed 模型生成失败 检查 Prompt 是否触发内容过滤,适当修改描述
Webhook 未收到回调 服务端地址不可公网访问 确认 URL 可公网访问;本地调试使用 ngrok

八、完整调用流程图

用户请求
    │
    ▼
提交 Task(POST /tasks/image 或 /tasks/video)
    │
    ├── 图像生成 ──► 同步返回(status: succeeded)
    │                        │
    │                        ▼
    │               获取 output.url(建议长期托管)
    │
    └── 视频生成 ──► 返回 task_id(status: pending)
                             │
              ┌──────────────┴──────────────┐
              ▼                             ▼
         轮询模式                      Webhook 模式
    GET /tasks/{task_id}          服务端接收 POST 回调
         每 5 秒查询一次             解析 status + output
              │                             │
              └──────────────┬──────────────┘
                             ▼
                    获取 output.url → 写库更新业务状态

九、总结

本文围绕图像 + 视频双模型接入的工程实践,覆盖了以下核心问题:

  • ✅ 统一任务协议接入多模型,减少各自维护的适配代码
  • ✅ 同步图像生成的调用与返回处理
  • ✅ 异步视频生成的两种方案:轮询 vs Webhook 对比与选型建议
  • ✅ 文件存储的两种策略及注意事项
  • ✅ 生产级重试封装(指数退避 + 随机抖动)
  • ✅ 高频报错原因与解决方案汇总

这套代码结构在实际项目中经过验证,可以直接作为你自己项目的参考基础。有问题欢迎在评论区交流。


如果这篇文章对你有帮助,欢迎点赞收藏,后续会持续更新 AI 工程实践相关内容。

Logo

AtomGit 是由开放原子开源基金会联合 CSDN 等生态伙伴共同推出的新一代开源与人工智能协作平台。平台坚持“开放、中立、公益”的理念,把代码托管、模型共享、数据集托管、智能体开发体验和算力服务整合在一起,为开发者提供从开发、训练到部署的一站式体验。

更多推荐