Spring Boot异步导出CSV时PipedInputStream阻塞问题排查
Spring Boot CSV导出阻塞问题排查与解决
问题2:为何加@Async仍阻塞?关于pin返回的理解是否正确?
你的理解不正确。当前代码里,就算给exportCsv加了@Async,如果注解没生效,写入操作会在main线程同步执行——这时候你还没返回ResponseEntity,客户端根本没开始读取PipedInputStream,当PipedOutputStream写入的数据填满默认缓冲区(1024字节)后,就会触发awaitSpace阻塞,直到有数据被读取腾出空间。
@Async没生效的常见原因:
exportCsv方法不是public修饰:@Async需要方法是public才能被Spring代理拦截,你写的是void exportCsv(...),缺少public修饰符。- 调用方和
exportCsv在同一个类:如果端点类和服务类是同一个类,Spring动态代理无法拦截内部方法调用,@Async直接失效,方法还是同步运行。 - 服务类没被Spring管理:比如没加
@Service或@Component注解,导致@Async不生效。
就算@Async真生效了,也可能因为客户端读取速度远慢于写入速度,导致缓冲区填满,写入线程阻塞,但此时阻塞的是异步线程,不会卡住main线程和整个应用。但你的日志显示main线程在阻塞,说明@Async肯定没生效,写入操作完全是在main线程同步执行的。
问题1:如何将exportCsv操作纳入try-catch处理?
首先得修正exportCsv的方法签名确保@Async能生效,然后在方法内部完善异常处理,同时调用方也要监听异步任务的异常:
服务类方法修改
@Async public void exportCsv(OutputStream outputStream, Long id) { BeanWriter writer = null; try { writer = StreamFactory.newInstance().createWriter("output", new OutputStreamWriter(outputStream)); AtomicInteger index = new AtomicInteger(); try (Stream<Data> stream = repository.streamData(id)) { stream.filter(Objects::nonNull).forEach(e -> { writer.write(e); int i = index.incrementAndGet(); if (i % FLUSH_SIZE == 0) { writer.flush(); // 手动调用System.gc()没必要,JVM会自动处理垃圾回收 } }); } writer.flush(); } catch (Exception e) { // 处理导出异常:打印日志、关闭流等 log.error("CSV导出失败", e); try { outputStream.close(); } catch (IOException ex) { log.error("关闭输出流失败", ex); } } finally { if (writer != null) { try { writer.close(); } catch (Exception e) { log.error("关闭BeanWriter失败", e); } } } }
端点类调用优化(支持异常监听)
如果需要在调用方感知异步任务的异常,可以把服务类方法改成返回CompletableFuture:
// 服务类修改后 @Async public CompletableFuture<Void> exportCsv(OutputStream outputStream, Long id) { try { // 上述写入逻辑 return CompletableFuture.completedFuture(null); } catch (Exception e) { log.error("CSV导出失败", e); // 异常时关闭流 try { outputStream.close(); } catch (IOException ex) { log.error("关闭输出流失败", ex); } return CompletableFuture.failedFuture(e); } } // 端点类代码 public ResponseEntity<InputStreamResource> exportCsv(Long id) { PipedInputStream pin = new PipedInputStream(); PipedOutputStream pout = new PipedOutputStream(pin); // 调用异步方法并监听异常 service.exportCsv(pout, id).exceptionally(ex -> { log.error("异步导出CSV失败", ex); try { pout.close(); pin.close(); } catch (IOException ioEx) { log.error("关闭流失败", ioEx); } return null; }); return ResponseEntity.ok() .headers(headers) .contentType(mediaType) .body(new InputStreamResource(pin)); }
另外要记住:PipedInputStream和PipedOutputStream必须在不同线程读写,否则必然阻塞。所以一定要确保@Async生效,让写入操作在独立线程执行,main线程尽快返回InputStreamResource给客户端,让客户端开始读取数据,这样读写形成流,就不会阻塞了。
原问题代码与日志
端点类代码
public ResponseEntity<InputStreamResource> exportCsv(Long id) { PipedInputStream pin = new PipedInputStream(); PipedOutputStream pout = new PipedOutputStream(pin); service.exportCsv(pout, id); return ResponseEntity.ok().headers(headers).contentType(mediaType).body(new InputStreamResource(pin)); }
服务类代码
@Async void exportCsv(OutputStream outputStream, Long id) { BeanWriter writer = StreamFactory.newInstance().createWriter("output", new OutputStreamWriter(outputStream)); AtomicInteger index = new AtomicInteger(); try (Stream<Data> stream = repository.streamData(id)) { stream.filter(Objects::nonNull).forEach(e -> { writer.write(e); int i = index.incrementAndGet(); if (i % FLUSH_SIZE == 0) {writer.flush(); System.gc();} }); } writer.flush(); writer.close(); }
日志信息
"main" prio=10 tid=0x08066000 nid=0x48d2 in Object.wait() [0xb7fd2000..0xb7fd31e8] java.lang.Thread.State: TIMED_WAITING (on object monitor) at java.lang.Object.wait(Native Method) - waiting on <0xa5c28be8> (a java.io.PipedInputStream) at java.io.PipedInputStream.awaitSpace(PipedInputStream.java:257) at java.io.PipedInputStream.receive(PipedInputStream.java:215)
内容的提问来源于stack exchange,提问作者Luke
相关产品推荐
相关产品推荐

