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

Java 11 HTTP服务器管道流解压时出现Pipe closed异常原因排查

解决Java管道流处理Multipart压缩数据时的IOException: Pipe closed异常

问题场景

在Java 11小型HTTP服务器中处理包含压缩数据的multipart流请求时,使用管道流实现边读边解压逻辑,解压后的文件内容正常,但偶尔会抛出IOException: Pipe closed异常,堆栈信息如下:

[java] 27016 [HTTP-Dispatcher] ERROR c.example.AhFormatter - e9137279-7c15-41a3-9783-eb029a975767 - Error while processing finish request 
[java] java.io.IOException: Pipe closed
[java]  at java.base/java.io.PipedInputStream.checkStateForReceive(PipedInputStream.java:260)
[java]  at java.base/java.io.PipedInputStream.receive(PipedInputStream.java:226)
[java]  at java.base/java.io.PipedOutputStream.write(PipedOutputStream.java:149)
[java]  at java.base/java.io.InputStream.transferTo(InputStream.java:705)
[java]  at com.example.upload.MultipartStream.readBodyData(MultipartStream.java:469)
[java]  at com.example.AhFormatter.processPdfData(AhFormatter.java:291)

相关核心代码如下:

处理Multipart流的方法

private void processPdfData(MultipartStream multipartStream) throws IOException, InterruptedException {
    String headers = multipartStream.readHeaders();
    if (!headers.contains("name=\"pdf\"")) {
        logger.warn("{} - Header with name pdf not found", taskId);
        throw new IllegalStateException();
    }

    PipedInputStream pipedInputStream = new PipedInputStream();
    PipedOutputStream pipedOutputStream = new PipedOutputStream();
    pipedInputStream.connect(pipedOutputStream);

    Thread unzipThread = new Thread(() -> {
        try {
            Zips.unzip(pipedInputStream, targetPath, true);
        } catch (IOException e) {
            logger.error(format("%s - Error during unzip", taskId), e);
            throw new ZipException(e);
        }
    }, "PipedZipStream");
    
    unzipThread.setName("UnzipThread");
    unzipThread.setUncaughtExceptionHandler((thread, throwable) -> {
        logger.error(format("%s - Uncaught exception while unzipping", taskId), throwable);
        server.stop(0);
        stopServer.set(true);
    });
    
    unzipThread.start();
    multipartStream.readBodyData(pipedOutputStream);
    unzipThread.join();

    pipedOutputStream.close();
    pipedInputStream.close();
}

解压方法

public static void unzip(InputStream zip, Path targetDirectory, boolean close) throws IOException {
    ZipInputStream zipStream = new ZipInputStream(zip);
    try {
        ZipEntry zipEntry;
        while ((zipEntry = zipStream.getNextEntry()) != null) {
            File targetFile = guardAgainstZipSlip(targetDirectory, zipEntry);

            if (zipEntry.isDirectory()) {
                if (!targetFile.isDirectory() && !targetFile.mkdirs()) {
                    throw new IOException("Failed to create directory " + targetFile);
                }
                continue;
            }

            // fix for Windows-created archives
            File parent = targetFile.getParentFile();
            if (!parent.isDirectory() && !parent.mkdirs()) {
                throw new IOException("Failed to create directory " + parent);
            }

            copy(zipStream, targetFile.toPath(), StandardCopyOption.REPLACE_EXISTING);
        }
    } finally {
        if (close) {
            zipStream.close();
        }
    }
}

private static File guardAgainstZipSlip(Path destinationDir, ZipEntry zipEntry) throws IOException {
    File targetFile = new File(destinationDir.toFile(), zipEntry.getName());
    if (!targetFile.getCanonicalPath().startsWith(destinationDir.toFile().getCanonicalPath() + File.separator)) {
        throw new IOException("Entry is outside of the target dir: " + zipEntry.getName());
    }
    return targetFile;
}

异常原因分析

