Java 21实践:用虚拟线程处理大模型SSE流式响应,简单得像同步代码

2026-09-09 0 239

最近在做一个小助手工具,想着能不能绕开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>”了。

Java 21实践:用虚拟线程处理大模型SSE流式响应,简单得像同步代码
收藏 (0) 打赏

感谢您的支持,我会继续努力的!

打开微信/支付宝扫一扫,即可进行扫码打赏哦,分享从这里开始,精彩与您同在
点赞 (0)

版权声明:
本站资源有的来自互联网收集整理,本站纯免费分享提供学习使用,如果侵犯了您的合法权益,请发送邮件1506151422@qq.com联系,将会及时下架删除。
本站资源仅供研究、学习交流之用,免费开源项目不代表完全可商用,若商业用途请先咨询开发企业能否商用,否则产生的一切后果将由下载用户自行承担。
原创板块未经允许不得转载,否则将追究法律责任。

淘吗网 java Java 21实践:用虚拟线程处理大模型SSE流式响应,简单得像同步代码 https://www.taomawang.com/server/java/2732.html

下一篇:

已经没有下一篇了!

常见问题

相关文章

猜你喜欢
发表评论
暂无评论
官方客服团队

为您解决烦忧 - 24小时在线 专业服务