同时接入图像 + 视频双模型生成:从异步回调到文件管理的完整避坑指南
一、写在前面:这几个坑你踩过吗?
在做 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 工程实践相关内容。
AtomGit 是由开放原子开源基金会联合 CSDN 等生态伙伴共同推出的新一代开源与人工智能协作平台。平台坚持“开放、中立、公益”的理念,把代码托管、模型共享、数据集托管、智能体开发体验和算力服务整合在一起,为开发者提供从开发、训练到部署的一站式体验。
更多推荐



所有评论(0)