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

WebFlux实现MP3上传S3存元数据到Mongo报Only one connection receive subscriber错误咨询

报错根因

  • 核心原因是Flux<ByteBuffer>类型的请求体被重复订阅:Service层先遍历body的分片,在分片处理逻辑中又把完整的body传给工具方法再次订阅,WebFlux中请求体是仅支持单次订阅的热流,多次订阅直接触发该报错。
  • 次要问题:代码中调用audioRepository.save(audioEntity).block()阻塞方法,破坏WebFlux的非阻塞线程模型,容易引发线程阻塞和流程异常。
  • 逻辑错误:字节转换逻辑是每个分片单独生成MultipartFile,最终拿到的是分片文件而非完整MP3,元数据解析必然失败。
  • 阻塞操作未做线程隔离:mp3agic的文件解析、本地IO、S3上传都是同步阻塞操作,直接在event loop线程执行会拖垮服务性能。

修复方案

1. 重写ByteBuffer转MultipartFile工具方法

先把所有分片聚合成完整字节数组,仅订阅一次请求体:

public static Mono<MultipartFile> byteBufferToMultipartFile(String fileName, Flux<ByteBuffer> body) {
    return body
            // 合并所有ByteBuffer为完整的字节缓存
            .reduce(ByteBuffer.allocate(0), (total, current) -> {
                ByteBuffer merged = ByteBuffer.allocate(total.remaining() + current.remaining());
                merged.put(total);
                merged.put(current);
                merged.flip();
                return merged;
            })
            .map(fullBuffer -> {
                byte[] bytes = new byte[fullBuffer.remaining()];
                fullBuffer.get(bytes);
                return new MockMultipartFile(fileName, fileName, "audio/mpeg", bytes);
            });
}

2. 重构Service层逻辑

去掉重复订阅、删除阻塞调用,把同步逻辑包装成响应式类型,用弹性线程池隔离阻塞操作:

public Mono<ResponseEntity<AudioDto>> saveAudioTrack(HttpHeaders headers, String fileName, Flux<ByteBuffer> body) {
    // 先转换为完整MultipartFile,仅订阅一次请求体
    return AppUtils.byteBufferToMultipartFile(fileName, body)
            // 把阻塞的元数据解析、S3上传逻辑放到弹性线程池执行,避免阻塞event loop
            .flatMap(multipartFile -> Mono.fromCallable(() -> storeAudioMeta(multipartFile, headers))
                    .subscribeOn(Schedulers.boundedElastic()))
            // 用响应式MongoRepository保存,不要调用block()
            .flatMap(audioDto -> {
                AudioEntity audioEntity = new AudioEntity();
                BeanUtils.copyProperties(audioDto, audioEntity);
                return audioRepository.save(audioEntity);
            })
            // 构造成功响应
            .map(storedEntity -> {
                AudioDto storedAudio = new AudioDto();
                BeanUtils.copyProperties(storedEntity, storedAudio);
                return ResponseEntity.ok(storedAudio);
            })
            // 统一异常处理
            .onErrorResume(e -> {
                e.printStackTrace();
                return Mono.just(ResponseEntity.status(HttpStatus.INTERNAL_SERVER_ERROR).build());
            });
}

3. 修复元数据解析方法

修复临时文件路径、空指针、资源泄漏问题:

private AudioDto storeAudioMeta(MultipartFile multipartFile, HttpHeaders headers) throws IOException, InvalidDataException, UnsupportedTagException {
    AudioDto audioDto = new AudioDto();
    audioDto.setFile(multipartFile); // 补充file字段赋值,避免后续空指针
    if (multipartFile.isEmpty()) {
        throw new IllegalStateException("Cannot upload empty file");
    }
    if(!"audio/mpeg".equals(multipartFile.getContentType())) {
        throw new IllegalStateException("File uploaded is not a mp3");
    }
    // 生成系统临时文件,不要写入项目资源目录(打包后目录不存在)
    Path tempPath = Files.createTempFile("temp_", multipartFile.getOriginalFilename());
    File targetFile = tempPath.toFile();
    try {
        multipartFile.transferTo(targetFile);
        Mp3File mp3file  = new Mp3File(targetFile.getPath());
        audioDto.setLength((int) mp3file.getLengthInSeconds());
        audioDto.setBitrate(mp3file.getBitrate());
        audioDto.setSampleRate(mp3file.getSampleRate());
        if (mp3file.hasId3v1Tag()) {
            ID3v1 id3v1Tag = mp3file.getId3v1Tag();
            audioDto.setTrack(id3v1Tag.getTrack());
            audioDto.setArtist(id3v1Tag.getArtist());
            audioDto.setTitle(id3v1Tag.getTitle());
            // 其他ID3v1字段赋值
        }
        if (mp3file.hasId3v2Tag()) {
            ID3v2 id3v2Tag = mp3file.getId3v2Tag();
            audioDto.setTrack(id3v2Tag.getTrack());
            audioDto.setAlbumImage(id3v2Tag.getAlbumImage());
            byte[] albumImageData = id3v2Tag.getAlbumImage();
            if (albumImageData != null) {
                audioDto.setAlbumImageSize(albumImageData.length);
                audioDto.setAlbumImageMimeType(id3v2Tag.getAlbumImageMimeType());
            }
            // 其他ID3v2字段赋值
        }
        // S3上传逻辑恢复
        Map<String, String> metadata = new HashMap<>();
        metadata.put("Content-Type", audioDto.getFile().getContentType());
        metadata.put("Content-Length", String.valueOf(audioDto.getFile().getSize()));
        String bucketName = environment.getProperty("amazon.aws.s3.audioBucket");
        String path = String.format("%s/%s", bucketName, UUID.randomUUID());
        String s3FileName = audioDto.getFile().getOriginalFilename();
        try {
            fileStore.upload(path, s3FileName, Optional.of(metadata), audioDto.getFile().getInputStream());
        } catch (IOException e) {
            throw new IllegalStateException("Failed to upload file to S3", e);
        }
        audioDto.setFileName(s3FileName);
        audioDto.setFilePath(path);
        return audioDto;
    } finally {
        // 用完删除临时文件,避免磁盘泄漏
        Files.deleteIfExists(tempPath);
    }
}

内容的提问来源于stack exchange,提问作者davidvera

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.29 14:24:01