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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.14 01:24:52