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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 15:32:40