2026-06-08 23:08:57 +08:00
|
|
|
|
---
|
2026-06-09 23:15:17 +08:00
|
|
|
|
tags: [流式响应, SSE, streaming, API, 多Provider]
|
|
|
|
|
|
create time: 2026-06-09 22:15
|
2026-06-08 23:08:57 +08:00
|
|
|
|
---
|
|
|
|
|
|
|
2026-06-09 23:15:17 +08:00
|
|
|
|
# 流式响应机制 - Claude Code 打字机效果原理
|
|
|
|
|
|
|
|
|
|
|
|
## 概述
|
|
|
|
|
|
|
|
|
|
|
|
Claude Code 通过 SSE(Server-Sent Events)实现流式响应,逐 token 接收 AI 输出,实时呈现"打字机"效果。本文解析事件状态机、内容块交织、错误处理和多 Provider 适配的完整实现。
|
|
|
|
|
|
|
|
|
|
|
|
## 正文
|
|
|
|
|
|
|
|
|
|
|
|
### 为什么需要流式
|
2026-06-08 23:08:57 +08:00
|
|
|
|
|
|
|
|
|
|
想象 AI 需要 30 秒才能生成完整回答——如果等 30 秒后才一次性显示,用户体验是灾难性的。
|
|
|
|
|
|
|
2026-06-09 23:15:17 +08:00
|
|
|
|
流式响应让用户**实时看到 AI 的思考过程**:
|
2026-06-08 23:08:57 +08:00
|
|
|
|
- 文字逐字出现,用户能提前判断方向是否正确
|
|
|
|
|
|
- 工具调用的参数在生成过程中就能预览
|
|
|
|
|
|
- 长时间任务不会让用户觉得"卡死了"
|
|
|
|
|
|
|
2026-06-09 23:15:17 +08:00
|
|
|
|
### BetaRawMessageStreamEvent 核心事件类型
|
2026-06-08 23:08:57 +08:00
|
|
|
|
|
2026-06-09 23:15:17 +08:00
|
|
|
|
流式 API 返回的是一系列 `BetaRawMessageStreamEvent`,每种事件类型对应流式响应的不同阶段(`src/services/api/claude.ts`):
|
2026-06-08 23:08:57 +08:00
|
|
|
|
|
2026-06-09 23:15:17 +08:00
|
|
|
|
```mermaid
|
|
|
|
|
|
flowchart TD
|
|
|
|
|
|
A["message_start 消息开始"] --> B["content_block_start 内容块开始"]
|
|
|
|
|
|
B --> C["content_block_delta 增量数据"]
|
|
|
|
|
|
C --> C
|
|
|
|
|
|
C --> D["content_block_stop 内容块结束"]
|
|
|
|
|
|
D --> B
|
|
|
|
|
|
D --> E["message_delta stop_reason + 最终usage"]
|
|
|
|
|
|
E --> F["message_stop 消息结束"]
|
2026-06-08 23:08:57 +08:00
|
|
|
|
```
|
|
|
|
|
|
|
2026-06-09 23:15:17 +08:00
|
|
|
|
#### 事件处理状态机
|
2026-06-08 23:08:57 +08:00
|
|
|
|
|
2026-06-09 23:15:17 +08:00
|
|
|
|
`src/services/api/claude.ts` 中 `queryModelWithStreaming()` 函数的事件处理循环实现了一个基于 `switch(part.type)` 的状态机:
|
2026-06-08 23:08:57 +08:00
|
|
|
|
|
|
|
|
|
|
| 事件类型 | 处理逻辑 | 状态变更 |
|
|
|
|
|
|
|----------|----------|----------|
|
|
|
|
|
|
| `message_start` | 初始化 `partialMessage`,记录 TTFT(首字节延迟) | `usage` 初始化 |
|
|
|
|
|
|
| `content_block_start` | 按 `part.index` 创建对应类型的内容块 | `contentBlocks[index]` 初始化 |
|
|
|
|
|
|
| `content_block_delta` | 按子类型增量追加数据 | text / thinking / input 累加 |
|
|
|
|
|
|
| `content_block_stop` | 构建完整 `AssistantMessage` 并 yield | 消息推入 `newMessages` |
|
|
|
|
|
|
| `message_delta` | 更新 stop_reason 和最终 usage | 写回最后一条消息 |
|
2026-06-09 23:15:17 +08:00
|
|
|
|
| `message_stop` | 无操作(流结束标记) | -- |
|
2026-06-08 23:08:57 +08:00
|
|
|
|
|
2026-06-09 23:15:17 +08:00
|
|
|
|
#### 内容块类型及其增量数据
|
2026-06-08 23:08:57 +08:00
|
|
|
|
|
2026-06-09 23:15:17 +08:00
|
|
|
|
`content_block_start` 中的 `content_block.type` 决定了如何处理后续 delta:
|
2026-06-08 23:08:57 +08:00
|
|
|
|
|
|
|
|
|
|
| 内容块类型 | Delta 类型 | 累加逻辑 |
|
|
|
|
|
|
|-----------|-----------|----------|
|
|
|
|
|
|
| `text` | `text_delta` | `text += delta.text` |
|
|
|
|
|
|
| `thinking` | `thinking_delta` + `signature_delta` | `thinking += delta.thinking`,`signature = delta.signature` |
|
|
|
|
|
|
| `tool_use` | `input_json_delta` | `input += delta.partial_json`(JSON 字符串增量拼接) |
|
|
|
|
|
|
| `server_tool_use` | `input_json_delta` | 同 tool_use |
|
|
|
|
|
|
| `connector_text` | `connector_text_delta` | 特殊连接器文本(feature flag 控制) |
|
|
|
|
|
|
|
2026-06-09 23:15:17 +08:00
|
|
|
|
> [!tip] 设计细节
|
|
|
|
|
|
> `content_block_start` 时所有文本字段初始化为空字符串,只通过 `content_block_delta` 累加。这是因为 SDK 有时在 start 和 delta 中重复发送相同文本。
|
2026-06-08 23:08:57 +08:00
|
|
|
|
|
2026-06-09 23:15:17 +08:00
|
|
|
|
### 文本 chunk 和 tool_use block 的交织
|
2026-06-08 23:08:57 +08:00
|
|
|
|
|
2026-06-09 23:15:17 +08:00
|
|
|
|
一次 AI 响应可能包含多个内容块,交替出现:
|
2026-06-08 23:08:57 +08:00
|
|
|
|
|
|
|
|
|
|
```
|
|
|
|
|
|
content_block_start (text, index=0) "我来帮你修复这个 bug。"
|
|
|
|
|
|
content_block_delta (text_delta) "首先..."
|
|
|
|
|
|
content_block_stop (index=0)
|
|
|
|
|
|
content_block_start (tool_use, index=1) { name: "Read", input: "..." }
|
2026-06-09 23:15:17 +08:00
|
|
|
|
content_block_delta (input_json_delta) '{"file_p' -> 'ath":' -> '"src/foo.ts"}'
|
2026-06-08 23:08:57 +08:00
|
|
|
|
content_block_stop (index=1)
|
|
|
|
|
|
content_block_start (text, index=2) "我已经看到了问题所在..."
|
|
|
|
|
|
content_block_stop (index=2)
|
|
|
|
|
|
```
|
|
|
|
|
|
|
|
|
|
|
|
每个 `content_block_stop` 触发一次 `yield`,将完整的 AssistantMessage 推送给消费者。这意味着一个 AI 响应会产生**多条** `AssistantMessage`——文本消息和工具调用消息交替产出。
|
|
|
|
|
|
|
2026-06-09 23:15:17 +08:00
|
|
|
|
`stop_reason` 要等到 `message_delta` 才确定(可能是 `end_turn`、`tool_use`、`max_tokens` 等),所以最后一条消息的 `stop_reason` 是**回写**的:
|
2026-06-08 23:08:57 +08:00
|
|
|
|
|
|
|
|
|
|
```typescript
|
2026-06-09 23:15:17 +08:00
|
|
|
|
// claude.ts -- stop_reason 回写逻辑(直接属性修改,不用对象替换)
|
2026-06-08 23:08:57 +08:00
|
|
|
|
// 因为 transcript 写队列持有 message.message 的引用
|
|
|
|
|
|
const lastMsg = newMessages.at(-1)
|
|
|
|
|
|
if (lastMsg) {
|
|
|
|
|
|
lastMsg.message.usage = usage
|
|
|
|
|
|
lastMsg.message.stop_reason = stopReason
|
|
|
|
|
|
}
|
|
|
|
|
|
```
|
|
|
|
|
|
|
2026-06-09 23:15:17 +08:00
|
|
|
|
### 流式中的错误处理
|
2026-06-08 23:08:57 +08:00
|
|
|
|
|
2026-06-09 23:15:17 +08:00
|
|
|
|
#### 网络断开
|
2026-06-08 23:08:57 +08:00
|
|
|
|
|
2026-06-09 23:15:17 +08:00
|
|
|
|
流式连接依赖 SSE。当连接中断时,系统有三层检测机制:
|
2026-06-08 23:08:57 +08:00
|
|
|
|
|
2026-06-09 23:15:17 +08:00
|
|
|
|
1. **被动停滞检测**: 当下一个事件到达时,计算与上一个事件的时间间隔。超过阈值(30 秒,`STALL_THRESHOLD_MS = 30_000`)记录为一次 stall,累积计数并写入遥测日志。
|
|
|
|
|
|
2. **主动空闲超时看门狗**: 使用 `setTimeout` 设置 90 秒(可通过 `CLAUDE_STREAM_IDLE_TIMEOUT_MS` 环境变量覆盖)的硬性超时。如果在此期间没有收到任何事件,主动终止流并抛出错误进入重试流程。
|
|
|
|
|
|
3. **非流式降级**: 作为最后手段,设置 `didFallBackToNonStreaming` 标志,通过 `executeNonStreamingRequest()` 回退到非流式请求。
|
2026-06-08 23:08:57 +08:00
|
|
|
|
|
|
|
|
|
|
```typescript
|
2026-06-09 23:15:17 +08:00
|
|
|
|
// claude.ts -- 被动停滞检测
|
2026-06-08 23:08:57 +08:00
|
|
|
|
const STALL_THRESHOLD_MS = 30_000 // 30 秒无事件视为停滞
|
|
|
|
|
|
|
2026-06-09 23:15:17 +08:00
|
|
|
|
// claude.ts -- 主动空闲超时
|
2026-06-08 23:08:57 +08:00
|
|
|
|
const STREAM_IDLE_TIMEOUT_MS =
|
|
|
|
|
|
parseInt(process.env.CLAUDE_STREAM_IDLE_TIMEOUT_MS || '', 10) || 90_000
|
|
|
|
|
|
```
|
|
|
|
|
|
|
2026-06-09 23:15:17 +08:00
|
|
|
|
#### API 限流
|
2026-06-08 23:08:57 +08:00
|
|
|
|
|
2026-06-09 23:15:17 +08:00
|
|
|
|
当 API 返回限流错误时,系统使用 `withRetry` 包装器进行指数退避重试。重试逻辑考虑了:
|
2026-06-08 23:08:57 +08:00
|
|
|
|
- 错误类型(429 限流 vs 500 服务器错误)
|
|
|
|
|
|
- 重试次数上限
|
|
|
|
|
|
- 退避间隔
|
|
|
|
|
|
|
2026-06-09 23:15:17 +08:00
|
|
|
|
#### Token 超限
|
2026-06-08 23:08:57 +08:00
|
|
|
|
|
2026-06-09 23:15:17 +08:00
|
|
|
|
两种 token 超限场景有不同的处理:
|
2026-06-08 23:08:57 +08:00
|
|
|
|
|
|
|
|
|
|
| 场景 | stop_reason | 处理方式 |
|
|
|
|
|
|
|------|------------|----------|
|
|
|
|
|
|
| **输出超限** | `max_tokens` | 生成错误消息,建议设置 `CLAUDE_CODE_MAX_OUTPUT_TOKENS` |
|
|
|
|
|
|
| **上下文窗口超限** | `model_context_window_exceeded` | 触发 compaction 压缩对话历史后重试 |
|
|
|
|
|
|
|
|
|
|
|
|
```typescript
|
2026-06-09 23:15:17 +08:00
|
|
|
|
// claude.ts -- stop_reason 处理
|
2026-06-08 23:08:57 +08:00
|
|
|
|
if (stopReason === 'max_tokens') {
|
|
|
|
|
|
yield createAssistantAPIErrorMessage({ error: 'max_output_tokens', ... })
|
|
|
|
|
|
}
|
|
|
|
|
|
if (stopReason === 'model_context_window_exceeded') {
|
|
|
|
|
|
// 复用 max_output_tokens 的恢复路径
|
|
|
|
|
|
yield createAssistantAPIErrorMessage({ error: 'max_output_tokens', ... })
|
|
|
|
|
|
}
|
|
|
|
|
|
```
|
|
|
|
|
|
|
2026-06-09 23:15:17 +08:00
|
|
|
|
### 工具执行的流式反馈
|
2026-06-08 23:08:57 +08:00
|
|
|
|
|
2026-06-09 23:15:17 +08:00
|
|
|
|
BashTool 的命令执行也是流式的——通过 `onProgress` 回调逐行推送输出:
|
2026-06-08 23:08:57 +08:00
|
|
|
|
|
2026-06-09 23:15:17 +08:00
|
|
|
|
```mermaid
|
|
|
|
|
|
flowchart TD
|
|
|
|
|
|
A["BashTool.call"] --> B["runShellCommand AsyncGenerator"]
|
|
|
|
|
|
B --> C["每秒轮询输出文件"]
|
|
|
|
|
|
C --> D["onProgress lastLines, allLines"]
|
|
|
|
|
|
D --> E["yield progress, output, fullOutput"]
|
|
|
|
|
|
B --> F["return code, stdout, interrupted"]
|
2026-06-08 23:08:57 +08:00
|
|
|
|
```
|
|
|
|
|
|
|
|
|
|
|
|
UI 层通过 `useToolCallProgress` hook 实时展示命令输出,而不是等命令完全结束。长时间运行的命令还支持自动后台化(`shouldAutoBackground`)。
|
|
|
|
|
|
|
2026-06-09 23:15:17 +08:00
|
|
|
|
### 多 Provider 适配
|
2026-06-08 23:08:57 +08:00
|
|
|
|
|
|
|
|
|
|
| Provider | 流式协议 | 特殊处理 |
|
|
|
|
|
|
|----------|----------|----------|
|
2026-06-09 23:15:17 +08:00
|
|
|
|
| **firstParty**(Anthropic Direct) | 原生 SSE | 延迟最低,TTFT 最快 |
|
2026-06-08 23:08:57 +08:00
|
|
|
|
| **AWS Bedrock** | AWS SDK 流式接口 | 需要额外的 beta header 和认证 |
|
2026-06-09 23:15:17 +08:00
|
|
|
|
| **Google Vertex** | gRPC -> 事件流 | 通过 `getMergedBetas()` 适配 |
|
2026-06-08 23:08:57 +08:00
|
|
|
|
| **foundry** | Anthropic 兼容 API | 内部部署 |
|
|
|
|
|
|
| **openai** | OpenAI 流式适配器 | 转换为 Anthropic 内部格式 |
|
|
|
|
|
|
| **gemini** | Gemini 流式适配器 | 转换为 Anthropic 内部格式 |
|
2026-06-09 23:15:17 +08:00
|
|
|
|
| **grok**(xAI) | Grok 流式适配器 | 转换为 Anthropic 内部格式 |
|
2026-06-08 23:08:57 +08:00
|
|
|
|
|
|
|
|
|
|
所有 Provider 通过统一的 `Stream<BetaRawMessageStreamEvent>` 抽象层屏蔽差异。上层代码(QueryEngine、REPL)不需要关心底层用的是哪个 Provider。
|
|
|
|
|
|
|
2026-06-09 23:15:17 +08:00
|
|
|
|
> [!info] Provider 选择
|
|
|
|
|
|
> `src/utils/model/providers.ts` 中的 `getAPIProvider()` 根据配置决定使用哪个 Provider。每个 Provider 需要适配认证方式、beta header、请求参数格式、错误码映射——但这些差异在 `claude.ts` 的 `queryStream()` 函数中被统一处理。
|
2026-06-08 23:08:57 +08:00
|
|
|
|
|
2026-06-09 23:15:17 +08:00
|
|
|
|
## 关联笔记
|
2026-06-08 23:08:57 +08:00
|
|
|
|
|
2026-06-09 23:15:17 +08:00
|
|
|
|
- [[the-loop]] - Agentic Loop 核心机制
|
|
|
|
|
|
- [[multi-turn]] - 多轮对话管理与 QueryEngine
|
|
|
|
|
|
- [[../context/token-budget]] - Token 预算与输出限制
|