使用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))
问题原因分析
- DataBuffer未正确释放:WebFlux的
DataBuffer基于内存池实现,如果使用后不手动释放,会导致内存泄漏,大文件上传时会引发解码器状态异常。 - 流未及时终止:当
PartEvent.isLast()为true时,代表当前是文件的最后一个分片,但你的代码没有在此时终止流,导致解码器继续等待后续数据,最终引发连接关闭异常。 - 异常处理不规范:在
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 + ")"; } } }
额外优化建议
- 配置内存缓冲区大小:在
application.properties中添加以下配置,避免大文件上传时内存溢出:spring.codec.max-in-memory-size=10MB - 避免依赖Content-Length:大文件上传可能使用分块编码,此时
Content-Length为-1,建议通过PartEvent的元数据做兼容处理。 - 直接写入文件系统:对于大文件,建议直接将
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
相关产品推荐
相关产品推荐

