Java异步与多线程结合性能低下问题排查与优化
我有如下函数,返回音频文件下载完成后的文件名列表:
List<CompletableFuture<String>> getAudios(List<String> text)
在函数内部,我通过简单的for循环调用服务类的talk方法:
for (String s: text) { futures.add(this.talkService.talk(s)); }
talk方法会向Azure API发送请求,完成两项操作:
- 获取字节数组形式的SpeechSynthesisResult音频数据
- 将数据写入文件
Azure API获取音频的方法仅返回Future而非CompletableFuture:
Future<SpeechSynthesisResult> SpeakSsmlAsync(final String var1)
因此我需要调用阻塞方法.get()来获取音频。为解决阻塞问题,我启动了一个并行线程执行该操作:
CompletableFuture<String> talk(String text) { CompletableFuture<String> completableFuture = new CompletableFuture<>(); SpeechSynthesizer speechSynthesizer = new SpeechSynthesizer(config, null); ... taskExecutor.execute(() -> { ... Future<SpeechSynthesisResult> result = speechSynthesizer.SpeakSsmlAsync(ssmlText); SpeechSynthesisResult audioResult = result.get(); if (audioResult.getReason() == ResultReason.SynthesizingAudioCompleted) { byte[] data = audioResult.getAudioData(); String filename = this.fileService.writeToFile(data, "mp3"); completableFuture.complete(filename); } }); return completableFuture; }
writeToFile是同步方法,使用FileOutputStream将byte[]数据写入文件。
我的控制器调用action方法合并所有音频文件:
public void action(List<String> text) { List<CompletableFuture<String>> audios = getAudios(text); // 等待所有音频下载完成后获取文件名列表 List<String> filenames = CompletableFuture.allOf(audios.toArray(new CompletableFuture[0])) .thenApply((v) -> { return audios.stream().map(CompletableFuture::join) .collect(Collectors.toList()); }).join(); }
但查看控制台日志时发现,所有音频文件总是按顺序下载,几乎没有并发执行。而且每个音频文件的下载耗时相同,我原本期望触发所有文件下载后,它们能几乎同时完成。请问代码是否存在问题?如何进行性能优化?
编辑补充:
@SpringBootApplication(exclude = {SecurityAutoConfiguration.class}) public class MyApplication { @Bean public ThreadPoolTaskExecutor taskExecutor() { ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); return executor; } ... }
核心问题:线程池未正确配置
你定义的ThreadPoolTaskExecutor没有设置任何参数,使用的是默认配置:
- 核心线程数默认是
1 - 队列容量默认是
Integer.MAX_VALUE
这就导致所有任务都会被放到无限大的队列里,只有核心的1个线程在逐个执行任务,完全没有并发能力,所以任务是串行执行的。
优化步骤
1. 正确配置线程池
修改线程池Bean的配置,根据业务需求设置合理的线程数(注意不要超过Azure Speech API的并发限制,避免触发限流):
@Bean public ThreadPoolTaskExecutor taskExecutor() { ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); // 核心线程数,可根据CPU核心数或API并发限制设置,例如8 executor.setCorePoolSize(8); // 最大线程数,可设为核心线程数的2倍左右 executor.setMaxPoolSize(16); // 队列容量,不要设太大,避免任务堆积 executor.setQueueCapacity(100); // 线程空闲超时时间 executor.setKeepAliveSeconds(60); // 线程名称前缀,便于日志排查 executor.setThreadNamePrefix("audio-task-"); // 初始化线程池 executor.initialize(); return executor; }
2. 简化CompletableFuture创建方式
手动创建CompletableFuture并通过taskExecutor.execute提交任务的方式不够优雅,建议使用CompletableFuture.supplyAsync直接指定线程池,同时注意释放资源:
CompletableFuture<String> talk(String text) { return CompletableFuture.supplyAsync(() -> { SpeechSynthesizer speechSynthesizer = new SpeechSynthesizer(config, null); try { Future<SpeechSynthesisResult> result = speechSynthesizer.SpeakSsmlAsync(ssmlText); SpeechSynthesisResult audioResult = result.get(); if (audioResult.getReason() == ResultReason.SynthesizingAudioCompleted) { byte[] data = audioResult.getAudioData(); return this.fileService.writeToFile(data, "mp3"); } // 处理合成失败的情况 throw new RuntimeException("音频合成失败: " + audioResult.getReason()); } catch (Exception e) { throw new RuntimeException("音频合成过程出错", e); } finally { // 释放SpeechSynthesizer资源 if (speechSynthesizer != null) { speechSynthesizer.close(); } } }, taskExecutor); }
3. 复用SpeechSynthesizer实例
每次调用talk方法都创建新的SpeechSynthesizer会带来不必要的资源开销,建议将其改为单例(需确认实例线程安全,若不安全则改用对象池管理):
// 在TalkService中注入SpeechSynthesizer单例 private final SpeechSynthesizer speechSynthesizer; public TalkService(SpeechConfig config) { this.speechSynthesizer = new SpeechSynthesizer(config, null); } CompletableFuture<String> talk(String text) { return CompletableFuture.supplyAsync(() -> { try { Future<SpeechSynthesisResult> result = speechSynthesizer.SpeakSsmlAsync(ssmlText); SpeechSynthesisResult audioResult = result.get(); if (audioResult.getReason() == ResultReason.SynthesizingAudioCompleted) { byte[] data = audioResult.getAudioData(); return this.fileService.writeToFile(data, "mp3"); } throw new RuntimeException("音频合成失败: " + audioResult.getReason()); } catch (Exception e) { throw new RuntimeException("音频合成过程出错", e); } }, taskExecutor); }
4. 优化文件写入操作
同步的FileOutputStream写入在高并发下可能成为瓶颈,可以:
- 使用
Files.write方法优化IO操作 - 为文件写入单独分配线程池,避免占用音频合成的线程资源
// 新增文件写入线程池Bean @Bean public ThreadPoolTaskExecutor fileWriteExecutor() { ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); executor.setCorePoolSize(4); executor.setMaxPoolSize(8); executor.setQueueCapacity(50); executor.setThreadNamePrefix("file-write-"); executor.initialize(); return executor; } // 修改talk方法,拆分音频合成与文件写入任务 CompletableFuture<String> talk(String text) { return CompletableFuture.supplyAsync(() -> { // 音频合成逻辑 SpeechSynthesizer speechSynthesizer = new SpeechSynthesizer(config, null); try { Future<SpeechSynthesisResult> result = speechSynthesizer.SpeakSsmlAsync(ssmlText); SpeechSynthesisResult audioResult = result.get(); if (audioResult.getReason() == ResultReason.SynthesizingAudioCompleted) { return audioResult.getAudioData(); } throw new RuntimeException("音频合成失败: " + audioResult.getReason()); } catch (Exception e) { throw new RuntimeException("音频合成过程出错", e); } finally { if (speechSynthesizer != null) { speechSynthesizer.close(); } } }, taskExecutor) .thenCompose(data -> CompletableFuture.supplyAsync(() -> { // 异步写入文件 return this.fileService.writeToFile(data, "mp3"); }, fileWriteExecutor)); }
5. 简化结果收集逻辑
原代码的结果收集逻辑可以简化,保持功能不变:
List<String> filenames = CompletableFuture.allOf(audios.toArray(new CompletableFuture[0])) .thenApply(v -> audios.stream() .map(CompletableFuture::join) .collect(Collectors.toList())) .join();
内容的提问来源于stack exchange,提问作者parsecer

