练习 2:流式输出

练习 1 里你把问题发出去,然后等。几秒钟后,答案整块砸下来。 可你每天用的每一个 AI 产品都不是这样——字是一个一个蹦出来的。

不是产品做了什么特效。模型本来就是一个 token 一个 token 生成的, "整块"才是加工出来的假象:练习 1 里,服务端替你把流攒成了一份完整 JSON 才发货。 这一章,把攒的动作拿回你自己手里。

敲进去

新建目录 ex02,在里面跑 go mod init ex02,然后新建 main.go。 大部分和练习 1 相同——注意不同的那几处,每一处都有注释:

// Learn Agent the Hard Way — 练习 2:流式输出
//
// 和练习 1 同一个请求,多一个字段:stream。
// 响应从一份 JSON 变成一条流——你的 harness 从此有了"边生成边看"的感官。
package main

import (
	"bufio"
	"bytes"
	"encoding/json"
	"fmt"
	"io"
	"net/http"
	"os"
	"strings"
)

// request 比练习 1 多了 Stream 和 StreamOptions 两个字段。
type request struct {
	Model         string         `json:"model"`
	Messages      []message      `json:"messages"`
	MaxTokens     int            `json:"max_tokens,omitempty"`
	Stream        bool           `json:"stream"`
	StreamOptions *streamOptions `json:"stream_options,omitempty"`
}

// streamOptions.IncludeUsage 让服务端在流的最后补一个带 usage 的块。
// 不带这个选项,OpenAI 和多数兼容服务商在流式下不报 token 数——账单直接消失。
type streamOptions struct {
	IncludeUsage bool `json:"include_usage"`
}

type message struct {
	Role    string `json:"role"`
	Content string `json:"content"`
}

// chunk 是流里每条 data: 行的形状。对照练习 1 的 response:
// message 变成了 delta——同一个位置,从"完整的话"变成"新吐出的几个字"。
// 其余字段都还在原地。
type chunk struct {
	Choices []struct {
		Delta struct {
			Content string `json:"content"`
		} `json:"delta"`
		FinishReason string `json:"finish_reason"` // 只在最后一个内容块上非空
	} `json:"choices"`
	Usage *struct { // 只在 include_usage 补发的终块上非空
		PromptTokens     int `json:"prompt_tokens"`
		CompletionTokens int `json:"completion_tokens"`
	} `json:"usage"`
	Error *struct { // 200 之后服务端也可能在流里报错——形状和练习 1 相同
		Message string `json:"message"`
		Type    string `json:"type"`
	} `json:"error"`
}

