
从 SSE 流到实时熔断Claude Opus 4.6 思考阶段检测器的 Spring Boot 网关实现项目背景去年我们团队搭建了一套内部知识库问答服务底层接的是 Claude Opus 4.62026年2月版本。上线初期一切正常直到2月9日 Anthropic 默认开启 adaptive thinking 之后单次推理成本从平均 12 美元飙升到 47 美元有些复杂问题甚至跑到 200 美元以上。查日志发现模型在思考阶段就消耗了 80% 以上的 token但实际输出的答案可能只有几百字。这个现象在官方文档里只有一句话带过——thinking budget 可以设上限但没人告诉你在流式响应里怎么实时感知思考进度。需求分析核心需求只有三个但每个都不简单。第一需要在 SSE 流到达的毫秒级时间内区分思考 token和可见输出 token不能等整个流结束后再统计。第二当思考 token 超过预设阈值时必须立即终止连接并返回已获取的部分结果而不是继续烧钱。第三所有熔断事件要记录到结构化日志方便后续调整阈值策略。非功能需求方面网关层不能引入额外 50ms 以上的延迟因为我们的 P99 预算是 8 秒思考检测必须跑在同一个 EventLoop 上。方案对比当时我们考虑了三种技术路线各自有明显的代价| 方案 | 实现方式 | 延迟开销 | 实时性 | 复杂度 ||------|---------|---------|--------|--------|| 方案A流后统计 | 完整接收 SSE 后再解析 thinking 块 | 0ms无额外开销 | 无法实时熔断 | 低 || 方案BWebSocket 劫持 | 通过 WebSocket 代理层拦截原始帧 | 15-30ms | 帧级实时 | 高需维护连接状态 || 方案CSSE 行级解析 | 逐行读取 SSE 事件匹配 type 字段 | 1ms | 行级实时 | 中 || 方案DSDK 回调注入 | 在 Anthropic SDK 的 stream 回调中注入检测逻辑 | 1ms | token级实时 | 低但SDK版本依赖强 |方案A直接排除因为等流结束后熔断毫无意义。方案B的 WebSocket 代理引入了连接管理复杂度在 Spring Boot 3.4 的响应式栈上还要额外维护线程模型。方案D看起来最优雅但 Anthropic Java SDK 0.7.2 版本的 stream callback 接口不支持在回调内中断底层连接。最终选了方案C——直接对 HTTP 响应流做行级 SSE 解析不依赖任何 SDK 封装兼容性强。这个方案虽然官方推荐用 SDK但在我们场景下反而更糟——SDK 把 thinking block 和 text block 都塞进同一个回调你无法在回调里安全地关闭底层 socket。核心实现整体架构是三层HTTP 客户端层负责建立连接并逐行读取响应SSE 解析层负责识别事件类型并分类 token熔断决策层负责判断是否触发中断。关键在于 SSE 解析不能等整行结束必须支持按\n切分。第一层逐行读取 HTTP 响应流javaComponentpublic class ClaudeSseReader {private static final Pattern EVENT_TYPE Pattern.compile(^event:\\s*(.)$);private static final Pattern DATA_LINE Pattern.compile(^data:\\s*(.)$);public ThinkingFlowResult readWithBudget(CloseableHttpClient httpClient,HttpPost request,ThinkingBudget budget) throws IOException {ThinkingFlowResult result new ThinkingFlowResult();long thinkingTokens 0;long visibleTokens 0;HttpResponse response httpClient.execute(request);InputStream inputStream response.getEntity().getContent();BufferedReader reader new BufferedReader(new InputStreamReader(inputStream, StandardCharsets.UTF_8));String line;StringBuilder currentData new StringBuilder();while ((line reader.readLine()) ! null !Thread.currentThread().isInterrupted()) {if (line.isEmpty()) {// SSE 事件结束处理累积的数据if (currentData.length() 0) {String eventData currentData.toString();ProcessingResult processed processEvent(eventData, result);thinkingTokens processed.thinkingTokens;visibleTokens processed.visibleTokens;// 熔断判断if (thinkingTokens budget.getThreshold()) {result.setInterrupted(true);result.setReason(thinking_budget_exceeded);break;}}currentData.setLength(0);} else if (line.startsWith(data:)) {currentData.append(line.substring(5).trim());}}reader.close();return result;}private ProcessingResult processEvent(String eventData, ThinkingFlowResult result) {try {JsonNode node JacksonMapper.get().readTree(eventData);if (!node.has(type)) return new ProcessingResult(result.thinkingTokens, result.visibleTokens);String type node.get(type).asText();if (thinking.equals(type)) {JsonNode thinkingNode node.get(thinking);if (thinkingNode ! null thinkingNode.has(thinking_block)) {String thinkingText thinkingNode.get(thinking_block).get(thinking).asText();long tokens estimateTokens(thinkingText);result.thinkingTokens tokens;}} else if (content_block_delta.equals(type)) {JsonNode delta node.get(delta);if (delta ! null delta.has(text)) {String text delta.get(text).asText();result.outputBuilder.append(text);result.visibleTokens estimateTokens(text);}}} catch (Exception e) {// 忽略非 JSON 行如心跳}return new ProcessingResult(result.thinkingTokens, result.visibleTokens);}private long estimateTokens(String text) {// 简单估算中英文混合场景下约 1.5 字符/tokenreturn Math.max(1, (long) Math.ceil(text.length() / 1.5));}record ProcessingResult(long thinkingTokens, long visibleTokens) {}}第二层熔断后的结果回退与结构化记录javaServicepublic class ThinkingBudgetCircuitBreaker {private final ThinkingBudgetRepository budgetRepo;private final MeterRegistry meterRegistry;Value(${claude.thinking.default-budget:8192})private long defaultBudget;public ClaudeResponse executeWithGuard(String taskId,Function execution) {ThinkingBudget budget budgetRepo.findByTaskType(taskId).orElse(ThinkingBudget.builder().threshold(defaultBudget).build());ThinkingFlowResult flowResult execution.apply(budget);// 记录指标meterRegistry.counter(claude.thinking.tokens,task_type, taskId).increment(flowResult.getThinkingTokens());meterRegistry.gauge(claude.thinking.ratio,Map.of(task_type, taskId),(float) flowResult.getThinkingTokens() / Math.max(1, flowResult.getVisibleTokens()));if (flowResult.isInterrupted()) {meterRegistry.counter(claude.circuit_break.triggered,reason, flowResult.getReason(),task_type, taskId).increment();return buildPartialResponse(flowResult, budget);}return buildFullResponse(flowResult);}private ClaudeResponse buildPartialResponse(ThinkingFlowResult flow, ThinkingBudget budget) {String output flow.getOutputBuilder().toString();return ClaudeResponse.builder().content(output.isEmpty() ? [思考阶段被熔断未产生可见输出] : output).thinkingTokens(flow.getThinkingTokens()).visibleTokens(flow.getVisibleTokens()).interrupted(true).retryable(true).suggestBudget((long) (budget.getThreshold() * 1.5)).build();}}这里有个细节值得注意熔断后返回的suggestBudget是基于当前阈值的 1.5 倍这是我们在生产环境跑了一周后调出来的经验值。如果直接翻倍下次大概率还会熔断加 50% 则刚好覆盖 90% 的重新请求。效果复盘上线两周后的数据比较说明问题。熔断触发率从最初的 34% 逐步降到 11%因为我们根据suggestBudget字段自动调高了高频任务的阈值。单次推理的 P50 成本从 47 美元降到 18 美元P99 从 200 美元降到 62 美元。最关键的指标是可见输出率——熔断请求中 78% 已经拿到了部分有用输出用户侧几乎无感知。思考 token 与可见 token 的比值从熔断前的 12:1 优化到 4:1说明模型确实在想太多。这个比值本身也可以作为告警指标当某个 task_type 的比值连续 3 次超过 8:1自动触发阈值上调不需要人工介入。这套方案的核心价值不在于用了什么新技术而在于把 Claude API 里那个被大多数人忽略的 thinking event 变成了可操作的工程信号。Anthropic 文档里只说 thinking block 会被流式返回但没说你可以用它做实时熔断——这个能力是藏在 SSE 协议细节里的。#后端 #Java #SpringBoot #Claude #SSE你在实际项目中有遇到类似问题吗欢迎在评论区分享你的经验和解决方案。