如何基于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头,或者你需要对原始压缩流做自定义处理,可以通过**流式转换FluxFlux<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); // 必须释放缓冲区,避免内存泄漏 } });
关键注意事项
- 线程池选择:必须用
Schedulers.boundedElastic()处理GZIPInputStream的阻塞操作,不能用Netty的IO线程,否则会拖垮整个响应式管道。 - 资源释放:所有
DataBuffer必须手动释放(或使用asInputStream(true)自动释放),否则会造成内存泄漏。 - 内存控制:通过固定大小的字节缓冲区(8KB)限制单次解压的数据量,确保内存占用始终在可控范围内。
你提到的Reactor Netty issue-251和Spring Integration issue-2300,主要是早期版本中流式压缩的适配问题,现在新版本的框架已经支持原生自动处理,但如果需要自定义流处理,上面的方案可以完美解决。
内容的提问来源于stack exchange,提问作者Abhijit Sarkar
相关产品推荐
相关产品推荐