func main() {
	if len(os.Args) < 2 {
		fmt.Fprintln(os.Stderr, `用法: ./ex02 "你的问题"`)
		os.Exit(1)
	}
	apiKey := os.Getenv("OPENAI_API_KEY")
	model := os.Getenv("MODEL")
	if apiKey == "" || model == "" {
		fmt.Fprintln(os.Stderr, "需要环境变量 OPENAI_API_KEY 和 MODEL")
		fmt.Fprintln(os.Stderr, `例: export OPENAI_API_KEY=sk-xxxx`)
		fmt.Fprintln(os.Stderr, `    export MODEL=deepseek-v4-flash`)
		fmt.Fprintln(os.Stderr, `    export OPENAI_BASE_URL=https://api.deepseek.com/v1  # 不设则默认 OpenAI 官方`)
		os.Exit(1)
	}
	base := os.Getenv("OPENAI_BASE_URL")
	if base == "" {
		base = "https://api.openai.com/v1"
	}

	body, _ := json.Marshal(request{
		Model:         model,
		MaxTokens:     1024,
		Messages:      []message{{Role: "user", Content: os.Args[1]}},
		Stream:        true,
		StreamOptions: &streamOptions{IncludeUsage: true},
	})

	req, err := http.NewRequest("POST", base+"/chat/completions", bytes.NewReader(body))
	if err != nil {
		fmt.Fprintln(os.Stderr, err)
		os.Exit(1)
	}
	req.Header.Set("Content-Type", "application/json")
	req.Header.Set("Authorization", "Bearer "+apiKey)
	req.Header.Set("Accept", "text/event-stream") // 声明:我要的是事件流

	resp, err := http.DefaultClient.Do(req)
	if err != nil {
		fmt.Fprintln(os.Stderr, "请求失败:", err)
		os.Exit(1)
	}
	defer resp.Body.Close()

	// 流式请求失败在"开流之前":非 200 时响应体是普通 JSON,不是流。
	if resp.StatusCode != 200 {
		raw, _ := io.ReadAll(resp.Body)
		fmt.Fprintf(os.Stderr, "HTTP %d: %s\n", resp.StatusCode, raw)
		os.Exit(1)
	}

	var (
		finish        string
		inTok, outTok int
	)
	scanner := bufio.NewScanner(resp.Body)
	for scanner.Scan() {
		line := scanner.Text()
		// SSE 的全部语法就这一条:以 "data:" 开头的行,后面跟一份 JSON。
		// 冒号后那个空格按规范是可选的——OpenAI 和 DeepSeek 会发,
		// 有的兼容服务商不发,所以两段都要剥。
		if !strings.HasPrefix(line, "data:") {
			continue
		}
		data := strings.TrimPrefix(strings.TrimPrefix(line, "data:"), " ")
		if data == "" {
			continue
		}
		if data == "[DONE]" { // 终止哨兵:流到头了
			break
		}

		var c chunk
		if err := json.Unmarshal([]byte(data), &c); err != nil {
			fmt.Fprintf(os.Stderr, "\n解析失败: %v\n原始行: %s\n", err, data)
			os.Exit(1)
		}
		if c.Error != nil {
			fmt.Fprintf(os.Stderr, "\nAPI 错误 [%s]: %s\n", c.Error.Type, c.Error.Message)
			os.Exit(1)
		}
		if c.Usage != nil {
			inTok, outTok = c.Usage.PromptTokens, c.Usage.CompletionTokens
		}
		if len(c.Choices) == 0 { // include_usage 的终块没有 choices
			continue
		}
		if c.Choices[0].Delta.Content != "" {
			fmt.Print(c.Choices[0].Delta.Content) // 到手就打,不攒
		}
		if c.Choices[0].FinishReason != "" {
			finish = c.Choices[0].FinishReason
		}
	}
	if err := scanner.Err(); err != nil {
		fmt.Fprintln(os.Stderr, "\n读流失败:", err)
		os.Exit(1)
	}

	fmt.Println()
	fmt.Fprintf(os.Stderr, "\n[输入 %d tokens · 输出 %d tokens · finish_reason=%s]\n",
		inTok, outTok, finish)
}

先别问为什么。敲完,跑起来,我们再回头讲。

跑起来

环境变量和练习 1 完全一样,云端 DeepSeek 或本机 Ollama 任选:

go build -o ex02 . && ./ex02 "用三句话介绍一下 Go 语言"

你应该看到什么

和练习 1 一样的回答——但这次字是一个一个蹦出来的, 和你用过的所有 AI 产品一样。最后仍然是那行账单:

[输入 90 tokens · 输出 117 tokens · finish_reason=stop]

故意问个长问题,盯着屏幕看。然后回去用练习 1 的 ex01 问同一个问题, 感受那几秒的干等。这就是所有 AI 产品都用流式的原因: 生成速度一个字都没变,变的是你什么时候开始看到。

发生了什么

stream: true 让服务端不再攒——每生成几个 token 就发一行。 这个格式叫 SSE(Server-Sent Events),听着像什么大协议,其实就是 一个一直不关的 HTTP 响应,里面一行一行地写。这是我本机 Ollama 对 "数到3" 的真实完整回应:

