生成式媒体服务通常不是一个模型解决所有问题:图像模型负责产出高质量关键帧,视频模型再把静态画面扩展成连续镜头。使用 AWS vLLM-Omni Deep Learning Container,可以在 Amazon SageMaker AI 上分别部署 FLUX.2-klein 与 Wan2.1-VACE,并根据任务特点选择两种调用方式:图片走实时推理,视频走异步推理,最终从 Amazon S3 获取 MP4。
同一个容器镜像,不等于同一个推理端点
这里的关键是复用同一种 AWS vLLM-Omni DLC,而不是把两个大型生成模型强行塞进同一个进程。更清晰的部署方式是创建两个 SageMaker 模型和端点:
- FLUX.2-klein 端点:处理文本到图像请求,使用实时推理。调用方保持 HTTP 连接,立即取得生成结果。
- Wan2.1-VACE 端点:接收图片与动画提示词,使用异步推理。请求体、状态和结果通过 S3 解耦。
- S3 存储桶:保存 FLUX 生成的 PNG、视频端点的请求 JSON,以及最终 MP4 或包含视频地址的响应文件。
这两条链路的资源特征不同。图片生成适合交互式体验,但仍要控制分辨率、采样步数和并发;视频生成时间更长、输出更大,用异步端点可以避免让客户端长时间占用连接。
部署时应让两个 SageMaker Model 指向同一个 DLC 镜像 URI,但分别配置模型文件、启动参数、实例类型和推理模式。不要因为共享容器镜像,就默认两个模型也必须共享扩缩容策略。
从关键帧到 MP4 的数据流
一条实用的处理流水线可以拆成五步:
- 应用向 FLUX.2-klein 实时端点提交提示词。
- 将返回的 PNG 上传到 S3,形成稳定的中间产物地址。
- 创建 Wan2.1-VACE 请求 JSON,其中包含图片 S3 URI、运动描述和视频参数。
- 调用 SageMaker
InvokeEndpointAsync,轮询输出位置。 - 从异步输出对象中读取 MP4,或者继续解析模型返回的视频 S3 URI。
把中间图片落到 S3 有两个好处:视频任务可以重试而不必重新生成图片,同时也方便记录一次生成任务的输入、关键帧和最终视频。生产环境中可使用任务 ID 组织对象,例如:
s3://media-bucket/jobs/01J.../keyframe.png
s3://media-bucket/jobs/01J.../vace-request.json
s3://media-bucket/jobs/01J.../output.out
s3://media-bucket/jobs/01J.../result.mp4
可改造的端到端 Python 示例
下面的脚本展示完整编排逻辑。运行前需要:
- 已创建一个 FLUX 实时端点和一个 Wan 异步端点;
- 当前 AWS 身份拥有调用两个端点及读写目标 S3 前缀的权限;
- Wan 端点使用的 SageMaker 执行角色能够读取关键帧所在的 S3 对象;
- 本地已安装并配置
boto3。
由于不同模型包的请求与响应字段可能不同,示例假设 FLUX 接收 prompt、width、height,并在 image 字段返回 Base64 PNG;Wan 请求则使用 image_s3_uri、prompt 和 num_frames。实际部署时,只需调整 build_flux_request、decode_flux_response 与 build_vace_request 三个适配函数。
python -m pip install boto3
export AWS_REGION=us-east-1
export FLUX_ENDPOINT=flux2-klein-realtime
export VACE_ENDPOINT=wan21-vace-async
export MEDIA_BUCKET=my-generative-media-bucket
export MEDIA_PREFIX=jobs/demo-001
python generate_media.py
将以下内容保存为 generate_media.py:
import base64
import json
import os
import time
from pathlib import Path
from urllib.parse import urlparse
import boto3
REGION = os.getenv('AWS_REGION', 'us-east-1')
FLUX_ENDPOINT = os.environ['FLUX_ENDPOINT']
VACE_ENDPOINT = os.environ['VACE_ENDPOINT']
BUCKET = os.environ['MEDIA_BUCKET']
PREFIX = os.getenv('MEDIA_PREFIX', 'jobs/demo-001').strip('/')
runtime = boto3.client('sagemaker-runtime', region_name=REGION)
s3 = boto3.client('s3', region_name=REGION)
def split_s3_uri(uri):
parsed = urlparse(uri)
if parsed.scheme != 's3':
raise ValueError(f'Not an S3 URI: {uri}')
return parsed.netloc, parsed.path.lstrip('/')
def build_flux_request(prompt):
# 按实际 FLUX 服务契约修改字段。
return {
'prompt': prompt,
'width': 1024,
'height': 576,
'num_inference_steps': 28,
}
def decode_flux_response(raw):
# 假设响应为 {'image': '<base64>'}。
document = json.loads(raw)
encoded = document['image']
if encoded.startswith('data:'):
encoded = encoded.split(',', 1)[1]
return base64.b64decode(encoded)
def build_vace_request(image_uri, prompt):
# 按实际 Wan2.1-VACE 服务契约修改字段。
return {
'image_s3_uri': image_uri,
'prompt': prompt,
'num_frames': 81,
'fps': 16,
}
def wait_for_async_object(output_uri, failure_uri=None, timeout=1800):
deadline = time.time() + timeout
output_bucket, output_key = split_s3_uri(output_uri)
failure = split_s3_uri(failure_uri) if failure_uri else None
while time.time() < deadline:
try:
response = s3.get_object(Bucket=output_bucket, Key=output_key)
return response['Body'].read()
except s3.exceptions.NoSuchKey:
pass
if failure:
try:
response = s3.get_object(Bucket=failure[0], Key=failure[1])
message = response['Body'].read().decode('utf-8', errors='replace')
raise RuntimeError(f'Asynchronous inference failed: {message}')
except s3.exceptions.NoSuchKey:
pass
time.sleep(10)
raise TimeoutError(f'No async result after {timeout} seconds: {output_uri}')
def download_video(async_body, destination):
# 容器可能直接返回 MP4,也可能返回包含 S3 URI 或 Base64 的 JSON。
if not async_body.lstrip().startswith(b'{'):
Path(destination).write_bytes(async_body)
return
document = json.loads(async_body)
if 'video_s3_uri' in document:
bucket, key = split_s3_uri(document['video_s3_uri'])
s3.download_file(bucket, key, destination)
elif 'video' in document:
encoded = document['video']
if encoded.startswith('data:'):
encoded = encoded.split(',', 1)[1]
Path(destination).write_bytes(base64.b64decode(encoded))
else:
raise ValueError(f'Unknown video response fields: {list(document)}')
image_prompt = 'A compact red research rover crossing an icy moon, cinematic light'
flux_response = runtime.invoke_endpoint(
EndpointName=FLUX_ENDPOINT,
ContentType='application/json',
Accept='application/json',
Body=json.dumps(build_flux_request(image_prompt)).encode('utf-8'),
)
png = decode_flux_response(flux_response['Body'].read())
image_key = f'{PREFIX}/keyframe.png'
s3.put_object(Bucket=BUCKET, Key=image_key, Body=png, ContentType='image/png')
image_uri = f's3://{BUCKET}/{image_key}'
print(f'Keyframe: {image_uri}')
vace_request = build_vace_request(
image_uri,
'The rover moves slowly forward while fine ice particles drift across the frame',
)
request_key = f'{PREFIX}/vace-request.json'
s3.put_object(
Bucket=BUCKET,
Key=request_key,
Body=json.dumps(vace_request).encode('utf-8'),
ContentType='application/json',
)
async_response = runtime.invoke_endpoint_async(
EndpointName=VACE_ENDPOINT,
InputLocation=f's3://{BUCKET}/{request_key}',
ContentType='application/json',
Accept='application/json',
InvocationTimeoutSeconds=1800,
)
print(f"Inference ID: {async_response.get('InferenceId')}")
print(f"Output: {async_response['OutputLocation']}")
result = wait_for_async_object(
async_response['OutputLocation'],
async_response.get('FailureLocation'),
)
download_video(result, 'result.mp4')
print('Saved result.mp4')
如果容器直接返回 video/mp4,可以把异步调用的 Accept 改为 video/mp4,并保留脚本中处理原始字节的分支。如果响应是 JSON,则应根据实际模型服务契约修改字段名。
生产化时不要忽略的边界
权限要尽量收窄。 调用方只需要特定端点的调用权限和指定任务前缀的 S3 读写权限;端点执行角色也不应获得整个存储桶的无条件访问权。启用 S3 加密,并根据数据敏感度配置访问日志与保留周期。
重试不能盲目重复生成。 实时图像生成和异步视频生成都会消耗较多计算资源。为每个业务请求生成幂等任务 ID,重试前先检查关键帧、请求文件和结果对象是否存在。
把超时与业务状态分开。 客户端轮询超时不代表 SageMaker 任务已经失败。应用数据库应保存推理 ID、输出位置和当前状态,由后台任务继续查询。
分别扩缩两个端点。 图片端点关注交互延迟和并发,视频端点更关注队列长度、单任务时间和 GPU 利用率。两者使用同一 DLC,并不意味着应该采用相同实例或自动扩缩指标。
控制存储成本。 中间 PNG、请求 JSON、失败输出和 MP4 会不断累积。可以为 jobs/ 前缀配置 S3 Lifecycle,在满足审计和复现周期后转入低频存储或自动删除。
采用建议
落地时先用固定分辨率、固定帧数和单并发验证整条链路,再分别压测实时端点与异步端点。确认响应格式、S3 权限和失败输出都可观测后,再加入自动扩缩、任务取消、内容审核和生命周期策略。
这套架构最重要的价值不是简单地把两个模型放进云端,而是用实时与异步推理匹配不同负载:FLUX.2-klein 快速交付关键帧,Wan2.1-VACE 在后台完成耗时动画,S3 则成为二者之间可重试、可审计的数据边界。