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.ClientTransport 应由进程长期复用。每请求新建 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 所需语义,并拒绝超过应用上限的单行;生产实现还应限制 idevent 长度。

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 执行 gofmtgo test ./...go test -race ./...go vet ./...

15. 生产部署与故障清单

入口、模型上游和负载均衡器的 idle timeout 要大于允许的流静默时间,必要时发送注释心跳;代理关闭 response buffering 与压缩需按环境验证。readiness 只在实例能接收新流时成功,优雅关闭先摘流量,再取消或限时等待在途生成,最终关闭空闲连接。

容量按“并发流 × 平均持续时间 × 每流缓冲与连接开销”测量,不能只看 QPS。限制单用户并发,避免一个租户耗尽文件描述符。发布时灰度比较 TTFT、错误分类、token 和成本分布;回滚配置与二进制一起版本化。

上线前最后检查:非 2xx body 有界;流有最大事件和总输出;所有阻塞点响应 context;完成与异常 EOF 可区分;首输出后不重试;部分结果有标识;密钥和内容不进遥测;工具参数只在完整验证后执行。做到这些,模型流才是一个可治理的生产协议,而不是碰巧能打印字符的长连接。


系列导航与关联阅读

官方资料

本文依据 Go 官方规范、标准库文档和 Go 官方博客重新梳理;正文与示例由 WR BLOG 编写。