Micronaut保留HTTP响应压缩数据及性能优化问题
解决方案
要避免Netty自动解压gzip响应、省去“解压再压缩”的性能损耗,核心是禁用客户端的自动解码逻辑,直接获取服务器返回的原始压缩字节流,再直接用于S3分片上传。以下是具体实现步骤:
1. 禁用Micronaut客户端的自动gzip解压
Micronaut默认会根据Accept-Encoding头自动添加gzip解码器,需要手动关闭这个行为,有两种方式可选:
方式一:配置文件快速设置
在application.yml中针对目标客户端关闭压缩解码支持:
micronaut: http: client: services: # 对应@Client("${url}")中的服务标识,比如url配置的是data-server,这里就写data-server ${url对应的服务名}: codec: content-encoding: enabled: false
如果用application.properties:
micronaut.http.client.services.${服务名}.codec.content-encoding.enabled=false
方式二:自定义HttpClient(配置文件不生效时用)
手动创建HttpClient实例,移除Netty的GzipDecoder:
import io.micronaut.http.client.HttpClient; import io.micronaut.http.client.netty.NettyHttpClientBuilder; import io.netty.handler.codec.http.HttpContentDecoder; import jakarta.inject.Singleton; @Singleton public class CustomHttpClientProvider { public HttpClient createNoGzipClient() { return new NettyHttpClientBuilder() .configureChannel(channel -> { // 移除所有HTTP内容解码器(包含GzipDecoder) channel.pipeline().remove(HttpContentDecoder.class); }) .build(); } }
然后在客户端接口中指定使用这个自定义HttpClient:
@Client(value = "${url}", httpClient = "customHttpClientProvider") @Header(name = "Accept-encoding", value = "gzip") public interface MyClient { @Get(value = "/data-feeds", consumes = MediaType.APPLICATION_OCTET_STREAM) Flux<ByteBuffer> getData(@QueryValue("from") String from, @QueryValue("minute") String minute, @QueryValue("filters") List<Object> filters); }
2. 调整客户端接口返回类型
将返回类型设为Flux<ByteBuffer>(或Flux<byte[]>),同时consumes指定为APPLICATION_OCTET_STREAM,确保框架不会尝试解析压缩后的JSON内容,直接返回原始字节流。
3. 优化S3分片上传流程
现在拿到的是服务器返回的原始gzip压缩字节,不需要再用GZIPOutputStream重新压缩,直接攒够分片大小后上传:
import software.amazon.awssdk.services.s3.S3AsyncClient; import software.amazon.awssdk.services.s3.model.UploadPartRequest; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; import java.nio.ByteBuffer; import java.io.ByteArrayOutputStream; import java.util.Collections; import java.util.concurrent.atomic.AtomicInteger; // 假设已初始化S3AsyncClient实例 S3AsyncClient s3Client = ...; String bucketName = "your-backup-bucket"; String s3Key = "daily-data/2022-01-01.gz"; // 先初始化分片上传,获取uploadId String uploadId = s3Client.createMultipartUpload(req -> req .bucket(bucketName) .key(s3Key) ).join().uploadId(); AtomicInteger partCounter = new AtomicInteger(1); Flux.range(0, 1439) .flatMapSequential(minute -> client.getData("2022-01-01", String.valueOf(minute), Collections.emptyList()), 20) // 保持20并发,可根据服务器限流/带宽调整 .bufferUntil(byteBuffers -> { // 累计字节数达到5Mb时触发分片 long totalSize = byteBuffers.stream().mapToLong(ByteBuffer::remaining).sum(); return totalSize >= 5 * 1024 * 1024; }) .flatMap(buffer -> { // 将ByteBuffer列表合并为单个字节数组 byte[] partData = buffer.stream() .map(ByteBuffer::array) .collect(ByteArrayOutputStream::new, ByteArrayOutputStream::write, ByteArrayOutputStream::writeTo) .toByteArray(); // 构建分片上传请求 UploadPartRequest partReq = UploadPartRequest.builder() .bucket(bucketName) .key(s3Key) .uploadId(uploadId) .partNumber(partCounter.getAndIncrement()) .contentLength((long) partData.length) .build(); // 异步上传分片 return Mono.fromFuture(s3Client.uploadPart(partReq, AsyncRequestBody.fromBytes(partData))); }) .then() .doOnSuccess(v -> { // 完成分片上传(需收集所有分片的ETag,此处省略收集逻辑) s3Client.completeMultipartUpload(req -> req .bucket(bucketName) .key(s3Key) .uploadId(uploadId) .multipartUpload(mpu -> mpu.parts(/* 传入所有分片的ETag和编号 */)) ).join(); }) .subscribe();
额外性能优化
- 移除不必要的线程池操作:现在直接处理原始压缩字节,无需在
boundedElastic()中做压缩映射,减少线程切换开销 - 动态调整并发数:根据服务器限流规则和自身网络带宽,调整
flatMapSequential的并发参数(比如20~50),平衡吞吐量和服务端压力
内容的提问来源于stack exchange,提问作者Stephane Paulus
相关产品推荐
相关产品推荐

