Go 基础体系 · 第 103/113 篇。示例统一基于 Go 1.26.4;核心片段可能省略 package 与 import,完整程序可直接按文中结构运行。
Go 调用 LLM API:统一客户端、SSE 流式输出、取消与重试
本文以 Go 1.26.4、标准库 net/http、HTTP/1.1、UTF-8 JSON 与 WHATWG Server-Sent Events(SSE)为稳定基线。示例采用 OpenAI 风格的 /v1/chat/completions 线格式说明通用问题,但领域层不依赖某家 SDK;实际字段、模型名、限额和结束事件必须以部署时锁定的提供商 API 版本为准。
调用模型不是“发 JSON、读字符串”这么简单。生产客户端要同时管理请求预算、连接复用、流的逐事件解析、增量状态合并、用户取消、有限重试、用量核算和敏感数据边界。尤其是流式响应,一旦首个字节已经交给用户,就不再等价于可透明重试的普通 HTTP 请求。
1. 先定义稳定的领域接口
把供应商 request/response 类型扩散到 handler、业务服务和测试,会让切换模型或升级 SDK 变成全仓修改。应用真正稳定的语义通常只有消息、生成参数、文本增量、工具调用增量、结束原因和用量。接口接收 context.Context,流必须显式关闭。
type Model interface {
Generate(context.Context, Request) (Response, error)
Stream(context.Context, Request) (Stream, error)
}
type Stream interface {
Next() bool
Event() Event
Err() error
Close() error
}
type Request struct {
Model string
Messages []Message
Temperature float64
MaxTokens int
}
type Event struct {
TextDelta string
ToolDelta *ToolCallDelta
Usage *Usage
Finish string
}
Next 为 false 后检查 Err,与 bufio.Scanner、数据库 rows 的约定一致。Close 应幂等;调用方始终 defer stream.Close()。领域错误至少区分无效请求、认证、限流、上游不可用、协议损坏、取消和超时,不能让业务依赖供应商错误字符串。
2. 配置、密钥与可重复版本基线
模型名、Base URL、API key、组织或项目标识、超时和上限由配置注入。Base URL 只允许启动时校验过的 HTTPS 主机,不能直接接受用户 URL,否则服务会成为 SSRF 代理。密钥来自 secret manager 或进程环境,绝不写入仓库、URL、日志和 trace 属性。
llm:
provider: openai-compatible
base_url: https://api.example.com/v1
model: production-chat-2026-08
request_timeout: 45s
response_header_timeout: 10s
max_error_body_bytes: 65536
max_sse_event_bytes: 1048576
max_output_tokens: 2048
配置中的模型别名应解析成运维批准的实际模型与版本快照。禁止生产使用含义会漂移的 latest。启动日志可记录 provider、模型别名和 API 版本,但只记录 key 的 secret 引用名,不记录值。不同租户若使用不同凭据,缓存和连接指标也必须按低基数 provider 维度组织。
3. HTTP Client 的所有权与连接复用
http.Client 和 Transport 应由进程长期复用。每请求新建 Transport 会丢失连接池并放大 TLS 开销。连接、TLS、响应头和整体生成分别设预算;流式调用不能用过短的 Client.Timeout,通常由请求 context 控制整条流。
transport := &http.Transport{
Proxy: http.ProxyFromEnvironment,
MaxIdleConns: 100,
MaxIdleConnsPerHost: 20,
IdleConnTimeout: 90 * time.Second,
TLSHandshakeTimeout: 5 * time.Second,
ResponseHeaderTimeout: 10 * time.Second,
ForceAttemptHTTP2: true,
}
client := &http.Client{Transport: transport}
应用关闭时可调用 transport.CloseIdleConnections(),但不能在单次请求结束时调用。代理是否支持长响应、空闲多久回收、是否缓冲 SSE 都应实测。HTTP/2 可以承载流式响应,但 SSE 的“事件”仍由线格式定义,绝不能把一次 socket read 或 DATA frame 当成一个 token。
4. 构造请求并有界读取错误
先在内存中编码经过大小限制的请求,再用 NewRequestWithContext 绑定生命周期。只有生成请求成功后才发出;响应无论成功失败都关闭 body。非 2xx 的错误体可能是 HTML、巨大内容或恶意数据,只读固定上限并进行 JSON 尝试解析。
payload, err := json.Marshal(wireRequest{
Model: req.Model, Messages: req.Messages, Stream: true,
})
if err != nil {
return nil, fmt.Errorf("encode model request: %w", err)
}
httpReq, err := http.NewRequestWithContext(ctx, http.MethodPost,
baseURL+"/chat/completions", bytes.NewReader(payload))
if err != nil {
return nil, fmt.Errorf("build model request: %w", err)
}
httpReq.Header.Set("Authorization", "Bearer "+apiKey)
httpReq.Header.Set("Content-Type", "application/json")
httpReq.Header.Set("Accept", "text/event-stream")
resp, err := client.Do(httpReq)
if err != nil {
return nil, fmt.Errorf("call model: %w", err)
}
if resp.StatusCode < 200 || resp.StatusCode >= 300 {
defer resp.Body.Close()
body, _ := io.ReadAll(io.LimitReader(resp.Body, 64<<10))
return nil, classifyHTTP(resp.StatusCode, body, resp.Header)
}
不要记录完整 payload。需要排障时记录经过批准的模板版本、字符数、消息数和稳定请求 ID。服务端返回的 request ID 可以进入日志,但先限制长度并去除控制字符。
5. SSE 不是逐行 JSON
SSE 事件由空行结束;字段形式为 name:value,冒号后的一个空格可忽略。多个 data: 行以换行连接,event: 指定事件类型,id: 提供恢复游标,冒号开头是注释。CRLF 与 LF 都可能出现。一个 JSON 对象可以跨多个 data: 行,但不会跨事件边界。
event: response.output_text.delta
id: evt_1042
data: {"delta":"Go 的"}
: keep-alive
event: response.completed
data: {"usage":{"input_tokens":81,"output_tokens":34}}
bufio.Scanner 默认 token 上限约 64 KiB,模型工具参数或结构化结果可能超过它。必须调用 Scanner.Buffer(initial, maxEventBytes),达到上限就以协议错误终止。只用 strings.TrimPrefix(line, "data: ") 会漏掉多行、无空格、注释与事件类型。
6. 一个有界 SSE 解析器
解析器累计当前事件字段,遇到空行才派发。下面保留 SSE 所需语义,并拒绝超过应用上限的单行;生产实现还应限制 id 和 event 长度。
type SSEEvent struct {
Type string
ID string
Data []byte
}
func scanSSE(r io.Reader, emit func(SSEEvent) error) error {
s := bufio.NewScanner(r)
s.Buffer(make([]byte, 32<<10), 1<<20)
var event SSEEvent
var data []string
flush := func() error {
if len(data) == 0 {
return nil
}
event.Data = []byte(strings.Join(data, "\n"))
if err := emit(event); err != nil {
return err
}
event, data = SSEEvent{}, data[:0]
return nil
}
for s.Scan() {
line := strings.TrimSuffix(s.Text(), "\r")
if line == "" {
if err := flush(); err != nil { return err }
continue
}
if strings.HasPrefix(line, ":") { continue }
name, value, _ := strings.Cut(line, ":")
value = strings.TrimPrefix(value, " ")
switch name {
case "event": event.Type = value
case "id": event.ID = value
case "data": data = append(data, value)
}
}
if err := s.Err(); err != nil { return fmt.Errorf("scan SSE: %w", err) }
return flush()
}
EOF 不等于业务完成:只有提供商定义的完成事件或明确 finish reason 才表示完整生成。EOF 之前无完成标志,应返回 io.ErrUnexpectedEOF 类别,让上层知道已有文本可能只是部分结果。
7. 从线事件映射成领域事件
不同 API 可能发送 [DONE]、类型化 completed 事件,或在最后一个 choice 中放 finish reason。适配器负责严格解码和归一化;未知事件可以按兼容策略忽略并计数,已知事件字段损坏则应失败。
{
"id": "gen_01",
"choices": [{
"index": 0,
"delta": {"content": "并发"},
"finish_reason": null
}]
}
不要假设每个 delta 都含文本:它可能只声明角色、工具调用 ID、函数名、参数片段或用量。多 choice 必须按 index 分开聚合;不支持多候选时在请求中固定为一,并对意外 index 报协议错误。空 delta 也可能是合法心跳或状态变化。
8. 增量文本与工具调用合并
文本增量按收到顺序追加字节即可,因为 JSON 解码已还原 UTF-8 字符串。工具参数更棘手:arguments 通常是 JSON 文本片段,在完成前不是合法 JSON。按 (choice index, tool-call index/id) 保存 builder,函数名只接受一致值,直到 finish 后再整体解码。
type callState struct {
ID string
Name string
Args strings.Builder
}
func (s *callState) Merge(delta ToolCallDelta) error {
if delta.ID != "" && s.ID != "" && delta.ID != s.ID {
return errors.New("tool call id changed")
}
if delta.Name != "" && s.Name != "" && delta.Name != s.Name {
return errors.New("tool name changed")
}
if s.ID == "" { s.ID = delta.ID }
if s.Name == "" { s.Name = delta.Name }
if s.Args.Len()+len(delta.Arguments) > 256<<10 {
return errors.New("tool arguments too large")
}
_, _ = s.Args.WriteString(delta.Arguments)
return nil
}
不能对每个参数片段 json.Unmarshal,也不能自己补括号猜测模型意图。完成后再用严格 decoder 验证一次,并交给工具层进行 schema、业务和权限校验。若流中途断开,所有未完成调用都作废,绝不执行“看起来已经够完整”的参数。
9. 向浏览器转发流与背压
网关可把领域事件重新编码为自己的 SSE 协议,避免前端绑定供应商格式。写 header 后,每个事件写入并通过 http.NewResponseController(w).Flush() 刷新。浏览器断开会取消 r.Context(),该 context 必须原样传播至模型请求。
func relay(w http.ResponseWriter, r *http.Request, model Model, req Request) error {
stream, err := model.Stream(r.Context(), req)
if err != nil { return err }
defer stream.Close()
w.Header().Set("Content-Type", "text/event-stream")
w.Header().Set("Cache-Control", "no-cache")
controller := http.NewResponseController(w)
for stream.Next() {
data, err := json.Marshal(stream.Event())
if err != nil { return err }
if _, err := fmt.Fprintf(w, "event: delta\ndata: %s\n\n", data); err != nil { return err }
if err := controller.Flush(); err != nil { return err }
}
return stream.Err()
}
下游写慢时,上游读取也会变慢,这是一种自然背压;但提供商可能继续占用连接并计费。若加入 channel 解耦,必须有界并定义满载策略。文本生成通常不能随意丢 delta,所以队列满时应取消整轮,而不是制造缺字答案。
10. 取消、代际编号与资源清理
每轮生成创建独立 context 和 generation ID。用户点击停止、同会话发起新轮、浏览器断开或服务关闭都取消该轮。发送端在写出前检查 generation ID 仍是会话当前值,防旧轮迟到事件串入新轮。
取消后 client.Do 或 body read 通常返回包装错误。分类时优先检查 context.Cause(ctx)、errors.Is(err, context.Canceled) 和 DeadlineExceeded,不要把用户停止记成上游故障。关闭顺序是:触发 cancel、停止接受事件、关闭 response body、等待拥有的 goroutine 退出。不得为一次流启动无人等待的 goroutine。
取消不是费用撤销保证。供应商可能已生成或计费部分 token;最终 usage 缺失时记录“未知/估算”,不能记为零。若业务需要保存部分答案,明确标记 partial=true 与终止原因,不能冒充完整消息进入后续上下文。
11. 重试只能发生在明确边界
连接建立失败、响应头前 429/502/503/504、且剩余 deadline 足够时,可以对幂等生成请求做有限重试。生成虽然不修改传统数据库,却可能触发计费、缓存或供应商 batch,因此最好携带供应商支持的幂等键。退避使用全抖动并尊重合法 Retry-After。
func waitRetry(ctx context.Context, delay time.Duration) error {
timer := time.NewTimer(delay)
defer timer.Stop()
select {
case <-timer.C:
return nil
case <-ctx.Done():
return context.Cause(ctx)
}
}
一旦任何用户可见 delta 已写出,就不透明重试。第二次生成不会从同一随机状态继续,直接拼接会重复、缺失或改变工具调用。除非 API 提供服务端游标和精确定义的 resume 协议,否则应结束该轮并告诉客户端结果不完整。认证失败、无效参数、内容策略拒绝也不重试。
12. 结构化输出与流式 JSON
请求结构化结果时发送受支持的 JSON Schema,并固定 additionalProperties: false、required、枚举和长度。流中的文本仍可能是半个 JSON token;只做增量展示,不在完成前执行业务动作。完成后先限制总字节,再严格解析。
func decodeStrict[T any](raw []byte, max int64) (T, error) {
var value T
decoder := json.NewDecoder(io.LimitReader(bytes.NewReader(raw), max+1))
decoder.DisallowUnknownFields()
if err := decoder.Decode(&value); err != nil { return value, err }
if decoder.Decode(&struct{}{}) != io.EOF {
return value, errors.New("trailing JSON value")
}
return value, nil
}
JSON 语法正确不代表符合 schema,更不代表业务有效。日期范围、资源存在性、租户归属仍由应用验证。解析失败可发起一次“修复请求”,但它是新的可计费生成,应计入总步骤与 token 预算,且不能把原始秘密数据完整回显。
13. 可观测性、用量与成本
每轮记录低基数字段:provider、模型版本、是否流式、结果类别、finish reason;指标包括请求数、首事件延迟 TTFT、总耗时、输入/输出 token、取消率、流中断率、重试次数和估算成本。trace 将 DNS/TLS/响应头、首 token、工具阶段分开,但不创建“每 token 一个 span”。
提示、模型输出和工具参数默认属于敏感内容,不放日志、metric label 或 span attribute。需要抽样留存时经过用户授权、脱敏、加密、访问审计和短保留期。token 数以供应商最终 usage 为结算事实;本地 tokenizer 只用于请求前预算和 usage 缺失时估算,并标明估算来源。
成本保护同时限制输入字符/token、输出 token、并发流、每租户日预算和工具循环总预算。接近限额时在开始生成前拒绝,比生成一半再中断更可解释。
14. 测试协议而不依赖真实模型
使用 httptest.Server 精确发送分片 SSE:把一个 JSON 拆成多次 Write、多行 data、CRLF、注释、超长行、未知事件、无完成 EOF 和延迟事件。测试必须证明解析与 TCP 分片无关,取消能让 handler 和客户端及时退出,首 delta 后不会重试。
func TestStreamingFragments(t *testing.T) {
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "text/event-stream")
controller := http.NewResponseController(w)
_, _ = io.WriteString(w, "data: {\"delta\":\"Go")
_ = controller.Flush()
_, _ = io.WriteString(w, " 并发\"}\n\ndata: [DONE]\n\n")
}))
defer server.Close()
// 调用真实适配器,断言只得到一个完整 delta 和 completed。
}
合并器用表驱动测试覆盖工具 ID/名称变化、交错 call index、空片段和大小上限;用 fuzz 测 SSE 与 JSON 输入永不 panic、永不无限分配。CI 执行 gofmt、go test ./...、go test -race ./... 和 go vet ./...。
15. 生产部署与故障清单
入口、模型上游和负载均衡器的 idle timeout 要大于允许的流静默时间,必要时发送注释心跳;代理关闭 response buffering 与压缩需按环境验证。readiness 只在实例能接收新流时成功,优雅关闭先摘流量,再取消或限时等待在途生成,最终关闭空闲连接。
容量按“并发流 × 平均持续时间 × 每流缓冲与连接开销”测量,不能只看 QPS。限制单用户并发,避免一个租户耗尽文件描述符。发布时灰度比较 TTFT、错误分类、token 和成本分布;回滚配置与二进制一起版本化。
上线前最后检查:非 2xx body 有界;流有最大事件和总输出;所有阻塞点响应 context;完成与异常 EOF 可区分;首输出后不重试;部分结果有标识;密钥和内容不进遥测;工具参数只在完整验证后执行。做到这些,模型流才是一个可治理的生产协议,而不是碰巧能打印字符的长连接。
系列导航与关联阅读
- 系列入口:Go 完整技术体系学习路线:从语法、并发到框架、中间件与 AI
- 上一篇:Go AI 应用学习路线:LLM、RAG、Agent、MCP 与生产治理
- 下一篇:Go LLM 结构化输出与工具调用:JSON Schema、循环和权限
- 延伸:Go net/http 基础:Server、Handler、Middleware 与 Client 超时
- 延伸:Go context 完整指南:取消、超时、Deadline 与 Value
- 延伸:Go WebSocket 与 SSE:实时通信、心跳、背压和断线恢复
官方资料
本文依据 Go 官方规范、标准库文档和 Go 官方博客重新梳理;正文与示例由 WR BLOG 编写。

评论
0 条讨论