最近在做一个小助手工具,想着能不能绕开Python去调用大模型。翻了一堆教程,基本全是Python的,Java要么是死等整个响应体,要么就是上Spring WebFlux搞响应式。看到那一堆Mono、Flux,头就大。
后来我换了个思路:Java 21有虚拟线程,那直接每个请求开一个“虚拟线程”去同步阻塞读取不就好了吗?以前的痛苦是因为平台线程宝贵,不敢阻塞。现在虚拟线程便宜得很,阻塞IO时还会自动让出底层线程,为什么不直接用最同步的写法?
然后还真给我跑通了。我们用标准HttpClient发请求,拿到InputStream以后按行读SSE流,解析data字段,一个字符一个字符往外抛。整个过程没引入任何第三方库。
先把SSE流弄清楚
大模型那边为了“打字机效果”,大多采用SSE(Server-Sent Events)来返回数据。其实就是一个HTTP响应,body里都是类似下面这种格式:
data: {"choices":[{"delta":{"content":"你"}}]}
data: {"choices":[{"delta":{"content":"好"}}]}
data: [DONE]
每一行以 data: 开头,后面跟一段JSON字符串。两个空行之间不一定有严格规定,但我们直接按行read就能读到。遇到 [DONE] 表示结束。
明白了这一点就很简单了:我们只需要拿到HTTP响应的InputStream,然后逐行读取,再把每行中data后面的部分提取出来,用JSON工具解析出最终内容。这里我为了让教程零依赖,直接手动截取字符串。
先用JDK内置HttpServer搞一个模拟服务
为了让你能直接跑通,我不拿OpenAI真实网址做演示(免得你没有key)。咱们用Java内置的 com.sun.net.httpserver.HttpServer 在本地起一个mock接口,返回一段模拟GPT流。
话不多说,先上模拟服务代码。注意放置在一个单独的类里:
import com.sun.net.httpserver.HttpServer;
import java.io.IOException;
import java.io.OutputStream;
import java.net.InetSocketAddress;
import java.nio.charset.StandardCharsets;
public class MockLLMServer {
public static void start() throws IOException {
HttpServer server = HttpServer.create(new InetSocketAddress(9090), 0);
server.createContext("/v1/chat/completions", exchange -> {
// 设置响应头,标明是SSE
exchange.getResponseHeaders().set("Content-Type", "text/event-stream; charset=utf-8");
exchange.getResponseHeaders().set("Cache-Control", "no-cache");
exchange.sendResponseHeaders(200, 0);
String[] words = {"你", "好", ",", "今", "天", "天", "气", "不", "错", "。"};
OutputStream os = exchange.getResponseBody();
try {
for (String word : words) {
// 包装成类似OpenAI的响应结构
String json = "{"choices":[{"delta":{"content":"" + word + ""}}]}";
String event = "data: " + json + "nn";
os.write(event.getBytes(StandardCharsets.UTF_8));
os.flush();
Thread.sleep(100); // 模拟每字间隔100ms
}
os.write("data: [DONE]nn".getBytes(StandardCharsets.UTF_8));
os.flush();
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
} finally {
os.close();
}
});
server.start();
System.out.println("Mock服务已启动: http://localhost:9090");
}
}
这个模拟服务会一句一句地返回“你好,今天天气不错。”,每个字间隔100毫秒。看着像不像在打字?
然后写虚拟线程调用端
主要逻辑都在这个方法里了。我建了一个 callLLMAndPrint 方法,它做的事很简单:
1. 构造一个HTTP请求,设置必要的header。
2. 用HttpClient.newBuilder().executor(Executors.newVirtualThreadPerTaskExecutor()).build()构建一个支持虚拟线程的HttpClient?等等,HttpClient本身不需要指定executor,因为send是阻塞调用。官方文档推荐在虚拟线程中直接调用阻塞API,不需要给HttpClient配executor。咱们直接在每个虚拟线程里调用send。
来看完整代码:
import java.net.URI;
import java.net.http.HttpClient;
import java.net.http.HttpRequest;
import java.net.http.HttpResponse;
import java.io.BufferedReader;
import java.io.InputStreamReader;
import java.nio.charset.StandardCharsets;
import java.util.concurrent.Executors;
public class StreamingClient {
public static void main(String[] args) throws Exception {
// 启动mock服务
MockLLMServer.start();
// 用虚拟线程执行这次流式调用,不阻塞主线程等待输入
try (var executor = Executors.newVirtualThreadPerTaskExecutor()) {
executor.submit(() -> {
try {
callLLMAndPrint();
} catch (Exception e) {
e.printStackTrace();
}
});
// 让主线程休眠,防止虚拟机退出
Thread.sleep(10000);
}
}
private static void callLLMAndPrint() throws Exception {
HttpClient client = HttpClient.newHttpClient();
HttpRequest request = HttpRequest.newBuilder()
.uri(URI.create("http://localhost:9090/v1/chat/completions"))
.header("Content-Type", "application/json")
.header("Accept", "text/event-stream")
.POST(HttpRequest.BodyPublishers.ofString("""
{
"model": "gpt-3.5-turbo",
"stream": true,
"messages": [
{"role": "user", "content": "你好"}
]
}
""", StandardCharsets.UTF_8))
.build();
// 发送请求并拿到响应体,这里会一直阻塞到连接建立以及响应头到达
HttpResponse<InputStream> response = client.send(request, HttpResponse.BodyHandlers.ofInputStream());
System.out.println("HTTP状态码: " + response.statusCode());
if (response.statusCode() != 200) {
String err = new String(response.body().readAllBytes(), StandardCharsets.UTF_8);
throw new RuntimeException("请求失败: " + err);
}
System.out.print("AI回答: ");
try (BufferedReader reader = new BufferedReader(
new InputStreamReader(response.body(), StandardCharsets.UTF_8))) {
String line;
while ((line = reader.readLine()) != null) {
if (!line.startsWith("data:")) {
continue;
}
String data = line.substring(5).trim();
if (data.equals("[DONE]")) {
break;
}
// 直接用最原始的方式提取 content 字段 —— 只适合demo
if (data.contains("delta") && data.contains("content")) {
int idx = data.indexOf(""content":"");
if (idx != -1) {
int start = idx + ""content":"".length();
int end = data.indexOf(""", start);
String content = data.substring(start, end);
System.out.print(content);
System.out.flush();
}
}
}
}
System.out.println();
}
}
跑起来以后你就会看到“AI回答: 你好,今天天气不错。”一个字一个字地出现在终端里。虽然我这没有真实模型,但整个流程跟接OpenAI一模一样。
这里面有什么特别之处?
因为我们在main里开了一个虚拟线程去执行callLLMAndPrint,所以就算整个调用阻塞在client.send()或者readLine()上,都不会占着昂贵的平台线程。JVM会把这个虚拟线程挂起,释放下面的载体线程去干别的事。这完全就是Java 21要解决的核心痛点:让阻塞变得和异步一样省资源。
你可以试一下同时发起10个这种请求,每个都去读不同的大模型流,用虚拟线程实现,代码和同步写法完全一样,根本不需要回调地狱。要是放在Java 8时代,要么你用线程池一个线程长期占着一个连接,要么你用Netty去写一小串状态机,都麻烦死了。
如果你是真实调用OpenAI或兼容API
真实情况下,你需要改几个地方:
- 把uri改成实际的API地址,比如
https://api.openai.com/v1/chat/completions。 - 加上
Authorization: Bearer 你的KEY请求头。 - 有些兼容接口需要写
stream_options:{"include_usage":true},看服务商支持。 - 代码中的JSON解析绝对不能用substring了。建议换成Jackson或者Gson,但是那种基础代码网上很多,我这里不想把篇幅拖长,只展示用Java标准库就能连通流式。
还有一个小坑:如果响应头Content-Type不是text/event-stream,有些网关会缓冲整个body。在请求头里通常要加一条Accept: text/event-stream。注意这不一定所有后端都买账,最好的办法是自己先curl一下看看。
为什么我不建议用RestTemplate?
很多老项目用RestTemplate调用第三方接口,而且还会配上工厂让它不超时。但RestTemplate本身是IO同步阻塞,你只能一个线程一个线程去等响应。如果希望在流式过程中边返回边读,RestTemplate并不方便暴露InputStream。即便你拿到ResponseExtractor,那也绕不开对当前线程的阻塞。不是说不能用,而是相比之下,Java HttpClien的BodyHandlers.ofInputStream()天然适合流式。
关于取消操作
如果用户在界面上点击停止生成,你需要提前终止HTTP请求。用标准HttpClient,你可以调用HttpResponse<InputStream> response = client.send(...)返回后获取到body,但如果要中途取消,更优雅的做法是使用CompletableFuture。不过既然我们用了虚拟线程,最简单的是在循环读取时加一个标志位,当用户发出取消信号时直接break并关闭流。
例如:
volatile boolean stopRequested = false;
// 在循环里判断
while ((line = reader.readLine()) != null) {
if (stopRequested) {
break;
}
// ...
}
关闭reader后,底层socket连接会被释放,服务端会收到连接关闭通知,从而停止继续生成。这在真实应用里特别管用。
一个容易忽略的乱码问题
SSE里中文内容一定要用UTF-8解码。如果你在BufferedReader里忘了指定字符集,Windows本地运行时可能会乱码。解决方法已经写在上面的代码里了,就是使用new InputStreamReader(response.body(), StandardCharsets.UTF_8)。
是否可以用HttpClient的BodyHandlers.ofString()直接读?
不能用。如果调用了BodyHandlers.ofString(),HttpClient会等整个响应体全部读完才给你结果,也就是所有文字一次性全到,根本没法实现打字机效果。我们需要的是一边收到数据一边把内容展示给用户,只有ofInputStream()能做到按需拉取。
并发时如何提高效率
你可以创建一百个虚拟线程同时调用这一套逻辑,它不会消耗一百个系统线程。底层也许只有十几个平台线程在循环调度。我在自己的笔记本上试过同时开启十个流式请求,全部顺畅打印,CPU占用也不高。
那是不是说有了虚拟线程后,我们就可以无脑为每个请求开一个线程?也不是。因为真正限制并发数的不是线程,而是对方API的QPS和你的网络带宽。但至少虚拟线程让你不用再在一开始为了“性能”写复杂的异步代码,这一点挺香的。
总结一下,重点就一句话
遇到流式IO大模型接口,直接开一个虚拟线程,用最笨的同步读法,既简单又高效。
Java 21已经出来很久了,还是很多人只把它当成普通的LTS版本,没有真正去体验虚拟线程带来的编程模型变化。如果你还在维护Java 8项目,可以先自己写个demo练习一下这种写法。我估计,等你习惯了,就不会再想碰响应式那一堆“Flux<DataBuffer>”了。

