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

DataInputStream.read()读取大数组仅获部分字节的原因及解决办法

问题

我正在开发一款软件,用于以分块(CHUNKED)、压缩、原始三种格式传输游戏数据片段,并测试每种格式的传输速度。但在调试分块(CHUNKED)传输模式时遇到问题:当待读取的字节数组大小超过140000字节时,客户端使用DataInputStream.read()方法最多只能读取到约131072字节,无论数组实际容量多大。我使用的是SDK 17,调整字节数量后确认了该阈值,甚至在一个无冗余代码的测试客户端/服务器项目中复现了问题,且阈值更低。

服务器代码

/**
 *
 * @return Time taken to complete transfer.
 */
public int start(String mode, int length) throws IOException, InterruptedException {
    if(mode.equals("RAW")){
        byte[] all = new ByteCollector(ServerMain.FILES, length).collect();
        output.writeUTF("SENDING " + mode + " " + all.length);
        expect("RECEIVING " + mode);
        long start = System.currentTimeMillis();
        echoSend(all);
        return (int) (System.currentTimeMillis() - start);
    }else if(mode.equals("CHUNKED")){ /*the important part*/
        //split into chunks
        byte[] all = new ByteCollector(ServerMain.FILES, length).collect();
        int chunks = maxChunks(all);
        output.writeUTF("SENDING " + mode + " " + chunks);
        System.out.println("Expecting RECEIVING " + chunks + "...");
        expect("RECEIVING " + chunks);
        int ms = 0;
        for(int i = 0; i<chunks; i++){
            byte[] currentChunk = getChunk(i, all);
            System.out.println("My chunk length is " + currentChunk.length);
            long start = System.currentTimeMillis();
            System.out.println("Sending...");
            echoSend(currentChunk);
            ms += System.currentTimeMillis() - start;
        }
        if(chunks == 0) expect("0"); //still need to confirm, even though no data was sent
        return ms;
    }else if(mode.equals("COMPRESSED")){
        byte[] compressed = new ByteCollector(ServerMain.FILES, length).collect();
        compressed = ExperimentUtils.compress(compressed);
        output.writeUTF("SENDING " + mode + " " + compressed.length);
        expect("RECEIVING " + mode);
        long start = System.currentTimeMillis();
        echoSend(compressed, length);
        return (int) (System.currentTimeMillis() - start);
    }
    return -1;
}
public static void main(String[] args) throws IOException,InterruptedException{
    FILES = Files.walk(Paths.get(DIRECTORY)).filter(Files::isRegularFile).toArray(Path[]::new);
    SyncServer server = new SyncServer(new ServerSocket(12222).accept());
    System.out.println("--------[CH UNK ED]--------");
    short[] chunkedSpeeds = new short[FOLDER_SIZE_MB + 1/*for "zero" or origin*/];
    for(int i = 0; i<=FOLDER_SIZE_MB; i++){
        chunkedSpeeds[i] = (short) server.start("CHUNKED", i * MB);
        System.out.println(i + "MB, Chunked: " + chunkedSpeeds[i]);
    }
    short[] compressedSpeeds = new short[FOLDER_SIZE_MB + 1];
    for(int i = 0; i<=FOLDER_SIZE_MB; i++){
        compressedSpeeds[i] = (short) server.start("COMPRESSED", i * MB);
    }
    short[] rawSpeeds = new short[FOLDER_SIZE_MB + 1];
    for(int i = 0; i<=FOLDER_SIZE_MB; i++){
        rawSpeeds[i] = (short) server.start("RAW", i * MB);
    }
    System.out.println("Raw speeds: " + Arrays.toString(rawSpeeds));
    System.out.println("\n\nCompressed speeds: " + Arrays.toString(compressedSpeeds));
    System.out.println("\n\nChunked speeds: " + Arrays.toString(chunkedSpeeds));
}

客户端代码

