用Spring Boot搭建统一AI流式网关:SSE事件、模型路由、取消与审计
文章摘要
企业同时接入DeepSeek、OpenAI、Qwen或本地模型后,如果每个业务系统直接对接Provider,很快会出现流式协议不一致、错误码分散、取消无效、Token统计缺失和模型切换困难。本文实现一个轻量级Spring Boot AI流式网关,统一请求协议、SSE事件、模型路由、任务状态、取消接口、日志审计与Provider适配,并说明如何避免网关成为新的单点瓶颈。
一、为什么需要统一流式网关
直接接入多个模型时,业务系统需要分别处理:
不同请求参数
不同流式事件格式
不同错误结构
不同Token统计
不同取消方式
不同模型名称
不同鉴权
最终形成:
客服系统 → Provider A
售前系统 → Provider B
知识库系统 → Provider C
代码助手 → 本地模型
模型升级或切换时,需要修改多个系统。
统一网关:
业务应用
→ Enterprise AI Gateway
→ Provider Adapter
→ 模型服务
网关负责:
- 统一请求;
- 模型路由;
- SSE事件;
- 身份与配额;
- 取消;
- 错误转换;
- Token与成本;
- Trace;
- 降级。
二、项目结构
ai-streaming-gateway
├── controller
│ └── AiStreamController.java
├── domain
│ ├── AiStreamRequest.java
│ ├── AiStreamEvent.java
│ └── AiTaskStatus.java
├── provider
│ ├── AiProvider.java
│ ├── SpringAiProvider.java
│ └── ProviderRegistry.java
├── routing
│ └── ModelRouter.java
├── task
│ ├── AiTask.java
│ └── AiTaskService.java
├── security
│ └── AiAccessService.java
└── observability
└── AiMetrics.java
三、统一请求对象
public record AiStreamRequest(
String taskType,
String message,
String conversationId,
String preferredModel,
Map metadata
) {
public AiStreamRequest {
if (
message == null
|| message.isBlank()
) {
throw new IllegalArgumentException(
"message不能为空"
);
}
}
}
业务系统只提交业务意图,不直接传递Provider密钥和底层完整参数。
四、统一事件协议
public record AiStreamEvent(
String taskId,
long sequence,
String type,
Object data,
long timestamp
) {
}
事件类型:
public final class AiEventTypes {
public static final String START = "start";
public static final String DELTA = "delta";
public static final String TOOL_START = "tool_start";
public static final String TOOL_RESULT = "tool_result";
public static final String USAGE = "usage";
public static final String ERROR = "error";
public static final String DONE = "done";
private AiEventTypes() {
}
}
Provider的原始事件不能直接透传给前端,否则业务系统仍然耦合Provider。
五、Provider抽象
public interface AiProvider {
String name();
boolean supports(String model);
Flux stream(
ProviderRequest request,
ProviderCallContext context
);
Mono cancel(String providerRequestId);
}
Provider Chunk:
public record ProviderChunk(
String providerRequestId,
String type,
String content,
Usage usage,
Map metadata
) {
}
六、Spring AI适配器
@Component
public class SpringAiProvider
implements AiProvider {
private final Map clients;
public SpringAiProvider(
@Qualifier("fastChatClient")
ChatClient fastClient,
@Qualifier("powerfulChatClient")
ChatClient powerfulClient
) {
this.clients = Map.of(
"fast",
fastClient,
"powerful",
powerfulClient
);
}
@Override
public String name() {
return "spring-ai";
}
@Override
public boolean supports(String model) {
return clients.containsKey(model);
}
@Override
public Flux stream(
ProviderRequest request,
ProviderCallContext context
) {
ChatClient client = clients.get(
request.model()
);
if (client == null) {
return Flux.error(
new IllegalArgumentException(
"不支持模型:"
+ request.model()
)
);
}
return client.prompt()
.advisors(spec -> spec
.param(
"requestId",
context.requestId()
)
.param(
"tenantId",
context.tenantId()
)
)
.user(request.message())
.stream()
.content()
.map(content ->
new ProviderChunk(
null,
"delta",
content,
null,
Map.of()
)
);
}
@Override
public Mono cancel(
String providerRequestId
) {
return Mono.empty();
}
}
如果底层Provider支持显式取消,适配器应保存providerRequestId并调用真实取消API。
七、模型路由
@Component
public class ModelRouter {
public ModelRoute route(
AiStreamRequest request,
UserQuota quota
) {
if (
request.preferredModel() != null
&& quota.allowedModels()
.contains(
request.preferredModel()
)
) {
return new ModelRoute(
"spring-ai",
request.preferredModel()
);
}
if (
"complex_analysis".equals(
request.taskType()
)
) {
return new ModelRoute(
"spring-ai",
"powerful"
);
}
return new ModelRoute(
"spring-ai",
"fast"
);
}
}
路由条件可以包括:
- 任务复杂度;
- 用户套餐;
- 成本预算;
- 数据敏感性;
- 延迟目标;
- Provider可用性;
- 当前限流状态。
八、任务状态
public enum AiTaskStatus {
CREATED,
RUNNING,
CANCEL_REQUESTED,
CANCELLED,
COMPLETED,
FAILED,
TIMED_OUT
}
任务:
public class AiTask {
private final String taskId;
private final String tenantId;
private final String userId;
private volatile AiTaskStatus status;
private volatile String providerRequestId;
private volatile Disposable subscription;
// 构造方法和状态变更方法省略
}
生产环境应持久化核心任务元数据,不能只保存在单实例Map中。
九、任务服务
@Service
public class AiTaskService {
private final ConcurrentMap
tasks = new ConcurrentHashMap();
public AiTask create(
String tenantId,
String userId
) {
String taskId = UUID.randomUUID()
.toString();
AiTask task = new AiTask(
taskId,
tenantId,
userId
);
tasks.put(taskId, task);
return task;
}
public Mono cancel(
String taskId,
String userId
) {
AiTask task = requireTask(taskId);
task.checkOwner(userId);
task.requestCancel();
Disposable subscription =
task.subscription();
if (subscription != null) {
subscription.dispose();
}
return Mono.empty();
}
}
实际取消还要调用Provider Adapter。
十、Controller
@RestController
@RequestMapping("/api/ai")
public class AiStreamController {
private final GatewayStreamingService service;
private final AiTaskService taskService;
@PostMapping(
value = "/stream",
produces = MediaType.TEXT_EVENT_STREAM_VALUE
)
public Flux>
stream(
@AuthenticationPrincipal
AuthenticatedUser user,
@Valid @RequestBody
AiStreamRequest request
) {
return service.stream(user, request)
.map(event ->
ServerSentEvent
.builder()
.id(
Long.toString(
event.sequence()
)
)
.event(event.type())
.data(event)
.build()
);
}
@PostMapping("/tasks/{taskId}/cancel")
public Mono cancel(
@AuthenticationPrincipal
AuthenticatedUser user,
@PathVariable String taskId
) {
return taskService.cancel(
taskId,
user.userId()
);
}
}
十一、组装流式事件
@Service
public class GatewayStreamingService {
private final AtomicLong globalSequence =
new AtomicLong();
public Flux stream(
AuthenticatedUser user,
AiStreamRequest request
) {
AiTask task = taskService.create(
user.tenantId(),
user.userId()
);
ModelRoute route = router.route(
request,
quotaService.getQuota(user)
);
AiProvider provider =
providerRegistry.require(
route.provider()
);
Flux content =
provider.stream(
toProviderRequest(
request,
route
),
toContext(user, task)
)
.map(chunk ->
toGatewayEvent(
task.taskId(),
chunk
)
);
AiStreamEvent start = event(
task.taskId(),
"start",
Map.of(
"model",
route.model()
)
);
AiStreamEvent done = event(
task.taskId(),
"done",
Map.of()
);
return Flux.concat(
Mono.just(start),
content,
Mono.just(done)
)
.doOnSubscribe(subscription ->
task.markRunning()
)
.doOnComplete(task::markCompleted)
.doOnCancel(task::markCancelled)
.doOnError(task::markFailed)
.doFinally(signal ->
metrics.recordFinished(
task,
signal
)
);
}
}
十二、统一错误事件
不要把SDK堆栈返回浏览器。
.onErrorResume(error -> {
AiGatewayError gatewayError =
errorMapper.map(error);
return Flux.just(
event(
task.taskId(),
"error",
Map.of(
"code",
gatewayError.code(),
"message",
gatewayError.userMessage()
)
)
);
})
错误码:
MODEL_TIMEOUT
RATE_LIMITED
QUOTA_EXCEEDED
MODEL_UNAVAILABLE
INVALID_REQUEST
CONTENT_BLOCKED
STREAM_INTERRUPTED
十三、不要把错误事件和HTTP错误混淆
流建立之前的错误:
认证失败
参数错误
没有配额
应返回标准HTTP错误。
流建立之后的错误:
Provider中断
工具失败
输出解析失败
通过SSE error事件返回,然后结束连接。
十四、Token与成本
每次任务记录:
input_tokens
cached_input_tokens
output_tokens
audio_tokens
model
provider
estimated_cost
actual_cost
如果Provider只在流结束时返回Usage,需要在usage事件中补发。
成本不应由前端计算。
十五、权限与配额
网关必须在调用前检查:
用户是否允许使用AI
租户是否开通模型
任务类型是否允许
月度配额
并发配额
单次Token上限
不要仅依赖Provider的总账户上限。
十六、审计
记录:
request_id
task_id
user_id
tenant_id
task_type
model
prompt_version
start_time
first_token_ms
duration_ms
status
cancelled
error_code
token_usage
Prompt正文和用户数据应按隐私要求脱敏或只保存Hash。
十七、网关如何避免成为单点
建议:
- 服务无状态;
- 任务状态外置;
- 多实例部署;
- 不在本地保存长期会话;
- Provider连接可重建;
- 使用统一Trace;
- 限制每实例长连接数;
- 慢客户端保护;
- 健康检查和熔断。
SSE连接可分配到任意实例,但独立取消接口需要通过共享任务存储找到对应任务和Provider请求。
十八、何时不需要统一网关
小型项目只有:
- 一个业务;
- 一个模型;
- 少量用户;
- 无配额;
- 无审计;
- 无模型切换;
可以先直接使用Spring AI。
出现以下情况后,网关价值明显:
多个业务
多个Provider
多租户
统一计费
动态路由
统一安全
统一评测
总结
统一AI流式网关应提供:
统一请求
+统一事件
+模型路由
+任务取消
+错误转换
+成本与审计
它的目标不是再包一层HTTP,而是把不同模型的流式能力转换为企业内部稳定、可治理的AI服务协议。
延伸阅读
如果你正在关注企业级 AI 应用、Agent、RAG、MCP 与大模型工程化落地,欢迎访问 智元界:
https://www.zyentor.com/
智元界将持续分享可运行的技术实战、架构设计、问题排查与企业应用案例。