同步到异步:Java HTTP IO 编程模型演进实验
全同步
主要代码:
val t1 = System.currentTimeMillis()
val albumDetailDtos = newAlbums
.mapNotNull { a ->
try {
api.getAlbumDetail(a.cid)
} catch (e: Exception) {
logger.error("获取专辑详情失败.", e)
null
}
}
val t2 = System.currentTimeMillis()
...
val t5 = System.currentTimeMillis()
val songEntities = albumDetailDtos
.flatMap { it.songs ?: emptyList() }
.apply {
logger.info("新歌曲数量: ${this.size}, 准备同步...")
}
.mapNotNull { s ->
try {
val songDetailDto = api.getSongDetail(s.cid)
// s 中没有 albumCid ,只有 songDetailDto 中才有
val albumEntity = cidAlbumMap[songDetailDto.albumCid]
songDetailDto.toEntity(albumEntity)
} catch (e: Exception) {
logger.error("获取歌曲详情失败.", e)
null
}
}
val t6 = System.currentTimeMillis()
...
logger.info("获取新专辑用时:${t2 - t1}ms")
logger.info("获取新歌曲用时:${t6 - t5}ms")
logger.info("本次更新,专辑 ${persistedAlbums.size} 张, 歌曲 ${persistedSongs.size} 首.")结果:
新专辑数量: 285, 准备同步...
新歌曲数量: 877, 准备同步...
获取新专辑用时:15845ms
获取新歌曲用时:44440ms
本次更新,专辑 285 张, 歌曲 877 首.使用 CompletableFuture 异步
先定义 IO 密集型线程池:
/**
* 短时 IO 密集型
*/
@Bean
fun apiExecutor(): ExecutorService {
// CPU 密集型:核心数 + 1(多1个是为了防止偶尔的缺页中断等导致CPU空闲)
val cpus = Runtime.getRuntime().availableProcessors()
// IO 密集型:通常为核心数的 2 倍,或者根据阻塞比例计算:核心数 * (1 + W/C)
// W/C 为等待时间与计算时间的比值
return ThreadPoolExecutor(
cpus * 2,
cpus * 2,
10L,
TimeUnit.SECONDS,
ArrayBlockingQueue(2048),
)
.apply {
allowCoreThreadTimeOut(true)
}
}主要代码:
val t1 = System.currentTimeMillis()
val albumDetailDtos = newAlbums
.map { album ->
// 默认 ForkJoinPool.commonPool() 线程池是 CPU 密集型,换自定义的 IO 密集型 线程池
CompletableFuture.supplyAsync({
try {
api.getAlbumDetail(album.cid)
} catch (e: Exception) {
logger.error("获取专辑详情失败", e)
null
}
}, apiExecutor)
}
.mapNotNull { it.join() }
val t2 = System.currentTimeMillis()
...
val t5 = System.currentTimeMillis()
val songEntities = albumDetailDtos
.flatMap { it.songs ?: emptyList() }
.apply {
logger.info("新歌曲数量: ${this.size}, 准备同步...")
}
.map { song ->
CompletableFuture.supplyAsync({
try {
val detail = api.getSongDetail(song.cid)
val album = cidAlbumMap[detail.albumCid]
detail.toEntity(album)
} catch (e: Exception) {
logger.error("获取歌曲失败", e)
null
}
}, apiExecutor)
}
.mapNotNull { it.join() }
val t6 = System.currentTimeMillis()
...
logger.info("获取新专辑用时:${t2 - t1}ms")
logger.info("获取新歌曲用时:${t6 - t5}ms")
logger.info("本次更新,专辑 ${persistedAlbums.size} 张, 歌曲 ${persistedSongs.size} 首.")结果:
新专辑数量: 285, 准备同步...
新歌曲数量: 877, 准备同步...
获取新专辑用时:637ms
获取新歌曲用时:1950ms
本次更新,专辑 285 张, 歌曲 877 首.使用虚拟线程异步
主要代码
// 限制并发,避免触发服务器的 HTTP2 并发流过载错误
val semaphore = Semaphore(96)
...
val t1 = System.currentTimeMillis()
val albumDetailDtos = Executors.newVirtualThreadPerTaskExecutor().use { executor ->
val futures = newAlbums.map { album ->
executor.submit<SirenAlbumDetailDto?> {
semaphore.acquire()
try {
api.getAlbumDetail(album.cid)
} catch (e:Exception) {
logger.error("获取专辑详情失败: {}", album.cid, e)
null
} finally {
semaphore.release()
}
}
}
futures.mapNotNull { it.get() }
}
val t2 = System.currentTimeMillis()
...
val t5 = System.currentTimeMillis()
val songs = albumDetailDtos
.flatMap { it.songs ?: emptyList() }
.also {
logger.info("新歌曲数量: ${it.size}, 准备同步...")
}
val songEntities = Executors.newVirtualThreadPerTaskExecutor().use { executor ->
val futures = songs.map { song ->
executor.submit<SongEntity?> {
semaphore.acquire()
try {
val dto = api.getSongDetail(song.cid)
val album = cidAlbumMap[dto.albumCid]
dto.toEntity(album)
} catch (e: Exception) {
logger.error("获取歌曲详情失败: {}", song.cid, e)
null
} finally {
semaphore.release()
}
}
}
futures.mapNotNull { it.get() }
}
val t6 = System.currentTimeMillis()结果:
新专辑数量: 285, 准备同步...
新歌曲数量: 877, 准备同步...
获取新专辑用时:210ms
获取新歌曲用时:718ms
本次更新,专辑 285 张, 歌曲 877 首.结果对比
测试数据规模一致:
- 新专辑:285
- 新歌曲:877
核心耗时对比:
| 模型 | 获取专辑 | 获取歌曲 | 总耗时 |
|---|---|---|---|
| 全同步 | 15845ms | 44440ms | 60285ms |
| CompletableFuture 异步 | 637ms | 1950ms | 2587ms |
| 虚拟线程异步 | 210ms | 718ms | 928ms |
相对全同步的加速倍数:
| 模型 | 专辑阶段 | 歌曲阶段 | 总体 |
|---|---|---|---|
| CompletableFuture 异步 | 24.9x | 22.8x | 23.3x |
| 虚拟线程异步 | 75.5x | 61.9x | 65.0x |
相对 CompletableFuture 的进一步提升:
| 模型 | 专辑阶段 | 歌曲阶段 | 总体 |
|---|---|---|---|
| 虚拟线程异步 | 3.0x | 2.7x | 2.8x |
结论:
- 对 HTTP IO 这类“等待远大于计算”的场景,全同步模型的时间几乎全部耗在串行等待上。
CompletableFuture + IO 线程池已经能显著压缩总耗时,说明并发拉取本身就是这一类任务的主要优化点。- 虚拟线程在保持“阻塞式写法”的前提下继续降低了调度与线程管理成本,最终拿到更短的尾延迟和更高的吞吐。
- 这组数据里,虚拟线程方案总耗时从
60285ms压缩到928ms,总体验证了“同步写法 + 大规模轻量并发”在 Java HTTP IO 场景中的优势。 - 并发不能无限放大;实验里通过
Semaphore(96)做限流,否则容易触发服务端 HTTP/2 并发流过载错误。