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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.10 01:20:38