Micronaut声明式客户端服务间字节流传输:服务器端接收无数据问题
Micronaut声明式客户端字节流传输问题解决
核心问题定位
你通过Micronaut声明式客户端传递Flowable<byte[]>实现服务间字节流传输时,客户端侧数据流验证正常,但服务器端无法接收数据,最终触发空闲超时。问题出在服务器端对流式请求体的配置或消费逻辑上,导致请求未被正确处理。
解决方案
1. 统一客户端与服务器端的媒体类型
两端必须指定匹配的流式媒体类型,避免Micronaut默认按JSON解析请求体:
- 客户端添加
@Produces注解:
import io.micronaut.http.annotation.Body; import io.micronaut.http.annotation.Post; import io.micronaut.http.annotation.Produces; import io.reactivex.Flowable; @Post("/stream") @Produces("application/octet-stream") void sendStream(@Body Flowable<byte[]> content);
- 服务器端添加
@Consumes注解:
import io.micronaut.http.annotation.Consumes; import io.micronaut.http.annotation.Post; import io.reactivex.Flowable; @Post("/stream") @Consumes("application/octet-stream") public void receiveStream(@Body Flowable<byte[]> content) { // 消费逻辑 }
2. 服务器端正确触发Flowable消费
Flowable是冷流,必须通过实际的订阅动作触发数据传输,不能仅定义订阅而无消费逻辑:
content.subscribe( chunk -> { // 处理每个字节块,如写入文件或业务处理 System.out.println("收到字节块:" + chunk.length + " 字节"); }, error -> error.printStackTrace(), () -> System.out.println("数据流传输完成") );
3. 客户端确保数据流正常完成
从InputStream转换Flowable时,必须在流读取完毕后调用onComplete(),否则服务器会一直等待数据流结束:
import io.reactivex.Flowable; import io.reactivex.BackpressureStrategy; import java.io.InputStream; import java.util.Arrays; Flowable<byte[]> createStreamFromInputStream(InputStream inputStream) { return Flowable.create(emitter -> { try (inputStream) { byte[] buffer = new byte[8192]; int readLen; while ((readLen = inputStream.read(buffer)) != -1) { emitter.onNext(Arrays.copyOf(buffer, readLen)); } emitter.onComplete(); // 必须调用,否则服务器会一直等待 } catch (Exception e) { emitter.onError(e); } }, BackpressureStrategy.BUFFER); }
4. 调整超时配置避免提前中断
如果数据流较大,需修改Micronaut的超时参数,防止因传输时间过长触发超时。在application.yml中配置:
micronaut: server: idle-timeout: 5m http: client: read-timeout: 5m
内容的提问来源于stack exchange,提问作者niebula
相关产品推荐
相关产品推荐

