如何从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
相关产品推荐
相关产品推荐

