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

如何从Mono<FilePart>获取InputStream?WebFlux音视频处理遇阻

WebFlux结合FFprobe/FFmpeg提取音频的响应式实现问题

我是WebFlux新手,原本有个用FFprobe和FFmpeg提取视频音频的简单应用,想改成响应式实现,但多次尝试都失败了。

控制器代码

@PostMapping("/upload")
public String upload(@RequestPart("file") Mono<FilePart> filePartMono, final Model model) {
    Flux<String> filenameList = mediaComponent.extractAudio(filePartMono);
    model.addAttribute("filenameList", new ReactiveDataDriverContextVariable(filenameList));
    return "download";
}

从视频中获取音频流的函数

public Mono<FFprobeResult> getAudioStreams(InputStream inputStream) {
    try {
        return Mono.just(FFprobe.atPath(FFprobePath)
                .setShowStreams(true)
                .setSelectStreams(StreamType.AUDIO)
                .setLogLevel(LogLevel.INFO)
                .setInput(inputStream)
                .execute());
    } catch (JaffreeException e) {
        log.error(e.getMessage(), e);
        return Mono.error(new MediaException("Audio formats could not be identified."));
    }
}

尝试过程及问题

尝试1

public Flux<String> extractAudio(Mono<FilePart> filePartMono) {
    filePartMono.flatMapMany(Part::content)
            .map(dataBuffer -> dataBuffer.asInputStream(true))
            .flatMap(this::getAudioStreams)
            .subscribe(System.out::println);
    ...
}

尝试3

public Flux<String> extractAudio(Mono<FilePart> filePartMono) {
    DataBufferUtils.write(filePartMono.flatMapMany(Part::content), OutputStream.nullOutputStream())
            .map(dataBuffer -> dataBuffer.asInputStream(true))
            .flatMap(this::getAudioStreams)
            .subscribe(System.out::println);
    ...
}

尝试1和3结果一致,FFprobe报错:

2022-10-30 11:24:30.292  WARN 79049 --- [         StdErr] c.g.k.jaffree.process.BaseStdReader      : [mov,mp4,m4a,3gp,3g2,mj2 @ 0x7f9162702340] [warning] STSZ atom truncated
2022-10-30 11:24:30.292 ERROR 79049 --- [         StdErr] c.g.k.jaffree.process.BaseStdReader      : [mov,mp4,m4a,3gp,3g2,mj2 @ 0x7f9162702340] [error] stream 0, contradictionary STSC and STCO
2022-10-30 11:24:30.292 ERROR 79049 --- [         StdErr] c.g.k.jaffree.process.BaseStdReader      : [mov,mp4,m4a,3gp,3g2,mj2 @ 0x7f9162702340] [error] error reading header
2022-10-30 11:24:30.294 ERROR 79049 --- [         StdErr] c.g.k.jaffree.process.BaseStdReader      : [error] tcp://127.0.0.1:51532: Invalid data found when processing input
2022-10-30 11:24:30.295  INFO 79049 --- [oundedElastic-3] c.g.k.jaffree.process.ProcessHandler     : Process has finished with status: 1
2022-10-30 11:24:30.409 ERROR 79049 --- [oundedElastic-3] c.e.s.application.MediaComponent         : Process execution has ended with non-zero status: 1. Check logs for detailed error message.

尝试2

public Flux<String> extractAudio(Mono<FilePart> filePartMono) {
    filePartMono.flatMapMany(Part::content)
            .reduce(InputStream.nullInputStream(), (inputStream, dataBuffer) -> new SequenceInputStream(
                    inputStream, dataBuffer.asInputStream()
            ))
            .flatMap(this::getAudioStreams)
            .subscribe(System.out::println);
    ...
}

出现栈溢出错误:

Exception in thread "Runnable-0" java.lang.StackOverflowError
    at java.base/java.io.SequenceInputStream.read(SequenceInputStream.java:198)

尝试4

public Flux<String> extractAudio(Mono<FilePart> filePartMono) {
    try {
        FFprobeResult FFprobeResult = getAudioStreams(getInputStreamFromFluxDataBuffer(filePartMono.flatMapMany(Part::content))).subscribe(System.out::println);
    return Flux.just("file ", "file2").delayElements(Duration.ofMinutes(1));
    } catch (IOException e) {
        log.error(e.getMessage(), e);
        return Flux.error(new MediaException("Audio extraction failed"));
    }
}

public InputStream getInputStreamFromFluxDataBuffer(Flux<DataBuffer> dataBuffer) throws IOException {
    PipedOutputStream pipedOutputStream = new PipedOutputStream();
    PipedInputStream pipedInputStream = new PipedInputStream(pipedOutputStream);
    DataBufferUtils.write(dataBuffer, pipedOutputStream)
        .subscribeOn(Schedulers.boundedElastic())
        .doOnComplete(() -> {
            try {
                pipedOutputStream.close();
            } catch (IOException e) {
                log.error(e.getMessage(), e);
            }
        })
        .subscribe(DataBufferUtils.releaseConsumer());
    return pipedInputStream;
}

此次FFprobe启动后一直运行,仿佛InputStream是无限的。

问题咨询

请问正确的实现方向是什么?我哪里做错了?另外,目前只是输出到控制台查看结果,最终需要将所有输出流作为参数传给另一个函数完成音频提取,这部分也需要解决。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.14 23:45:39