java流式输出

最近做AI接口,都是流式输出,所以总结下。

这里是用springboot做流式接口转发,涉及到的功能点有:tomcat、NIO、MVC异步线程池、streamingResponseBody、restTemplate(连接池)。

先给个例子:

@RestController
@Slf4j
public class StreamProxyController {

    private final RestTemplate streamingRestTemplate;

    public StreamProxyController(RestTemplate streamingRestTemplate) {
        this.streamingRestTemplate = streamingRestTemplate;
    }

    @GetMapping(value = "/proxy/stream")
    public StreamingResponseBody proxyStream(HttpServletResponse response) {

        // 告诉客户端这是流式
        response.setContentType("text/plain;charset=UTF-8");
        response.setHeader("Cache-Control", "no-store");
        response.setHeader("Connection", "keep-alive");

        log.info("Controller running on thread: {}", Thread.currentThread().getName());

        return outputStream -> {
            log.info("Streaming write running on thread: {}", Thread.currentThread().getName());

            streamingRestTemplate.execute(
                    "https://httpbin.org/stream-bytes/1048576", // 1MB 字节流
                    HttpMethod.GET,
                    requestCallback -> {
                        // 如果需要,在这里加第三方 Header(token 等)
                        requestCallback.getHeaders().add("Accept", "*/*");
                    },
                    responseExtractor -> {
                        // ⚠️ 这里非常关键:不要读成 String / byte[]
                        try (InputStream thirdPartyIn = responseExtractor.getBody()) {
                            byte[] buf = new byte[8192];
                            int len;
                            long total = 0;
                            while ((len = thirdPartyIn.read(buf)) != -1) {
                                outputStream.write(buf, 0, len);
                                outputStream.flush(); // 触发 chunked 发送
                                total += len;
                                if (total % 16384 == 0) {
                                    log.debug("Written {} bytes", total);
                                }
                            }
                        }
                        return null;
                    }
            );
        };
    }
}

1. 上面最重要的就是streamingResponseBody。

这个封装很多:告诉tomcat这是个请求我用异步方式处理(注意说法:不是非阻塞请求)维持异步线程request.startAsync()、分块返回、自主判断流式输出是否结束asyncContext.complete()。我们要做就是在resttemplate.execute中组装请求和返回。

上面案例中response中的内容就是第三方单次返回的内容,byte[8192]指的是单次读取最大的内容。也是业界的标准。实际每次返回的数据块都是远小于这个的。所以接口返回流式跟第三方返回的数据块是同步的。

2.然后是tomcat.

streamingResponseBoyd使用request.startAsync()启用了异步处理,然后tomcat收到这个信号后就会把状态从resquest-》Async。这就意味着,当前tomcat的这个请求线程已经被回收到线程池了。但是AsyncContext还在,也就是异步请求的上下文还在。

Tomcat 的约定是:

只要你是在 AsyncContext里,

你可以在你自己的线程里随便 write,

哪怕阻塞也没关系。

这样异步处理方式不占用tomcat的connector线程,但是写操作还在进行。并且这个写操作是单线程阻塞的。

这种方式在tomcat中属于:

Blocking I/O + AsyncContext

3. 两个线程池,一个是mvc异步线程池,一个是restTemplate http连接池。

前者是维持客户端到本服务的资源读写线程池。后者是管理本服务到第三方服务的http连接池。

线程池很有必要的,因为流式接口一般较慢,在大量请求时,会占用线程时间较长,业务压力较大。所以使用线程池管理很有必要,还能设置线程上限个数,防止内存溢出。

同步 Controller:请求生命周期 = 线程生命周期,线程由 Tomcat 管

异步 Controller:请求生命周期 > 线程生命周期,线程由 Spring 管

4. 提个问题:StreamingResponseBody 是异步还是非阻塞?

综上的回答就是:

是异步,不是非阻塞 IO。

  • 异步:请求处理线程不阻塞,交给别的线程

  • 阻塞:读第三方接口、写 OutputStream 仍然是阻塞调用

5. 什么是socket?

socket是管道,是操作系统给的一条通信管道,它不是线程、连接、http、tcp、tomcat创建的。

所以回到流式输出中,为啥异步返回后请求线程已经没了,仍然能够继续写数据。因为socket管道还在。

socket是管道,所以更像个实物。而线程是干活的人。tomcat不能创建它,大家只是持有它而已。

概念

餐厅类比

Socket

桌上的电话线

TCP 连接

电话接通状态

线程

服务员

read()

听电话

write()

说话

flush()

按下话筒“发”键

async

服务员把话筒交给后厨,自己走了

complete()

挂电话

一个socket一个时间只能有一个请求在写。但时间上可以复用。来回切换就是。

6.再说下Nio吧。

tomcat Nio connector是用selector 监听一堆socket,不卡connector请求线程,而是交给自己的服务的线程来处理。让服务的瓶颈不在tomcat上。

记录nginx流式输出失效问题

直接调用接口是流式输出,但是走nginx就失效了。这就在nginx也要做相关配置,比如说关掉缓存等。nginx修改了,但是流式还是没起作用。 后面发现这是应为nginx代理了两次,也就是端口1-》代理到端口2-》代理到真正的服务。这就要求每次服务的转发都要做相关配置才行。否则就会失效。

tomcat线程池与http线程池的区别

同步 Controller:请求生命周期 = 线程生命周期,线程由 Tomcat 管

异步 Controller:请求生命周期 > 线程生命周期,线程由 Spring 管

因为同步的是阻塞式的,另外再用线程池就没必要了。

构造器注入有没有使用bean

使用了,并且最开始bean实例化时就注入了。要不说是构造器。 这样也能在启动时就能鉴别出循环依赖的问题。

怎么指定bean,@Resource  和@Autowired分别怎么指定bean的?

@Autowired按“类型优先,名称兜底”匹配;

@Resource按“名称优先,类型兜底”匹配。

@Autowired
@Qualifier("userServiceImplV2")
private UserService userService;
@Resource(name = "userServiceImplV2")
private UserService userService;

注意名称兜底指的是变量名,所以我们使用@Autowired注入时,变量名要跟类名一致喽。

服务器内存溢出排查

测试环境由于资源紧张,部署的服务过多,导致服务器内存不足,就会引发占用内存最大的那个服务经常被服务器自主停掉。确认是否是这个原因的日志是:

检查服务器级别的日志:dmesg |grep -i "out of memory\|ppm\|killed process"

这个命令会输出由于内存溢出别停掉的进场id.

Logo

openEuler 是由开放原子开源基金会孵化的全场景开源操作系统项目,面向数字基础设施四大核心场景(服务器、云计算、边缘计算、嵌入式),全面支持 ARM、x86、RISC-V、loongArch、PowerPC、SW-64 等多样性计算架构

更多推荐