Spring Boot 3.2.2接收WAV音频流异常及响应式改造咨询
问题修复与响应式改造方案
一、当前阻塞问题修复
1. 客户端问题根源与修复
客户端卡在send方法、服务端未调用echo的核心原因:
AudioInputStream包装TargetDataLine后,未正确触发流结束标记,导致HTTP客户端持续等待发送数据- 虚拟线程直接关闭
TargetDataLine,未确保剩余音频数据被读取完成
修改客户端代码:
public class AudioRecorder { private static final HttpClient HTTP_CLIENT = HttpClient.newHttpClient(); public static void main(String... args) throws InterruptedException, LineUnavailableException, IOException, UnsupportedAudioFileException { AudioFormat format = new AudioFormat(8000, 8, 1, true, true); try (TargetDataLine targetLine = AudioSystem.getTargetDataLine(format)) { targetLine.open(); targetLine.start(); System.out.println("Recording..."); // 调整录音停止逻辑:排空数据再关闭 Thread.startVirtualThread((ThrowingRunnable) () -> { Thread.sleep(1000); targetLine.stop(); targetLine.drain(); // 排空剩余音频数据 targetLine.close(); System.out.println("Stopped"); }); final HttpResponse<Void> httpResponse = HTTP_CLIENT.send( HttpRequest.newBuilder() .uri(URI.create("http://localhost:8080/echo")) .header("Content-Type", "audio/wav") .POST(HttpRequest.BodyPublishers.ofInputStream(() -> { AudioInputStream audioStream = new AudioInputStream(targetLine); // 监听TargetDataLine关闭,主动终止音频流 Thread.startVirtualThread(() -> { try { targetLine.join(); audioStream.close(); } catch (InterruptedException | IOException e) { Thread.currentThread().interrupt(); } }); return audioStream; })) .build(), HttpResponse.BodyHandlers.discarding() ); System.out.println("Got status code: " + httpResponse.statusCode()); } } }
2. 服务端问题根源与修复
- 原代码中
try-with-resources会在clip.start()后立即关闭资源,导致音频未播放完成就被终止 - 未明确指定接口接收的媒体类型,可能导致Spring无法正确匹配请求
修改服务端接口:
@PostMapping(value = "echo", consumes = "audio/wav") public ResponseEntity<?> echo(@RequestBody InputStream in) throws UnsupportedAudioFileException, IOException, LineUnavailableException, InterruptedException { System.out.println("Starting!!!!!!!!!!!!!!!"); try (AudioInputStream audioInputStream = AudioSystem.getAudioInputStream(in); Clip clip = AudioSystem.getClip()) { clip.open(audioInputStream); clip.start(); // 等待音频播放完成再释放资源 while (clip.isRunning()) { Thread.sleep(100); } } return ResponseEntity.ok().build(); }
额外配置:添加流式消息转换器,确保Spring正确处理InputStream请求体
@Configuration public class WebConfig implements WebMvcConfigurer { @Override public void configureMessageConverters(List<HttpMessageConverter<?>> converters) { converters.add(new InputStreamHttpMessageConverter()); } }
二、响应式实现改造(WebFlux)
1. 服务端改造(基于Spring WebFlux)
替换Spring MVC依赖为WebFlux,编写流式音频处理接口:
@RestController public class ReactiveAudioController { // 流式处理:边接收边播放 @PostMapping(value = "echo-stream", consumes = "audio/wav") public Mono<ResponseEntity<Void>> echoStream(@RequestBody Flux<DataBuffer> audioFlux) { System.out.println("Starting reactive stream echo!!!!!!!!!!!!!!!"); return Mono.create(sink -> { try { AudioFormat format = new AudioFormat(8000, 8, 1, true, true); Clip clip = (Clip) AudioSystem.getLine(new DataLine.Info(Clip.class, format)); clip.open(format); audioFlux.subscribe( dataBuffer -> { try { byte[] bytes = new byte[dataBuffer.readableByteCount()]; dataBuffer.read(bytes); clip.setMicrosecondPosition(0); clip.open(format, bytes, 0, bytes.length); clip.start(); } catch (LineUnavailableException | IOException e) { sink.error(e); } finally { DataBufferUtils.release(dataBuffer); } }, sink::error, () -> { // 等待最后一段音频播放完成 while (clip.isRunning()) { try { Thread.sleep(100); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } } clip.close(); sink.success(ResponseEntity.ok().build()); } ); } catch (LineUnavailableException e) { sink.error(e); } }); } }
2. 客户端响应式改造(基于WebClient)
public class ReactiveAudioRecorder { private static final WebClient WEB_CLIENT = WebClient.create("http://localhost:8080"); public static void main(String... args) throws InterruptedException, LineUnavailableException { AudioFormat format = new AudioFormat(8000, 8, 1, true, true); try (TargetDataLine targetLine = AudioSystem.getTargetDataLine(format)) { targetLine.open(); targetLine.start(); System.out.println("Recording..."); Thread.startVirtualThread((ThrowingRunnable) () -> { Thread.sleep(1000); targetLine.stop(); targetLine.drain(); targetLine.close(); System.out.println("Stopped"); }); // 将TargetDataLine转为响应式数据流 Flux<DataBuffer> audioFlux = Flux.create(sink -> { byte[] buffer = new byte[1024]; int readLen; try { while ((readLen = targetLine.read(buffer, 0, buffer.length)) != -1) { sink.next(DataBufferFactory.DEFAULT.wrap(buffer, 0, readLen)); } sink.complete(); } catch (IOException e) { sink.error(e); } }); WEB_CLIENT.post() .uri("/echo-stream") .contentType(MediaType.parseMediaType("audio/wav")) .body(audioFlux, DataBuffer.class) .retrieve() .toBodilessEntity() .doOnSuccess(resp -> System.out.println("Got status code: " + resp.getStatusCode())) .doOnError(err -> System.err.println("Error: " + err.getMessage())) .block(); } } }
内容的提问来源于stack exchange,提问作者Roy Ash
相关产品推荐
相关产品推荐

