Spring Boot微服务接入万亿参数MoE大模型的API网关与流控设计

微服务架构下大模型接入的核心挑战

万亿参数MoE大模型(如Kimi K3)开源后,企业面临将模型能力嵌入现有微服务体系的需求。Spring Boot微服务接入大模型不是简单的HTTP调用封装,需要解决三个核心问题:推理接口的高延迟与超时治理、Token计费与限流策略、MoE模型多轮工具调用的会话管理。本文基于Spring Boot 3.3 + Spring Cloud 2024技术栈,给出从API网关到推理服务的全链路设计方案。

API网关层的路由与协议转换

Spring Cloud Gateway作为微服务的统一入口,负责将业务请求路由到后端推理服务,同时处理SSE流式响应的透传:

// GatewayConfig.java - 推理路由配置
@Configuration
public class GatewayConfig {

    @Bean
    public RouteLocator agentRoutes(RouteLocatorBuilder builder) {
        return builder.routes()
            .route("agent-inference", r -> r
                .path("/api/agent/**")
                .filters(f -> f
                    .filter(tokenCountFilter)
                    .filter(rateLimitFilter)
                    .modifyRequestBody(String.class, (exchange, body) -> {
                        // 注入请求上下文:tenant_id, user_id, trace_id
                        JsonNode json = objectMapper.readTree(body);
                        ((ObjectNode) json).put("tenant_id",
                            exchange.getRequest().getHeaders().getFirst("X-Tenant-Id"));
                        return Mono.just(objectMapper.writeValueAsString(json));
                    })
                    .timeout(Duration.ofSeconds(120))
                )
                .uri("lb://agent-inference-service")
            )
            .build();
    }
}

网关层的关键配置:超时设置为120秒(MoE模型推理+工具调用链路可能较长),请求体改写注入租户和链路信息。SSE流式响应的透传需要关闭Gateway的响应缓冲,配置spring.cloud.gateway.httpclient.response-timeout=120s并使用ServerHttpResponse.flushBuffer()确保每个数据块即时下发。

Token计费与细粒度限流

MoE模型按Token使用量计费,限流策略不能简单用QPS,需要基于Token配额的令牌桶算法:

// TokenRateLimitFilter.java
@Component
public class TokenRateLimitFilter implements GlobalFilter, Ordered {

    private final RedisTemplate<String, String> redisTemplate;

    @Override
    public Mono<Void> filter(ServerWebExchange exchange, GatewayFilterChain chain) {
        String tenantId = exchange.getRequest().getHeaders().getFirst("X-Tenant-Id");
        String key = "token_quota:" + tenantId;

        // 1. 检查租户Token配额
        String remaining = redisTemplate.opsForValue().get(key);
        if (remaining != null && Long.parseLong(remaining) <= 0) {
            exchange.getResponse().setStatusCode(HttpStatus.TOO_MANY_REQUESTS);
            return exchange.getResponse().writeWith(Mono.just(
                exchange.getResponse().bufferFactory()
                    .wrap("{\"error\":\"Token quota exceeded\"}".getBytes())
            ));
        }

        // 2. 扣减配额在推理完成后异步执行(通过回调钩子)
        exchange.getAttributes().put("token_quota_key", key);

        return chain.filter(exchange);
    }

    /**
     * 推理完成后扣减Token配额(由InferenceCallback触发)
     */
    public void deductTokenQuota(String tenantId, long tokensUsed) {
        String key = "token_quota:" + tenantId;
        redisTemplate.opsForValue().decrement(key, tokensUsed);
    }

    @Override
    public int getOrder() { return -1; }
}

配额扣减采用后付费模式:请求放行时不拦截,推理完成后根据实际Token消耗扣减配额。配额不足时下一个请求被拒绝。这种方式避免了预扣减导致的配额浪费,也解决了MoE模型单次调用Token消耗不确定的问题。

推理服务封装与重试策略

推理服务客户端封装需要处理三种异常场景:模型冷启动(首次请求延迟30s+)、推理超时(MoE模型长文本场景)、服务端过载(503 Service Unavailable):

// InferenceClient.java
@Service
public class InferenceClient {

    private final WebClient webClient;

    public InferenceClient(WebClient.Builder webClientBuilder) {
        this.webClient = webClientBuilder
            .baseUrl("http://agent-inference-service:8000")
            .build();
    }

    public Flux<ServerSentEvent<String>> streamChat(ChatRequest request) {
        return webClient.post()
            .uri("/v1/chat/completions")
            .contentType(MediaType.APPLICATION_JSON)
            .bodyValue(request)
            .retrieve()
            .bodyToFlux(new ParameterizedTypeReference<>() {})
            .timeout(Duration.ofSeconds(120))
            .retryWhen(Retry.backoff(3, Duration.ofSeconds(2))
                .filter(this::isRetryable)
                .doBeforeRetry(signal -> log.warn(
                    "Retry attempt {} for inference call",
                    signal.totalRetries() + 1
                ))
            );
    }

