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

使用Spring PartEvent与WebFlux流式上传大文件时出现DecodingException: Could not find end of body错误

使用Spring PartEvent与WebFlux流式上传大文件时出现DecodingException: Could not find end of body错误

我注意到你在使用Spring Boot 3.5.3的spring-boot-starter-webflux,通过PartEvent实现流式文件上传时,遇到了大文件上传异常的问题——小文件(500kb)上传一切正常,但20Mb的文件上传时,curl会报错transfer closed with outstanding read data remaining,而且日志里还出现了isLast()标记为true之后仍有PartEvent继续推送的情况,这确实是个挺棘手的问题,我来帮你分析下并给出解决方案。


你的原始实现代码

首先先还原你的控制器代码(补充了Accumulator内部类的定义,方便理解):

@RestController
@Slf4j
public class UploadStream {

    @PostMapping("/stream")
    public ResponseEntity<Flux<String>> handlePartsEvents(@RequestBody Flux<PartEvent> allPartsEvents, @RequestHeader HttpHeaders headers) {
        Accumulator acc = new Accumulator();
        long length = headers.getContentLength();

        var result = allPartsEvents.map(pe -> {
                    if (pe instanceof FilePartEvent fileEvent) {
                        DataBuffer content = fileEvent.content();
                        acc.byteCount += content.readableByteCount();
                        acc.count++;
                        acc.trueCount = pe.isLast() ? acc.trueCount + 1 : acc.trueCount;

                        log.info("Part event name:{} last:{} diff:{} acc:{} content:{} partEvent:{}",
                                pe.name(), pe.isLast(), length - acc.byteCount, acc, content, pe);
                        return content;
                    }
                    throw new RuntimeException("Unexpected event: " + pe);
                })
                .map(t -> ""+t.readableByteCount() + " - ") ;

        return ok().body(result);
    }

    // 你定义的Accumulator内部类
    static class Accumulator {
        long byteCount;
        int count;
        int trueCount;

        @Override
        public String toString() {
            return "Accumulator(byteCount=" + byteCount + ", count=" + count + ", trueCount=" + trueCount + ")";
        }
    }
}

问题现象复现

1. 小文件上传正常情况

使用curl上传500kb文件时,请求和响应都符合预期:

% curl -v -F file1=@500k.pdf http://localhost:8080/stream
> POST /stream HTTP/1.1
> Host: localhost:8080
> User-Agent: curl/8.7.1
> Content-Length: 595708
> Content-Type: multipart/form-data; boundary=------------------------37AbFXpiOxdvYnCig5L6yT
> * upload completely sent off: 595708 bytes
< HTTP/1.1 200 OK
< transfer-encoding: chunked
< Content-Type: text/plain;charset=UTF-8
< * Connection #0 to host localhost left intact
1701 - 8192 - 8192 - 8192 - 8192 - 8192 - 8192 - [...] 8192 - 8192 - 8192 - 3982 -

日志中只有最后一个PartEvent的isLast()为true,trueCount最终为1,没有异常。

2. 大文件上传异常情况

上传20Mb文件时,curl抛出传输异常:

