视频生成管线后端:从“异步等待”到“流式进度推送”的工程重构
上周处理一个内容生成平台的并发瓶颈时,我发现传统的“提交任务-轮询状态”模式在视频生成场景下彻底失效了。当用户请求通过 Sora 或 Veo 2 级模型生成一段 10 秒的高清视频时,前端页面往往因为长时间无响应而让用户误以为系统崩溃。更糟糕的是,当并发量突破 500 QPS 时,数据库的连接池迅速耗尽,大量任务堆积在 Kafka 中导致延迟飙升至分钟级。这迫使我们重新审视后端架构,核心问题不在于模型推理本身,而在于如何高效管理长生命周期任务的上下文与状态同步。
背景:高吞吐视频生成的架构痛点
我们的技术栈基于 Spring Boot 3.2.5 运行在 JDK 17.0.12 环境上,中间件选用 Redis 7.2.5 进行缓存,消息队列使用 RabbitMQ 3.12.0(早期版本)并计划迁移至 Kafka 3.6.0。业务场景要求支持多模态输入(文本+参考图)生成视频,平均生成耗时从 30 秒到 3 分钟不等。
初期设计采用简单的“提交即返回”策略,将任务 ID 存入 Redis,前端通过 WebSocket 心跳查询。这种方案在小流量下尚可,但在大促期间暴露出严重问题:
- 状态一致性差:Redis 中的数据与 RabbitMQ 中的实际处理进度不同步,导致前端显示“处理中”但后端已报错。
- 连接泄露:每个活跃用户维持一个 WebSocket 长连接,线程池资源被大量空闲连接占用。
- 重试风暴:网络抖动时,前端无限重连,后端重复消费消息,产生大量脏数据。
我们需要一种机制,既能保证状态实时同步,又能降低后端资源消耗,同时确保在模型服务降级时,用户依然能获得清晰的反馈。
过程:引入事件溯源与流式进度推送
为了解决上述问题,我决定放弃单纯的状态轮询,转而采用“事件溯源(Event Sourcing)”结合“Server-Sent Events (SSE)”的方案。核心思路是将任务的生命周期拆解为一系列不可变的事件,后端仅负责发布事件,前端订阅并渲染状态。
1. 任务状态机设计
首先,定义严格的状态枚举,避免模糊的“处理中”状态。我们引入了TaskPhase接口,涵盖从接收到最终结果的全流程:
```java
public enum TaskPhase implements Serializable {
PENDING(0, "排队中"),
VALIDATING(10, "参数校验中"),
GENERATING_START(20, "开始生成"),
GENERATING_PROGRESS(50, "生成中"), // 支持分段进度
COMPLETED(100, "已完成"),
FAILED(-1, "失败");
private final int code;
private final String desc;
// getter...
}
```
2. SSE 推送服务实现
相较于 WebSocket 的双向通信,SSE 更适合“服务端主动推送、客户端只读”的场景。它基于 HTTP 长连接,天然支持断线重连,且防火墙兼容性更好。以下是核心的 Controller 实现:
```java
@RestController
@RequestMapping("/api/v1/videos")
@RequiredArgsConstructor
public class VideoGenerationController {
private final VideoService videoService;
private final EventPublisher eventPublisher;
@GetMapping(value = "/status/{taskId}", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public Flux> streamTaskStatus(@PathVariable String taskId) {
return eventPublisher.subscribeToTask(taskId)
.map(event -> ServerSentEvent.builder()
.event("taskUpdate")
.data(new Gson().toJson(event))
.build())
.onErrorResume(e -> Flux.just(
ServerSentEvent.builder()
.event("error")
.data("{\"message\":\"Connection lost\"}")
.build()))
.timeout(Duration.ofMinutes(10)); // 设置超时,防止资源永久占用
}
}
```
3. 异步解耦与背压控制
在 Service 层,调用大模型 API 是耗时操作。我们使用CompletableFuture配合自定义线程池videoGenExecutor(核心线程数 20,最大线程数 100)来异步执行。关键在于,当上游请求过快时,必须实施背压(Backpressure)。我们引入了 Redisson 的分布式锁来限制同一用户每秒最多发起 3 个生成请求,防止单个用户拖垮整个集群。
此外,为了优化模型响应,我们在发送请求前增加了预处理逻辑:利用本地轻量级 OCR 和 NLP 模型(Spring AI 0.8.1)对输入内容进行合规性检查,过滤掉明显违规或低质量的提示词,将无效请求拦截在服务入口,而非浪费 GPU 算力。
```yaml
application.yml 配置片段
spring:
ai:
openai:
api-key: ${OPENAI_API_KEY}
chat:
options:
model: gpt-4o-2024-05-13 # 实际调用时替换为视频生成模型端点
task:
execution:
pool:
core-size: 20
max-size: 100
queue-capacity: 500
```
4. 遭遇的坑:SSE 连接中断与状态丢失
这个方案虽然优雅,但在实际落地时遇到了一个棘手问题:移动网络环境下,SSE 连接极易中断。一旦断开,客户端需要重新建立连接,但如果此时任务正在生成,服务端如何知道该从哪个进度继续推送?
起初,我尝试在 Redis 中存储最新进度,但这导致了“竞态条件”:客户端重连读取旧进度,服务端已推进到新进度,造成状态跳跃。经过排查,发现根本原因是缺乏全局唯一的事务 ID。最终解决方案是引入Redis Stream。我们将每个任务的每一步状态变更都追加到对应的 Stream 中,客户端重连时携带最后收到的ID,服务端通过XREAD命令从该 ID 之后拉取所有未推送的事件。这不仅解决了断线重连的问题,还保留了完整的操作历史,便于后续审计。
效果:性能指标对比
重构后,我们进行了为期一周的压力测试,数据对比如下:
| 指标 | 重构前 (WebSocket 轮询) | 重构后 (SSE + Redis Stream) | 提升幅度 |
| :--- | :--- | :--- | :--- |
| 平均首屏响应时间 | 2.5s | 0.8s | ↓ 68% |
| 服务端内存占用 (峰值) | 4.2 GB | 2.1 GB | ↓ 50% |
| 网络带宽消耗 | 高 (频繁握手) | 低 (HTTP 复用) | ↓ 40% |
| 任务超时失败率 | 12% | 1.5% | ↓ 87.5% |
更重要的是,用户反馈显著改善。之前因“假死”导致的客服投诉下降了 90%。虽然官方推荐 WebSocket 用于实时交互,但在视频生成这种单向数据流场景下,SSE 配合事件溯源确实更稳定、更省资源。
总结
多模态视频生成的后端挑战,本质上是长事务管理与高并发状态的平衡艺术。通过引入 SSE 替代 WebSocket,并结合 Redis Stream 实现精准的状态回溯,我们成功构建了高可用的生成管线。记住,不要盲目追求技术潮流,适合场景的才是最好的。对于后端开发者而言,可观测性与容错机制的设计,往往比单纯的代码实现更能决定系统的生死。
#后端 #Java #SpringBoot #SSE #视频生成
你在实际项目中有遇到类似问题吗?欢迎在评论区分享你的经验和解决方案。