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

如何基于Spring WebClient无中间存储解压缩Flux<DataBuffer>并实现写入?

用Spring WebClient流式处理GZIP压缩流(无中间存储)

针对你提出的「在响应式场景下实时解压缩GZIP流、避免OOM或磁盘存储」的需求,我分两种场景给出解决方案,从最简单的自动处理到自定义流式解压缩都覆盖到:


场景1:服务器返回标准GZIP响应(推荐)

如果服务器正确设置了Content-Encoding: gzip响应头,WebClient可以直接借助Reactor Netty的自动压缩支持完成解压缩,不需要手动处理流,这是最省心的方式:

import org.springframework.web.reactive.function.client.WebClient;
import reactor.netty.http.client.HttpClient;

// 配置WebClient启用自动GZIP解压缩
WebClient webClient = WebClient.builder()
    .clientConnector(new ReactorClientHttpConnector(
        HttpClient.create().compress(true) // 自动识别并处理GZIP/DEFLATE压缩响应
    ))
    .build();

// 直接获取解压缩后的业务对象流
webClient.get()
    .uri("你的接口地址")
    .accept(MediaType.APPLICATION_JSON)
    .retrieve()
    .bodyToFlux(Stuff.class); // 自动完成解压缩 -> 反序列化

这种方式完全基于框架原生能力,不需要自己处理流的细节,内存占用由Reactor Netty的缓冲区机制自动控制,从根源避免OOM。


场景2:手动处理原始GZIP数据流(自定义需求)

如果服务器未设置Content-Encoding头,或者你需要对原始压缩流做自定义处理,可以通过**流式转换Flux**实现无内存堆积的解压缩。核心思路是把响应式的Flux<DataBuffer>适配成InputStream,再结合GZIPInputStream逐段解压,同时确保不会一次性加载所有数据到内存。

步骤1:实现响应式InputStream适配器

这个类会订阅Flux<DataBuffer>并逐段提供数据,避免一次性合并所有缓冲区:

import org.springframework.core.io.buffer.DataBuffer;
import org.springframework.core.io.buffer.DataBufferUtils;
import reactor.core.publisher.Flux;

import java.io.IOException;
import java.io.InputStream;
import java.util.Queue;
import java.util.concurrent.ConcurrentLinkedQueue;
import java.util.concurrent.atomic.AtomicBoolean;

public class ReactiveInputStream extends InputStream {
    private final Queue<byte[]> bufferQueue = new ConcurrentLinkedQueue<>();
    private final AtomicBoolean isComplete = new AtomicBoolean(false);
    private byte[] currentBuffer;
    private int currentPosition;

    public ReactiveInputStream(Flux<DataBuffer> dataBufferFlux) {
        dataBufferFlux.subscribe(
            dataBuffer -> {
                // 读取当前DataBuffer的字节并释放资源
                byte[] bytes = new byte[dataBuffer.readableByteCount()];
                dataBuffer.read(bytes);
                DataBufferUtils.release(dataBuffer);
                bufferQueue.add(bytes);
                synchronized (this) { notifyAll(); }
            },
            error -> {
                synchronized (this) { notifyAll(); }
            },
            () -> {
                isComplete.set(true);
                synchronized (this) { notifyAll(); }
            }
        );
    }

    @Override
    public int read(byte[] b, int off, int len) throws IOException {
        if (len == 0) return 0;
        int bytesRead = 0;

        while (bytesRead < len) {
            // 加载下一段缓冲区数据
            if (currentBuffer == null || currentPosition >= currentBuffer.length) {
                currentBuffer = bufferQueue.poll();
                if (currentBuffer == null) {
                    if (isComplete.get()) return bytesRead == 0 ? -1 : bytesRead;
                    try { synchronized (this) { wait(); } }
                    catch (InterruptedException e) {
                        Thread.currentThread().interrupt();
                        throw new IOException("等待数据时被中断", e);
                    }
                } else currentPosition = 0;
            }

            // 复制当前缓冲区的可用数据
            int available = currentBuffer.length - currentPosition;
            int toRead = Math.min(available, len - bytesRead);
            System.arraycopy(currentBuffer, currentPosition, b, off + bytesRead, toRead);
            currentPosition += toRead;
            bytesRead += toRead;
        }
        return bytesRead;
    }

    @Override
    public int read() throws IOException {
        byte[] singleByte = new byte[1];
        int result = read(singleByte, 0, 1);
        return result == -1 ? -1 : singleByte[0] & 0xFF;
    }
}

步骤2:实现流式解压缩工具类

借助上面的适配器,结合GZIPInputStream完成逐段解压:

import org.springframework.core.io.buffer.DataBuffer;
import org.springframework.core.io.buffer.DefaultDataBufferFactory;
import reactor.core.publisher.Flux;
import reactor.core.scheduler.Schedulers;

import java.io.IOException;
import java.util.zip.GZIPInputStream;

public class GzipStreamUtils {
    private static final DefaultDataBufferFactory BUFFER_FACTORY = new DefaultDataBufferFactory();
    private static final int CHUNK_SIZE = 8192; // 8KB缓冲区,可根据业务调整

    public static Flux<DataBuffer> decompress(Flux<DataBuffer> compressedBuffers) {
        return Flux.create(sink -> {
            // 在弹性线程池处理阻塞的GZIPInputStream,避免阻塞Netty IO线程
            Schedulers.boundedElastic().schedule(() -> {
                try (var reactiveIn = new ReactiveInputStream(compressedBuffers);
                     var gzipIn = new GZIPInputStream(reactiveIn)) {

                    byte[] chunk = new byte[CHUNK_SIZE];
                    int bytesRead;
                    while ((bytesRead = gzipIn.read(chunk)) != -1) {
                        // 将解压后的字节包装成DataBuffer发送到下游
                        DataBuffer buffer = BUFFER_FACTORY.wrap(chunk, 0, bytesRead);
                        sink.next(buffer);
                    }
                    sink.complete();
                } catch (IOException e) {
                    sink.error(e);
                }
            });
        });
    }
}

步骤3:在WebClient中使用

WebClient webClient = WebClient.builder().build();

webClient.get()
    .uri("你的接口地址")
    .header("Accept-Encoding", "gzip") // 告诉服务器返回GZIP压缩数据
    .retrieve()
    .bodyToFlux(DataBuffer.class) // 获取原始压缩数据流
    .transform(GzipStreamUtils::decompress) // 流式解压缩
    .map(dataBuffer -> {
        // 将解压后的DataBuffer转换为业务对象
        try {
            return objectMapper.readValue(dataBuffer.asInputStream(true), Stuff.class);
        } catch (IOException e) {
            throw new RuntimeException("反序列化失败", e);
        } finally {
            DataBufferUtils.release(dataBuffer); // 必须释放缓冲区,避免内存泄漏
        }
    });

关键注意事项

  1. 线程池选择:必须用Schedulers.boundedElastic()处理GZIPInputStream的阻塞操作,不能用Netty的IO线程,否则会拖垮整个响应式管道。
  2. 资源释放:所有DataBuffer必须手动释放(或使用asInputStream(true)自动释放),否则会造成内存泄漏。
  3. 内存控制:通过固定大小的字节缓冲区(8KB)限制单次解压的数据量,确保内存占用始终在可控范围内。

你提到的Reactor Netty issue-251和Spring Integration issue-2300,主要是早期版本中流式压缩的适配问题,现在新版本的框架已经支持原生自动处理,但如果需要自定义流处理,上面的方案可以完美解决。

内容的提问来源于stack exchange,提问作者Abhijit Sarkar

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 06:53:00