这个异常的核心是管道流的关闭顺序错误:

  1. 解压线程执行Zips.unzip时,因为传入了close=true,在finally块中会关闭ZipInputStream,而ZipInputStream底层关联的是PipedInputStream,这会直接导致PipedInputStream被关闭。
  2. 主线程此时可能还在执行multipartStream.readBodyData(pipedOutputStream),尝试往管道输出流写入剩余数据,但管道输入流已经被解压线程关闭,管道连接失效,触发Pipe closed异常。

你提到不显式关流时无异常,是因为此时JVM会在流对象被GC时自动关闭,不会出现线程间的关闭顺序冲突,但这种做法不符合资源管理的规范。

修复方案

方案1:调整流的关闭顺序

修改processPdfData方法,确保主线程写完数据后先关闭管道输出流,再等待解压线程完成,最后关闭管道输入流:

private void processPdfData(MultipartStream multipartStream) throws IOException, InterruptedException {
    String headers = multipartStream.readHeaders();
    if (!headers.contains("name=\"pdf\"")) {
        logger.warn("{} - Header with name pdf not found", taskId);
        throw new IllegalStateException();
    }

    PipedInputStream pipedInputStream = new PipedInputStream();
    PipedOutputStream pipedOutputStream = new PipedOutputStream();
    pipedInputStream.connect(pipedOutputStream);

    Thread unzipThread = new Thread(() -> {
        try {
            Zips.unzip(pipedInputStream, targetPath, true);
        } catch (IOException e) {
            logger.error(format("%s - Error during unzip", taskId), e);
            throw new ZipException(e);
        }
    }, "PipedZipStream");
    
    unzipThread.setName("UnzipThread");
    unzipThread.setUncaughtExceptionHandler((thread, throwable) -> {
        logger.error(format("%s - Uncaught exception while unzipping", taskId), throwable);
        server.stop(0);
        stopServer.set(true);
    });
    
    unzipThread.start();
    // 写入数据到管道输出流
    multipartStream.readBodyData(pipedOutputStream);
    // 写完后关闭输出流,告知解压线程数据已全部写入
    pipedOutputStream.close();
    // 等待解压线程处理完成
    unzipThread.join();
    // 最后关闭输入流
    pipedInputStream.close();
}

方案2:修改解压方法的流关闭逻辑

如果不想调整线程顺序,可以修改Zips.unzip方法,不对管道输入流进行关闭,由主线程统一管理:

public static void unzip(InputStream zip, Path targetDirectory, boolean close) throws IOException {
    ZipInputStream zipStream = new ZipInputStream(zip);
    try {
        ZipEntry zipEntry;
        while ((zipEntry = zipStream.getNextEntry()) != null) {
            File targetFile = guardAgainstZipSlip(targetDirectory, zipEntry);

            if (zipEntry.isDirectory()) {
                if (!targetFile.isDirectory() && !targetFile.mkdirs()) {
                    throw new IOException("Failed to create directory " + targetFile);
                }
                continue;
            }

            // fix for Windows-created archives
            File parent = targetFile.getParentFile();
            if (!parent.isDirectory() && !parent.mkdirs()) {
                throw new IOException("Failed to create directory " + parent);
            }

            copy(zipStream, targetFile.toPath(), StandardCopyOption.REPLACE_EXISTING);
        }
    } finally {
        if (close) {
            // 只关闭ZipInputStream的内部资源,不关闭底层的管道输入流
            zipStream.closeEntry();
            // 如果不是管道输入流,再关闭整个流
            if (!(zip instanceof PipedInputStream)) {
                zipStream.close();
            }
        }
    }
}

同时在调用Zips.unzip时保持close=true即可。

核心原理说明

管道流的线程协作规则:

  • 当管道输出流被关闭时,输入流读取到EOF后会正常结束读取逻辑;
  • 如果管道输入流先被关闭,输出流再尝试写入就会立即触发Pipe closed异常。
    调整关闭顺序后,能保证主线程写完所有数据后才通知解压线程结束,避免了线程间的资源冲突。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 01:07:02