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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.01 18:46:02