    private boolean isRetryable(Throwable error) {
        if (error instanceof WebClientResponseException e) {
            // 503过载可重试,400/422参数错误不重试
            return e.getStatusCode() == HttpStatus.SERVICE_UNAVAILABLE;
        }
        return error instanceof TimeoutException;
    }
}

重试策略的关键:Retry.backoff(3, Duration.ofSeconds(2))表示最多3次重试,指数退避起始2秒。只对503和超时重试,参数错误(4xx)直接失败。SSE流场景下Flux天然支持背压,不需要额外处理消费端限速。

MoE多轮工具调用的会话管理

MoE模型的Agent模式下,一次用户请求可能触发多轮工具调用。会话状态需要持久化到Redis,支持跨请求恢复:

// AgentSession.java
@Data
@RedisHash(value = "agent_session", timeToLive = 1800)
public class AgentSession {
    @Id
    private String sessionId;
    private String userId;
    private String tenantId;
    private List<ChatMessage> history;
    private List<ToolCall> pendingToolCalls;
    private SessionStatus status;
    private long tokensUsed;
    private Instant createdAt;
    private Instant lastActiveAt;

    public enum SessionStatus {
        ACTIVE, WAITING_TOOL_RESULT, COMPLETED, EXPIRED
    }
}

// SessionService.java
@Service
@RequiredArgsConstructor
public class SessionService {

    private final AgentSessionRepository sessionRepo;

    /**
     * 工具调用完成后恢复会话
     */
    public AgentSession resumeWithToolResult(
            String sessionId, String toolCallId, Object result) {
        AgentSession session = sessionRepo.findById(sessionId)
            .orElseThrow(() -> new SessionNotFoundException(sessionId));

        // 将工具结果追加到历史
        session.getHistory().add(ChatMessage.toolResult(toolCallId, result));
        session.setStatus(SessionStatus.ACTIVE);
        session.setLastActiveAt(Instant.now());

        return sessionRepo.save(session);
    }
}

会话TTL设为30分钟,超时自动清理。状态机:ACTIVE(推理中)到WAITING_TOOL_RESULT(等待工具返回)到ACTIVE(恢复推理)到COMPLETED。工具服务通过MCP协议回调resumeWithToolResult方法,触发下一轮推理。

服务治理与熔断降级

大模型推理服务属于慢速资源,熔断策略需要与普通微服务区别配置。Resilience4j的慢调用比例熔断更适合大模型场景:

# application.yml - Resilience4j配置
resilience4j:
  circuitbreaker:
    instances:
      inference-service:
        sliding-window-type: COUNT_BASED
        sliding-window-size: 20
        slow-call-duration-threshold: 60s
        slow-call-rate-threshold: 80
        failure-rate-threshold: 50
        wait-duration-in-open-state: 30s
        permitted-number-of-calls-in-half-open-state: 3
  timelimiter:
    instances:
      inference-service:
        timeout-duration: 120s

慢调用阈值设为60秒,MoE模型正常推理在30秒内完成,超过60秒判定为慢调用。当20次调用中80%为慢调用时触发熔断。熔断后降级策略返回缓存结果或转投轻量模型(如Qwen2.5-72B),避免业务完全中断。

分布式链路追踪与Token用量统计

Spring Boot 3.3原生支持Micrometer Tracing,配合OpenTelemetry导出到Jaeger。Token用量按模型和租户维度分别计数,接入Grafana看板后可以实时监控每个租户的Token消耗速率和累计用量,为成本分摊提供数据依据:

// TokenUsageAspect.java - Token用量切面
@Aspect
@Component
@Slf4j
public class TokenUsageAspect {

    @AfterReturning(
        pointcut = "execution(* com.example.service.InferenceClient.*(..))",
        returning = "result"
    )
    public void recordTokenUsage(JoinPoint joinPoint, Object result) {
        if (result instanceof InferenceResponse resp) {
            Metrics.counter("agent.token.usage",
                "model", resp.getModel(),
                "tenant", resp.getTenantId()
            ).increment(resp.getUsage().getTotalTokens());
        }
    }
}

Token用量切面自动拦截推理客户端的返回值,提取模型名和租户ID写入Micrometer Counter。Grafana看板按租户、模型、时间段三个维度聚合Token消耗,为费用分摊和容量规划提供精确数据。

原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/springboot-wei-fu-wu-jie-ru-wan-yi-can-shu-moe-da-mo-xing/

(0)
小编小编
上一篇 2小时前
下一篇 2小时前

相关推荐