Java SocketInflater/Deflater解压报错及多消息适配疑问
压缩Socket通信中的Inflater异常与多消息处理问题
问题描述
- 服务端先收到仅含zlib头部(
789c)的数据包,无法完成解压;后续收到完整payload后继续解压时,抛出java.util.zip.DataFormatException: invalid stored block lengths错误。 - 测试
InflaterInputStream时发现,它无法处理包含多个压缩消息的输入流:内部Inflater完成一个消息解压后,后续读取直接返回-1,无法处理下一个压缩消息,询问是否有适配方法,还是必须直接使用Inflater处理。
核心代码
服务端代码
public void listen() { try { InputStream inputStream = clientSocket.getInputStream(); byte[] readBuff = new byte[1024]; ByteArrayOutputStream readOutputStream = new ByteArrayOutputStream(); Inflater inflater = new Inflater(); byte[] inflateBuff = new byte[1024]; int readCount; int inflateCount = 0; while ((readCount = inputStream.read(readBuff)) != -1) { log.debug("readCount: {}", readCount); readOutputStream.write(readBuff, 0, readCount); log.debug("readOutputStream byte array: {}", hexFormat.formatHex(readOutputStream.toByteArray())); inflater.setInput(readOutputStream.toByteArray()); while (inflater.getRemaining() > 0 && !inflater.finished() && !inflater.needsInput()) { log.debug("inflater not finished and doesn't needs input, calling inflater.inflate"); inflateCount += inflater.inflate(inflateBuff); } if (inflater.finished()) { log.debug("inflater finished, handling message"); handleMessage(Arrays.copyOfRange(inflateBuff, 0, inflateCount)); inflater.reset(); inflateCount = 0; readOutputStream.close(); readOutputStream = new ByteArrayOutputStream(); } } } catch (IOException | DataFormatException e) { log.error("MessageListener error", e); throw new RuntimeException(e); } }
客户端代码
@Override public void run(ApplicationArguments args) throws Exception { try (Socket socket = new Socket("localhost", 4943)) { OutputStream socketOutputStream = socket.getOutputStream(); for (int i = 0; i < 1000; i++) { StringBuilder message = new StringBuilder("Hello ").append(i); writeMessage(socketOutputStream, message.toString().getBytes(StandardCharsets.UTF_8)); } } } private void writeMessage(OutputStream socketOutputStream, byte[] message) { log.debug("Writing {}", message); try { DeflaterOutputStream deflaterOutputStream = new DeflaterOutputStream(socketOutputStream); deflaterOutputStream.write(message); deflaterOutputStream.finish(); deflaterOutputStream.flush(); } catch (IOException e) { throw new UncheckedIOException(e); } }
问题分析与解决方案
1. 服务端Inflater异常修复
问题根源
你的服务端代码每次读取新数据后,都会调用inflater.setInput(readOutputStream.toByteArray()),这会重复传入已经处理过的字节数据。比如第一次收到789c头部后,Inflater已加载该数据但未完成解压;第二次收到完整payload时,代码将789c+新数据重新传入Inflater,导致Inflater重复解析头部,触发格式错误。
修复代码
public void listen() { try { InputStream inputStream = clientSocket.getInputStream(); byte[] readBuff = new byte[1024]; ByteArrayOutputStream pendingBytes = new ByteArrayOutputStream(); Inflater inflater = new Inflater(); byte[] inflateBuff = new byte[1024]; int inflateCount = 0; int readCount; while ((readCount = inputStream.read(readBuff)) != -1) { pendingBytes.write(readBuff, 0, readCount); byte[] pendingData = pendingBytes.toByteArray(); inflater.setInput(pendingData); // 尝试解压直到需要更多输入或完成 while (!inflater.finished() && !inflater.needsInput()) { int len = inflater.inflate(inflateBuff, inflateCount, inflateBuff.length - inflateCount); if (len == 0) break; inflateCount += len; } if (inflater.finished()) { // 处理完整消息 handleMessage(Arrays.copyOfRange(inflateBuff, 0, inflateCount)); // 重置Inflater状态 inflater.reset(); inflateCount = 0; // 保留未处理的剩余字节(如果存在) int remaining = inflater.getRemaining(); pendingBytes.reset(); if (remaining > 0) { pendingBytes.write(pendingData, pendingData.length - remaining, remaining); } } } } catch (IOException | DataFormatException e) { log.error("MessageListener error", e); throw new RuntimeException(e); } }
2. InflaterInputStream处理多压缩消息的方案
InflaterInputStream的设计目标是处理单个连续压缩流,当内部Inflater完成解压后,它会判定流已结束,后续读取直接返回-1,这是固有行为。若要处理多个独立压缩消息,有两种可行方案:
方案1:自定义支持多消息的InputStream
自己封装InputStream,每次处理完一个压缩消息后,重置内部Inflater,继续读取下一个压缩块。
方案2:给压缩消息添加长度前缀(推荐)
客户端发送消息时,先发送压缩数据的长度(固定4字节int),服务端先读取长度,再读取对应长度的压缩数据,最后用InflaterInputStream处理该固定长度的字节流,实现逐个处理独立压缩消息。
客户端修改代码
private void writeMessage(OutputStream socketOutputStream, byte[] message) throws IOException { // 压缩消息 ByteArrayOutputStream baos = new ByteArrayOutputStream(); try (DeflaterOutputStream dos = new DeflaterOutputStream(baos)) { dos.write(message); dos.finish(); } byte[] compressed = baos.toByteArray(); // 发送长度前缀 socketOutputStream.write(ByteBuffer.allocate(4).putInt(compressed.length).array()); // 发送压缩数据 socketOutputStream.write(compressed); socketOutputStream.flush(); }
服务端修改代码
public void listen() throws IOException { InputStream inputStream = clientSocket.getInputStream(); DataInputStream dis = new DataInputStream(inputStream); while (true) { try { // 读取压缩数据长度 int len = dis.readInt(); byte[] compressed = new byte[len]; dis.readFully(compressed); // 解压处理 ByteArrayInputStream bais = new ByteArrayInputStream(compressed); try (InflaterInputStream iis = new InflaterInputStream(bais)) { ByteArrayOutputStream baos = new ByteArrayOutputStream(); byte[] buff = new byte[1024]; int read; while ((read = iis.read(buff)) != -1) { baos.write(buff, 0, read); } handleMessage(baos.toByteArray()); } } catch (EOFException e) { // 连接关闭,退出循环 break; } catch (DataFormatException e) { log.error("解压失败", e); throw new RuntimeException(e); } } }
内容的提问来源于stack exchange,提问作者gadams00
相关产品推荐
相关产品推荐

