用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/

智元界将持续分享可运行的技术实战、架构设计、问题排查与企业应用案例。