public static void main(String[] args) throws IOException, InterruptedException {
    Socket socket = new Socket("localhost", 12222);
    DataInputStream input = new DataInputStream(socket.getInputStream());
    DataOutputStream output = new DataOutputStream(socket.getOutputStream());
    while(socket.isConnected()){
        String response = input.readUTF();
        String[] content = response.split(" ");
        if(response.startsWith("SENDING CHUNKED")){
            int chunks = Integer.parseInt(content[2]);
            System.out.println("Read chunk amount of " + chunks);
            output.writeUTF("RECEIVING " + chunks);
            for(int i = 0; i<chunks; i++){
                byte[] chunk = new byte[32 * MB];
                System.out.println("Ready to receive...");
                int read = input.read(chunk);
                System.out.println("Echoing read length of " + read);
                output.writeUTF(String.valueOf(read));
            }
            if(chunks == 0) output.writeUTF("0");
        }else if(response.startsWith("SENDING COMPRESSED")){
            byte[] compressed = new byte[Integer.parseInt(content[2])];
            output.writeUTF("RECEIVING " + compressed.length);
            input.read(compressed);
            decompress(compressed);
            output.writeInt(decompress(compressed).length);
        }else if(response.startsWith("SENDING RAW")){
            int length = Integer.parseInt(content[2]);
            output.writeUTF("RECEIVING " + length);
            byte[] received = new byte[length];
            input.read(received);
            output.writeInt(received.length);
        }
    }
}
public static byte[] decompress(byte[] in) throws IOException {
    try {
        ByteArrayOutputStream out = new ByteArrayOutputStream();
        InflaterOutputStream infl = new InflaterOutputStream(out);
        infl.write(in);
        infl.flush();
        infl.close();

        return out.toByteArray();
    } catch (Exception e) {
        System.out.println("Error decompressing byte array with length " + in.length);
        throw e;
    }
}

问题原因

  1. DataInputStream.read(byte[])的特性限制:该方法不会一次性读取所有预期数据,它仅返回当前网络接收缓冲区中可用的字节数。默认情况下,Socket接收缓冲区大小为128KB(131072字节),这就是你看到的阈值来源。当发送数据超过缓冲区容量时,剩余数据会滞留在TCP发送队列中,需要多次读取才能获取完整数据。
  2. 分块传输协议缺失关键信息:服务器未在发送每个chunk前告知客户端该chunk的具体长度,客户端无法判断需要读取多少字节才算完成一个chunk的接收。
  3. 单次读取的逻辑错误:客户端仅调用一次read()就认定chunk接收完成,没有处理“数据未读完”的情况,导致只获取了缓冲区中的部分数据。

优化实现方案

1. 服务器端:添加chunk长度前置通知

修改分块传输逻辑,在发送chunk字节数据前,先传递该chunk的长度,确保客户端明确知道需要读取的字节数:

// 分块传输循环内的修改
for(int i = 0; i<chunks; i++){
    byte[] currentChunk = getChunk(i, all);
    System.out.println("My chunk length is " + currentChunk.length);
    long start = System.currentTimeMillis();
    // 先发送chunk长度
    output.writeInt(currentChunk.length);
    output.flush(); // 强制刷新缓冲区,确保长度数据优先发送
    echoSend(currentChunk);
    ms += System.currentTimeMillis() - start;
    // 等待客户端确认该chunk接收完成
    String ack = input.readUTF();
    System.out.println("Chunk " + i + " ack received: " + ack);
}

2. 客户端:循环读取直到完成整个chunk

客户端先读取chunk长度,再循环调用read(),直到读取的总字节数等于chunk长度:

// 分块接收循环内的修改
for(int i = 0; i<chunks; i++){
    // 读取当前chunk的长度
    int chunkLength = input.readInt();
    byte[] chunk = new byte[chunkLength];
    int totalRead = 0;
    System.out.println("Ready to receive chunk of length: " + chunkLength);
    // 循环读取,直到获取完整chunk
    while(totalRead < chunkLength){
        int read = input.read(chunk, totalRead, chunkLength - totalRead);
        if(read == -1){
            throw new IOException("Socket closed before full chunk received");
        }
        totalRead += read;
    }
    System.out.println("Received full chunk, total bytes: " + totalRead);
    output.writeUTF("CHUNK_" + i + "_OK");
}

3. 通用优化点

  • 调用output.flush():在服务器发送数据后强制刷新缓冲区,避免Nagle算法导致的小数据块延迟。
  • 异常处理:添加Socket断开、数据不完整等异常的捕获逻辑,提升程序稳定性。

4. 进阶替代方案

如果追求更高的传输效率,可以改用Java NIO(SocketChannel)或Netty框架,它们提供异步IO、更灵活的缓冲区管理,更适合游戏数据这类对传输性能要求较高的场景。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.09 20:30:27