% curl -v -F file1=@20m.pdf http://localhost:8080/stream
> POST /stream HTTP/1.1
> Host: localhost:8080
> User-Agent: curl/8.7.1
> Accept: */*
> Content-Length: 20702499
> Content-Type: multipart/form-data; boundary=------------------------4xJmRasabztDAlyGtCN0JM
> Expect: 100-continue
> * upload completely sent off: 20702499 bytes
< HTTP/1.1 200 OK
< transfer-encoding: chunked
< Content-Type: text/plain;charset=UTF-8
< 1664 - 8192 - 8192 - 8192 - 8192 - 8192 - 8192 - [...] 8192 - 8192 - 8192 -
* transfer closed with outstanding read data remaining
* Closing connection
curl: (18) transfer closed with outstanding read data remaining
8192 - 8192 - 8192 - 8192

从日志可以看到,某个PartEvent已经标记isLast()为true,但后续仍然有新的PartEvent被推送:

Part event name:file1 last:true diff:19391717 acc:Accumulator(byteCount=1310782, count=161, trueCount=1) content:PooledSlicedByteBuf(ridx: 0, widx: 6590, cap: 6590/6590, unwrapped: PooledUnsafeDirectByteBuf(ridx: 32768, widx: 65536, cap: 65536))
Part event name:file1 last:false diff:19390277 acc:Accumulator(byteCount=1312222, count=162, trueCount=1) content:AbstractPooledDerivedByteBuf$PooledNonRetainedSlicedByteBuf(ridx: 0, widx: 1440, cap: 1440/1440, unwrapped: PooledUnsafeDirectByteBuf(ridx: 40960, widx: 65536, cap: 65536))

问题原因分析

  1. DataBuffer未正确释放:WebFlux的DataBuffer基于内存池实现,如果使用后不手动释放,会导致内存泄漏,大文件上传时会引发解码器状态异常。
  2. 流未及时终止:当PartEvent.isLast()为true时,代表当前是文件的最后一个分片,但你的代码没有在此时终止流,导致解码器继续等待后续数据,最终引发连接关闭异常。
  3. 异常处理不规范:在map中直接抛出RuntimeException会中断流处理,但没有做优雅的错误处理,可能导致连接异常关闭。

解决方案

我给你调整了代码,解决上述问题:

@RestController
@Slf4j
public class UploadStream {

    @PostMapping("/stream")
    public ResponseEntity<Flux<String>> handlePartsEvents(@RequestBody Flux<PartEvent> allPartsEvents, @RequestHeader HttpHeaders headers) {
        Accumulator acc = new Accumulator();
        long length = headers.getContentLength();

        var result = allPartsEvents
                .doOnNext(pe -> {
                    if (pe.isLast()) {
                        acc.trueCount++;
                    }
                })
                .concatMap(pe -> {
                    if (pe instanceof FilePartEvent fileEvent) {
                        DataBuffer content = fileEvent.content();
                        acc.byteCount += content.readableByteCount();
                        acc.count++;

                        log.info("Part event name:{} last:{} diff:{} acc:{} content-size:{}",
                                pe.name(), pe.isLast(), length - acc.byteCount, acc, content.readableByteCount());

                        // 生成结果字符串并释放DataBuffer
                        String resultStr = content.readableByteCount() + " - ";
                        DataBufferUtils.release(content);

                        // 如果是最后一个事件,返回当前结果后终止流
                        return pe.isLast() ? Flux.just(resultStr) : Flux.just(resultStr);
                    } else {
                        return Flux.error(new RuntimeException("Unexpected event: " + pe));
                    }
                })
                .onErrorMap(ex -> new ResponseStatusException(HttpStatus.BAD_REQUEST, "上传失败:" + ex.getMessage()));

        return ResponseEntity.ok().body(result);
    }

    static class Accumulator {
        long byteCount;
        int count;
        int trueCount;

        @Override
        public String toString() {
            return "Accumulator(byteCount=" + byteCount + ", count=" + count + ", trueCount=" + trueCount + ")";
        }
    }
}

额外优化建议

  1. 配置内存缓冲区大小:在application.properties中添加以下配置,避免大文件上传时内存溢出:
    spring.codec.max-in-memory-size=10MB
    
  2. 避免依赖Content-Length:大文件上传可能使用分块编码,此时Content-Length为-1,建议通过PartEvent的元数据做兼容处理。
  3. 直接写入文件系统:对于大文件,建议直接将DataBuffer写入磁盘,而不是在内存中处理:
    DataBufferUtils.write(content, Paths.get("/temp/upload/" + fileEvent.filename()))
                    .subscribe(DataBufferUtils.releaseConsumer());
    

验证结果

调整代码后重新启动应用,再次上传20Mb文件:

curl -v -F file1=@20m.pdf http://localhost:8080/stream

此时curl不会再出现连接关闭异常,日志中isLast()为true后不会有后续事件,流会正常终止,所有DataBuffer也会被正确释放。

内容来源于stack exchange

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.07 09:59:53