用Spring Boot搭建统一AI流式网关:SSE事件、模型路由、取消与审计

发布时间:2026/7/27 18:22:13

用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三、统一请求对象publicrecordAiStreamRequest(StringtaskType,Stringmessage,StringconversationId,StringpreferredModel,MapString,Objectmetadata){publicAiStreamRequest{if(messagenull||message.isBlank()){thrownewIllegalArgumentException(message不能为空);}}}业务系统只提交业务意图不直接传递Provider密钥和底层完整参数。四、统一事件协议publicrecordAiStreamEvent(StringtaskId,longsequence,Stringtype,Objectdata,longtimestamp){}事件类型publicfinalclassAiEventTypes{publicstaticfinalStringSTARTstart;publicstaticfinalStringDELTAdelta;publicstaticfinalStringTOOL_STARTtool_start;publicstaticfinalStringTOOL_RESULTtool_result;publicstaticfinalStringUSAGEusage;publicstaticfinalStringERRORerror;publicstaticfinalStringDONEdone;privateAiEventTypes(){}}Provider的原始事件不能直接透传给前端否则业务系统仍然耦合Provider。五、Provider抽象publicinterfaceAiProvider{Stringname();booleansupports(Stringmodel);FluxProviderChunkstream(ProviderRequestrequest,ProviderCallContextcontext);MonoVoidcancel(StringproviderRequestId);}Provider ChunkpublicrecordProviderChunk(StringproviderRequestId,Stringtype,Stringcontent,Usageusage,MapString,Objectmetadata){}六、Spring AI适配器ComponentpublicclassSpringAiProviderimplementsAiProvider{privatefinalMapString,ChatClientclients;publicSpringAiProvider(Qualifier(fastChatClient)ChatClientfastClient,Qualifier(powerfulChatClient)ChatClientpowerfulClient){this.clientsMap.of(fast,fastClient,powerful,powerfulClient);}OverridepublicStringname(){returnspring-ai;}Overridepublicbooleansupports(Stringmodel){returnclients.containsKey(model);}OverridepublicFluxProviderChunkstream(ProviderRequestrequest,ProviderCallContextcontext){ChatClientclientclients.get(request.model());if(clientnull){returnFlux.error(newIllegalArgumentException(不支持模型request.model()));}returnclient.prompt().advisors(spec-spec.param(requestId,context.requestId()).param(tenantId,context.tenantId())).user(request.message()).stream().content().map(content-newProviderChunk(null,delta,content,null,Map.of()));}OverridepublicMonoVoidcancel(StringproviderRequestId){returnMono.empty();}}如果底层Provider支持显式取消适配器应保存providerRequestId并调用真实取消API。七、模型路由ComponentpublicclassModelRouter{publicModelRouteroute(AiStreamRequestrequest,UserQuotaquota){if(request.preferredModel()!nullquota.allowedModels().contains(request.preferredModel())){returnnewModelRoute(spring-ai,request.preferredModel());}if(complex_analysis.equals(request.taskType())){returnnewModelRoute(spring-ai,powerful);}returnnewModelRoute(spring-ai,fast);}}路由条件可以包括任务复杂度用户套餐成本预算数据敏感性延迟目标Provider可用性当前限流状态。八、任务状态publicenumAiTaskStatus{CREATED,RUNNING,CANCEL_REQUESTED,CANCELLED,COMPLETED,FAILED,TIMED_OUT}任务publicclassAiTask{privatefinalStringtaskId;privatefinalStringtenantId;privatefinalStringuserId;privatevolatileAiTaskStatusstatus;privatevolatileStringproviderRequestId;privatevolatileDisposablesubscription;// 构造方法和状态变更方法省略}生产环境应持久化核心任务元数据不能只保存在单实例Map中。九、任务服务ServicepublicclassAiTaskService{privatefinalConcurrentMapString,AiTasktasksnewConcurrentHashMap();publicAiTaskcreate(StringtenantId,StringuserId){StringtaskIdUUID.randomUUID().toString();AiTasktasknewAiTask(taskId,tenantId,userId);tasks.put(taskId,task);returntask;}publicMonoVoidcancel(StringtaskId,StringuserId){AiTasktaskrequireTask(taskId);task.checkOwner(userId);task.requestCancel();Disposablesubscriptiontask.subscription();if(subscription!null){subscription.dispose();}returnMono.empty();}}实际取消还要调用Provider Adapter。十、ControllerRestControllerRequestMapping(/api/ai)publicclassAiStreamController{privatefinalGatewayStreamingServiceservice;privatefinalAiTaskServicetaskService;PostMapping(value/stream,producesMediaType.TEXT_EVENT_STREAM_VALUE)publicFluxServerSentEventAiStreamEventstream(AuthenticationPrincipalAuthenticatedUseruser,ValidRequestBodyAiStreamRequestrequest){returnservice.stream(user,request).map(event-ServerSentEvent.AiStreamEventbuilder().id(Long.toString(event.sequence())).event(event.type()).data(event).build());}PostMapping(/tasks/{taskId}/cancel)publicMonoVoidcancel(AuthenticationPrincipalAuthenticatedUseruser,PathVariableStringtaskId){returntaskService.cancel(taskId,user.userId());}}十一、组装流式事件ServicepublicclassGatewayStreamingService{privatefinalAtomicLongglobalSequencenewAtomicLong();publicFluxAiStreamEventstream(AuthenticatedUseruser,AiStreamRequestrequest){AiTasktasktaskService.create(user.tenantId(),user.userId());ModelRouterouterouter.route(request,quotaService.getQuota(user));AiProviderproviderproviderRegistry.require(route.provider());FluxAiStreamEventcontentprovider.stream(toProviderRequest(request,route),toContext(user,task)).map(chunk-toGatewayEvent(task.taskId(),chunk));AiStreamEventstartevent(task.taskId(),start,Map.of(model,route.model()));AiStreamEventdoneevent(task.taskId(),done,Map.of());returnFlux.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-{AiGatewayErrorgatewayErrorerrorMapper.map(error);returnFlux.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中断 工具失败 输出解析失败通过SSEerror事件返回然后结束连接。十四、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_usagePrompt正文和用户数据应按隐私要求脱敏或只保存Hash。十七、网关如何避免成为单点建议服务无状态任务状态外置多实例部署不在本地保存长期会话Provider连接可重建使用统一Trace限制每实例长连接数慢客户端保护健康检查和熔断。SSE连接可分配到任意实例但独立取消接口需要通过共享任务存储找到对应任务和Provider请求。十八、何时不需要统一网关小型项目只有一个业务一个模型少量用户无配额无审计无模型切换可以先直接使用Spring AI。出现以下情况后网关价值明显多个业务 多个Provider 多租户 统一计费 动态路由 统一安全 统一评测总结统一AI流式网关应提供统一请求 统一事件 模型路由 任务取消 错误转换 成本与审计它的目标不是再包一层HTTP而是把不同模型的流式能力转换为企业内部稳定、可治理的AI服务协议。

相关新闻