流式传输不只是性能优化,它改变了用户感知"思考"的方式。
5.1 核心问题
第 4 章搭好了工具引擎。但谁来决定"现在该调用哪个工具"?谁来理解用户输入并生成回复?答案是 LLM——我们需要把用户消息和工具列表发给 Claude,拿回它的回复。
本章要解决的问题是:
如何与 Anthropic API 通信、实时处理流式响应、把 tool_use 块重组为可执行的工具调用,并把 API 成本控制在最低?
这比看起来复杂。流式 API 不是一次性返回完整 JSON,而是把响应拆成几十乃至几百个"事件"逐条推送。tool_use 块的输入参数(本身是 JSON)也是分片流入的——你必须把这些 delta 碎片拼接成完整 JSON 再解析。
Claude Code 在此基础上还叠加了两个性能特性:Prompt Caching(缓存 System Prompt,降本最高 90%)和 Extended Thinking / Ultrathink(让模型生成内部推理 token,提升复杂任务质量)。
5.2 原理讲解
5.2.1 Messages API:消息格式
Anthropic 的 Messages API 是对话的基础协议。每次调用都发送完整的历史消息列表(无状态——服务端不记忆上下文):
// 最小请求结构
const response = await anthropic.messages.create({
model: 'claude-opus-4-5',
max_tokens: 4096,
system: '你是一个 Coding Agent...', // 系统提示
messages: [
{ role: 'user', content: '帮我写一个冒泡排序' },
{ role: 'assistant', content: '好的,我先...' },
{ role: 'user', content: '请加上类型注解' },
// ... 完整历史
],
tools: [/* 工具 JSON Schema 列表,见第4章 */],
})
响应的 stop_reason 有两种关键值:
'end_turn':正常回复结束'tool_use':模型要调用工具(content数组里有tool_use块)
5.2.2 流式 SSE 协议:事件序列
非流式调用等全部 token 生成完才返回,延迟高达数秒。流式调用将响应拆成事件序列,逐条 yield,实现"打字机效果":
服务端推送顺序:
message_start → { type: 'message_start', message: { usage: {...}, ... } }
content_block_start[0] → { type: 'content_block_start', index: 0, content_block: { type: 'text' } }
content_block_delta[0] → { type: 'content_block_delta', index: 0, delta: { type: 'text_delta', text: '好' } }
content_block_delta[0] → { type: 'content_block_delta', index: 0, delta: { type: 'text_delta', text: '的,' } }
...(更多 text_delta)
content_block_stop[0] → { type: 'content_block_stop', index: 0 }
message_delta → { type: 'message_delta', delta: { stop_reason: 'end_turn' }, usage: { output_tokens: 42 } }
message_stop → { type: 'message_stop' }
关键字段来源:
input_tokens(输入 token 数):在message_start事件中output_tokens(输出 token 数):在message_delta事件中cache_read_input_tokens:在message_start中,命中缓存时不为零
5.2.3 tool_use 块的流式重组
当模型决定调用工具时,响应里会出现 tool_use 类型的内容块。关键难点:工具的 input 字段(JSON 对象)是流式的——它通过 input_json_delta 事件一片一片推送:
content_block_start[1] → { type: 'content_block_start', index: 1,
content_block: { type: 'tool_use', id: 'toolu_01', name: 'Bash', input: '' } }
content_block_delta[1] → { delta: { type: 'input_json_delta', partial_json: '{"com' } }
content_block_delta[1] → { delta: { type: 'input_json_delta', partial_json: 'mand' } }
content_block_delta[1] → { delta: { type: 'input_json_delta', partial_json: '": "ls -la"}' } }
content_block_stop[1] → { type: 'content_block_stop', index: 1 }
重组逻辑:在内存中为每个 index 维护一个字符串累积器,content_block_stop 时调用 JSON.parse() 得到完整输入对象。
Claude Code 的实现(src/services/api/claude.ts,content_block_delta 分支):
// 简化版重组逻辑
case 'content_block_delta': {
const block = contentBlocks[part.index]
if (part.delta.type === 'input_json_delta') {
// block.type === 'tool_use',block.input 是累积字符串
block.input += part.delta.partial_json
} else if (part.delta.type === 'text_delta') {
block.text += part.delta.text
} else if (part.delta.type === 'thinking_delta') {
block.thinking += part.delta.thinking // Extended Thinking
}
break
}
// content_block_stop 时:JSON.parse(block.input) → 完整工具参数
5.2.4 Token 计数与成本估算
Anthropic API 按 token 计费。Claude Code 的 src/utils/modelCost.ts 维护了完整的定价表:
// 以 claude-sonnet-4-5 为例($3 input / $15 output per Mtok)
export const COST_TIER_3_15 = {
inputTokens: 3, // $3 / 1M tokens
outputTokens: 15, // $15 / 1M tokens
promptCacheWriteTokens: 3.75, // 写缓存
promptCacheReadTokens: 0.3, // 读缓存(仅为输入价格的 10%)
webSearchRequests: 0.01, // 每次搜索
}
成本计算公式:
$$\text{cost} = \frac{N_{in}}{10^6} \times p_{in} + \frac{N_{out}}{10^6} \times p_{out} + \frac{N_{cache_read}}{10^6} \times p_{cache_read} + \frac{N_{cache_write}}{10^6} \times p_{cache_write}$$
5.2.5 Prompt Caching:缓存断点
这是 Claude Code 成本优化的核心机制。原理:对请求中不变的大块内容(System Prompt、工具列表)加上 cache_control 标记,Anthropic 服务端会把这段 token 缓存起来,后续请求命中缓存时只收取约 10% 的费用。
// 在 system prompt 末尾注入缓存断点
const systemBlocks = [
{
type: 'text',
text: longSystemPrompt, // 可能有数千 token
cache_control: { type: 'ephemeral' } // ← 缓存断点
}
]
// 在工具列表末尾注入缓存断点
const tools = [
...allTools.slice(0, -1),
{
...lastTool,
cache_control: { type: 'ephemeral' } // ← 缓存断点
}
]
两种 TTL:
- 默认 5 分钟(活跃会话通常足够)
- 1 小时(
should1hCacheTTL()返回 true 时:Anthropic 员工 or claude.ai 订阅用户)
getPromptCachingEnabled() 的实现逻辑:
// src/services/api/claude.ts
export function getPromptCachingEnabled(model: string): boolean {
if (isEnvTruthy(process.env.DISABLE_PROMPT_CACHING)) return false
// ... 可按模型独立关闭
return true // 默认开启
}
5.2.6 Extended Thinking / Ultrathink
Extended Thinking 让模型在回复前生成"内心独白"(CoT 推理),对复杂任务有显著提升。
三种配置模式(src/utils/thinking.ts):
export type ThinkingConfig =
| { type: 'adaptive' } // 模型自决是否思考
| { type: 'enabled'; budgetTokens: number } // 固定 budget(越大越深思)
| { type: 'disabled' } // 关闭思考
触发机制:用户消息中含 ultrathink(大小写不敏感,须是完整单词)时自动启用最大 thinking budget:
// src/utils/thinking.ts
export function hasUltrathinkKeyword(text: string): boolean {
return /\bultrathink\b/i.test(text)
}
API 请求加上 thinking 参数后,响应流中会出现 thinking 类型的内容块,通过 thinking_delta 事件流式传输(内容对用户可见,类似"思考过程展示")。
5.3 Claude Code 源码细节
5.3.1 queryModelWithStreaming:主入口
// src/services/api/claude.ts(精简)
export async function* queryModelWithStreaming({
messages,
systemPrompt,
thinkingConfig,
tools,
signal,
options,
}: {
messages: Message[]
systemPrompt: SystemPrompt
thinkingConfig: ThinkingConfig
tools: Tools
signal: AbortSignal // 取消信号(用户按 Ctrl+C 时触发)
options: Options
}): AsyncGenerator<StreamEvent | AssistantMessage | SystemAPIErrorMessage, void>
这是一个 async generator——调用者用 for await (const event of queryModelWithStreaming(...)) 逐事件消费,不需要等待整个响应完成。
5.3.2 流式事件状态机
Claude Code 内部维护一个 contentBlocks 数组(按 index 索引),在流式处理过程中逐步填充:
message_start → 初始化 usage(input_tokens 在此确定)
content_block_start → contentBlocks[index] = 新建对应类型的块
content_block_delta → 累积到 contentBlocks[index](text/input_json/thinking)
content_block_stop → tool_use 块:JSON.parse(input 字符串) → 完整对象
message_delta → 更新 usage(output_tokens 在此确定)、stop_reason
message_stop → 构造完整 AssistantMessage yield 出去
5.3.3 getCacheControl():动态决定缓存 TTL
// src/services/api/claude.ts
export function getCacheControl({ scope, querySource } = {}): {
type: 'ephemeral'
ttl?: '1h'
scope?: CacheScope
} {
return {
type: 'ephemeral',
...(should1hCacheTTL(querySource) && { ttl: '1h' }),
...(scope === 'global' && { scope }),
}
}
TTL '1h' 只在满足以下所有条件时启用:Anthropic 内部员工 或 claude.ai 订阅用户 + 未超额 + GrowthBook 功能开关允许。这体现了"按用户等级差异化缓存"的精细化成本管理。
5.3.4 calculateUSDCost():成本计算
// src/utils/modelCost.ts
export function calculateUSDCost(resolvedModel: string, usage: Usage): number {
const modelCosts = getModelCosts(resolvedModel, usage)
return (
(usage.input_tokens / 1_000_000) * modelCosts.inputTokens +
(usage.output_tokens / 1_000_000) * modelCosts.outputTokens +
((usage.cache_read_input_tokens ?? 0) / 1_000_000) * modelCosts.promptCacheReadTokens +
((usage.cache_creation_input_tokens ?? 0) / 1_000_000) * modelCosts.promptCacheWriteTokens
)
}
注意 cache_read_input_tokens 字段:第一次请求时为 0(缓存未命中,走写入);后续命中缓存的请求中此值激增,对应成本大幅下降。
5.4 最小化产出物
代码骨架位于
../chapters/05/src/,参考实现位于../chapters/05/solution/。 前置条件:需要ANTHROPIC_API_KEY环境变量(验收脚本不调用 API,但npm start需要)。
本章要实现什么
在 ../chapters/05/src/llm.ts 中完成 LLM 客户端。
接口规范(已提供,不要修改):
export const MODEL: string
export const MODEL_COSTS: { inputTokens: number; outputTokens: number; promptCacheWriteTokens: number; promptCacheReadTokens: number }
export type ThinkingConfig = { type: 'enabled'; budget_tokens: number } | { type: 'disabled' }
export type ToolUse = { id: string; name: string; input: Record<string, unknown> }
export type Usage = { input_tokens: number; output_tokens: number; cache_creation_input_tokens?: number; cache_read_input_tokens?: number }
export function calculateCostUSD(usage: Usage): number
export function hasUltrathinkKeyword(text: string): boolean
export async function streamQuery(params: {
systemPrompt: string
messages: Anthropic.MessageParam[]
thinkingConfig: ThinkingConfig
tools?: Anthropic.Tool[]
enablePromptCaching: boolean
onText?: (text: string) => void
}): Promise<{ text: string; thinking?: string; toolUses: ToolUse[]; stopReason: string; usage: Usage }>
你需要实现:
calculateCostUSD():按 MODEL_COSTS 计算美元成本(含 cache_creation 和 cache_read)hasUltrathinkKeyword():检测\bultrathink\b(完整单词,大小写不敏感)streamQuery():- 创建 Anthropic client,组装 systemBlocks(可选注入
cache_control: { type: 'ephemeral' }) - 维护
contentBlocks状态机,处理message_start/content_block_start/content_block_delta/message_delta事件 text_delta→ 累积 text,调用onText回调input_json_delta→ 累积 tool_use 的 input 字符串- 最终重组:text、thinking、toolUses(
JSON.parseinput 字符串)
- 创建 Anthropic client,组装 systemBlocks(可选注入
关键约束:
streamQuery必须处理tool_use块的流式重组(input_json_delta累积 →JSON.parse)onText回调在每个text_delta时调用,不是最后一次性调用
验收
cd docs/chapters/05
npm install
npm test
卡住时查看 ../chapters/05/solution/llm.ts。
5.5 流式事件处理完整生命周期图
API 请求发出
│
▼
┌──────────────────┐
│ message_start │ ← 获取 input_tokens、cache_read/creation_tokens
└────────┬─────────┘
│
▼ 对每个内容块循环:
┌──────────────────────────────────────────────┐
│ content_block_start[i] │ ← 按类型初始化槽位
│ type=text → contentBlocks[i] = {text:''} │
│ type=tool_use → contentBlocks[i] = {input:''} (JSON 将分片到达)
│ type=thinking → contentBlocks[i] = {thinking:''}
└────────┬─────────────────────────────────────┘
│
▼ 多次:
┌──────────────────────────────────────────────┐
│ content_block_delta[i] │
│ text_delta → block.text += delta.text │ ← 实时打印
│ input_json_delta → block.input += partial │ ← 累积 JSON 碎片
│ thinking_delta → block.thinking += delta │
└────────┬─────────────────────────────────────┘
│
▼
┌──────────────────┐
│ content_block_stop[i] │ ← tool_use: JSON.parse(block.input) → 完整对象
└────────┬─────────┘
│
▼ (所有块完成后)
┌──────────────────┐
│ message_delta │ ← 获取 output_tokens、stop_reason
└────────┬─────────┘
│
▼
┌──────────────────┐
│ message_stop │ ← 流结束
└──────────────────┘
5.6 本章小结
本章构建了与 Claude API 通信的完整客户端:
-
流式 SSE 协议:响应被拆分为
message_start→content_block_*→message_delta→message_stop的事件序列,需要维护状态机来重组完整内容。 -
tool_use 块重组:工具参数(JSON)通过
input_json_delta分片推送,必须逐片累积、最终JSON.parse()得到完整对象——这是对话循环能执行工具的基础。 -
Prompt Caching:在 System Prompt 末尾注入
cache_control: { type: 'ephemeral' },让 Anthropic 服务端缓存不变的内容。缓存命中时费率降低约 90%,长会话的累计节省相当可观。 -
Extended Thinking:请求参数加
thinking: { type: 'enabled', budget_tokens: N }触发扩展思考,响应流中出现thinking块。Claude Code 通过ultrathink关键词检测自动启用最大 budget。 -
成本可观测性:
message_start事件中的cache_read_input_tokens字段是缓存效果的直接证据——第一次请求创建缓存,后续命中时此值激增、成本大幅下降。
下一章将把本章的 API 客户端嵌入多轮对话循环:维护消息历史、处理 tool_use → tool_result 的往返、在上下文接近上限时触发压缩。