给一个 LLM 流式系统写回归测试,最容易犯的错是用批处理系统的断言习惯去断言流:
- "最终收到了全部 N 条数据" —— 坏实现也能通过:攒到最后一次性爆发的实现,和逐条流动的实现,在这个断言下无法区分;
- "收到了正确的内容" —— 流根本没关闭也能通过:很多流式泄漏在生产上被客户端行为掩盖,内容全部送达,但资源已经泄露。
我们为一个 Agent 平台补流式测试基线时,把两类真正的流式缺陷各抓到了一个实例。这篇文章总结两条断言原则和一套可复用的测试 harness。
1. 两条核心断言原则
原则一:断言"流式",而不是断言"完成"
判据:首元素到达时间 ≪ 生产者发完全部数据的时间。
流式系统的本质承诺是"边产生边消费"。验证方式不是数收到多少条,而是比较两个时间:
T_first = 第一个元素到达消费者的时刻
T_total = 生产者发出全部数据所需的时间(由测试自己控制)
断言:T_first - T_start < 显著小于 T_total 的阈值
以我们的回归测试为例:模拟网关以 10ms 间隔吐 600 个 SSE chunk(T_total ≥ 6s),断言首 chunk 在 4 秒内到达。被修复前的坏实现(Reactor 背压饿死,见文末关联复盘)会在第 ~256 个 chunk 后断流、全部缓冲到 [DONE] 才爆发——首元素到达时刻 ≈ T_total,必然越过阈值。
阈值要留出 CI 抖动余量,但必须显著小于 T_total,否则失去判别力。
原则二:断言"终止",而不是只断言内容
流式泄漏会被客户端行为掩盖。
真实案例:我们的 JSON 契约 Agent 在"模型返回非法 JSON → 自动重试"的路径上,重试流成功发完 done 事件并完成持久化——但外层流从未发出 complete 信号。生产上毫无症状:前端收到 done 就主动断开 SSE 连接,取消信号顺带清理了一切。
直到测试用 collectList().block(timeout) 等待流正常结束,30 秒超时才暴露:外层 Flux.create 的 emitter 在重试分支里永远不会被 complete。
教训:对任何返回 Flux/Flowable/响应流的 API,测试必须等待并断言正常终止(onComplete),并给等待加一个远小于"无限"的超时。只断言内容的测试,对泄漏是盲的。
2. 测试 harness:本地模拟流式上游
不依赖真实 LLM 网关,用 JDK 自带的 com.sun.net.httpserver.HttpServer 模拟一个逐行吐 SSE 的上游:
server = HttpServer.create(new InetSocketAddress("127.0.0.1", 0), 0);
server.createContext("/chat/stream", exchange -> {
exchange.getRequestBody().readAllBytes();
exchange.getResponseHeaders().set("Content-Type", "text/event-stream");
exchange.sendResponseHeaders(200, 0); // chunked
try (var os = exchange.getResponseBody()) {
for (int i = 0; i < 600; i++) {
os.write("data: {\"choices\":[{\"delta\":{\"content\":\"x\"}}]}\n\n"
.getBytes(StandardCharsets.UTF_8));
os.flush(); // 关键:逐行刷出,模拟真实流
Thread.sleep(10); // 关键:控制 T_total ≥ 6s
}
os.write("data: [DONE]\n\n".getBytes(StandardCharsets.UTF_8));
os.flush();
}
exchange.close();
});
server.start();
被测流要刻意复现事故的三个条件(缺一不可,见关联复盘):
var firstItemAt = new AtomicLong();
long start = System.currentTimeMillis();
var items = chatModel.stream(new Prompt("hi"))
.subscribeOn(Schedulers.boundedElastic()) // 条件 ②:订阅在弹性线程池
.concatMap(cr -> Mono.delay(Duration.ofMillis(5)).thenReturn(cr), 32) // 慢消费 + 背压
.doOnNext(cr -> firstItemAt.compareAndSet(0, System.currentTimeMillis()))
.collectList()
.block(Duration.ofSeconds(60)); // 原则二:等待终止 + 有限超时
assertThat(items).hasSizeGreaterThanOrEqualTo(600); // 不丢数据
assertThat(firstItemAt.get() - start).isLessThan(4_000); // 原则一:流式,非爆发
要点:
- chunk 数必须显著超过链路中最大的 prefetch(Reactor 默认 256)。事故复盘里短回复永远不触发,就是因为没越过预取边界——测试数据量要按"越过边界"设计;
- 慢消费者用
concatMap(delay)而不是delayElements:前者保持背压,后者把背压换成定时器缓冲,测的不是一回事; - 生产者是真实 HTTP 连接 + 真实线程,不是 mock 出来的
Flux.just——背压问题恰恰发生在真实 IO 与调度器的交界处。
3. 超时兜底也要测"快失败"
流式读的超时不是"总时长超时",而是空闲看门狗:两个数据块之间超过阈值即判定上游挂起。配置值必须远大于正常 token 间隔(我们取 300s),否则会误杀正常的长思考流。
测试方式是让模拟网关接收请求后永不响应,断言调用方在读超时内快速失败,而不是无限挂起:
assertThatThrownBy(() -> model.stream(prompt)
.subscribeOn(Schedulers.boundedElastic())
.collectList().block(Duration.ofSeconds(20)))
.hasRootCauseInstanceOf(SocketTimeoutException.class);
// 并且整个过程远快于"永远"
附录 A:mock Spring AI ChatClient 流式链路
ChatClient.prompt().system(...).user(...).options(...).advisors(...).tools(...).stream().chatResponse() 是一条全 fluent 链路。手工逐个 stub 极其繁琐,标准姿势是:
var requestSpec = mock(ChatClient.ChatClientRequestSpec.class, Mockito.RETURNS_SELF);
var streamSpec = mock(ChatClient.StreamResponseSpec.class);
when(chatClient.prompt()).thenReturn(requestSpec);
when(requestSpec.stream()).thenReturn(streamSpec);
when(streamSpec.chatResponse()).thenReturn(cannedFlux1, cannedFlux2);
RETURNS_SELF 让链上所有返回自身类型的方法自动返回 mock 自己;只有返回其他类型的 stream()/call() 需要单独 stub。多轮 thenReturn(f1, f2) 可以模拟"第一次返回非法 JSON、重试后返回合法 JSON"这类序列。注意一个坑:thenReturn(array[0], array) 会把 array 按 varargs 展开导致第一个值重复——第二个参数要传 Arrays.copyOfRange(array, 1, ...)。
附录 B:FunctionToolCallback 返回值的 JSON 规范化
Spring AI 的 FunctionToolCallback 对函数返回值做 JSON 规范化:本身合法的 JSON 文本原样透传,非 JSON 文本被 JSON 字符串化(加引号)。即工具返回 "42" 时 LLM 看到 42(数字),返回 "表达式非法" 时 LLM 看到 "表达式非法"(带引号的 JSON 字符串)。写断言时按"规范化后"的形态断言,否则会出现同一条链路有的值带引号、有的不带的困惑。
关联阅读:2026-07-16 Reactor 背压饿死复盘 —— 本文原则一对应的事故本体。