You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.22 13:36:29