data: {"model":"qwen3:4b-instruct","choices":[{"delta":{"role":"assistant","content":"1"},"finish_reason":null}]}
data: {"model":"qwen3:4b-instruct","choices":[{"delta":{"content":"  \n"},"finish_reason":null}]}
data: {"model":"qwen3:4b-instruct","choices":[{"delta":{"content":"2"},"finish_reason":null}]}
data: {"model":"qwen3:4b-instruct","choices":[{"delta":{"content":"  \n"},"finish_reason":null}]}
data: {"model":"qwen3:4b-instruct","choices":[{"delta":{"content":"3"},"finish_reason":null}]}
data: {"model":"qwen3:4b-instruct","choices":[{"delta":{"content":""},"finish_reason":"stop"}]}
data: {"model":"qwen3:4b-instruct","choices":[],"usage":{"prompt_tokens":11,"completion_tokens":6,"total_tokens":17}}
data: [DONE]

(每行还有 id、created 等字段,为排版删了。你自己抓一份完整的——加分练习 1。)

盯着这八行,把协议读完:

delta 顶替了 message 练习 1 里 choices[0].message 装的是完整的话, 现在 choices[0].delta 装的是新吐出的几个字。名字换了,位置没换—— 这不是一个新协议,是同一个协议的分期付款版。

finish_reason 单独坐一班车。 它在一个 content 为空的块上到达。 练习 1 教你"判断说完了没有,永远看 finish_reason"——流式下这条纪律不变, 只是你要在一串块里等它出现。

账单最后才来,而且要你主动要。 那个 "choices":[] 的块只装 usage—— 没有 stream_options.include_usage,OpenAI 官方和 Ollama 都不发它 (有的服务商会好心补发,加分练习 3 你会看到 DeepSeek 和 Ollama 行为不一样)。 所以代码里 len(c.Choices) == 0 时要 continue 而不是报错:没有 choices 的块是合法的。

[DONE] 是纯文本哨兵,不是 JSON。 先判断它再解析,顺序反了就是一个解析错误。

最后想一层:你的 harness 从此有了感官——模型生成的每个字, 在它落地的瞬间你就看到了。现在流过这条管道的只有文字; 到练习 5,模型伸手调工具时,工具名和参数也是从这条管道里 一个碎片一个碎片流进来的。你今天写的这个 for 循环,就是将来 agent loop 的入口。

常见问题

  • 字没有逐个蹦,整段一次砸出来:中间有东西在攒——常见是公司代理或网关 开了响应缓冲。用 curl -N(加分练习 1)直接打服务商,如果裸流是逐行到的, 问题就在你和它之间。
  • 解析失败,原始行是半截 JSON:连接中途断了,最后一行没写完。 健康的服务端每行都是完整 JSON。
  • token 数是 0:忘了 stream_options.include_usage。另外有的兼容服务商 确实不发 usage 块——收不到就是 0,代码不报错,这是协议边缘的现实。
  • 等了很久一个字都没有:你八成用了思考型模型(练习 1 的坑)。它在想, 想的过程也在流里——只是藏在另一个字段(加分练习 2)。

加分练习

  1. 用 curl 看裸流,你的代码解析的就是这个东西:

    curl -N http://localhost:11434/v1/chat/completions \
      -H "Content-Type: application/json" \
      -d '{"model":"qwen3:4b-instruct","stream":true,"messages":[{"role":"user","content":"数到3"}]}'
    
  2. Delta 加一个 Reasoning string 字段(tag 写 json:"reasoning")并打印它, 换思考型模型 qwen3:4b 跑一次——练习 1 那个"回答是空的"之谜, 现在你能亲眼看着它把预算想光。顺便注意:这个字段名各家不统一 (Ollama 叫 reasoning,DeepSeek 叫 reasoning_content)—— OpenAI 官方协议里没有它,思考字段是各家自己长出来的,兼容的边缘从来没有看上去那么齐。

  3. 删掉 StreamOptions 那行,DeepSeek 和 Ollama 各跑一遍。Ollama 的账单变成 0, DeepSeek 的还在——include_usage 在协议里是"不主动要就没有", 但有的服务商无论如何都发。你的 harness 不能赌服务商的好心:要数据,就明说。

  4. data == "[DONE]" 的判断挪到 json.Unmarshal 之后,跑一次,看报什么错。 然后把它挪回来。协议里总有几个不是 JSON 的东西,解析顺序就是防御顺序。