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; }
异常原因分析
这个异常的核心是管道流的关闭顺序错误:
- 解压线程执行
Zips.unzip时,因为传入了close=true,在finally块中会关闭ZipInputStream,而ZipInputStream底层关联的是PipedInputStream,这会直接导致PipedInputStream被关闭。 - 主线程此时可能还在执行
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
相关产品推荐
相关产品推荐

