同步到异步: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

核心耗时对比:

模型获取专辑获取歌曲总耗时
全同步15845ms44440ms60285ms
CompletableFuture 异步637ms1950ms2587ms
虚拟线程异步210ms718ms928ms

相对全同步的加速倍数:

模型专辑阶段歌曲阶段总体
CompletableFuture 异步24.9x22.8x23.3x
虚拟线程异步75.5x61.9x65.0x

相对 CompletableFuture 的进一步提升:

模型专辑阶段歌曲阶段总体
虚拟线程异步3.0x2.7x2.8x

结论:

  • 对 HTTP IO 这类“等待远大于计算”的场景,全同步模型的时间几乎全部耗在串行等待上。
  • CompletableFuture + IO 线程池 已经能显著压缩总耗时,说明并发拉取本身就是这一类任务的主要优化点。
  • 虚拟线程在保持“阻塞式写法”的前提下继续降低了调度与线程管理成本,最终拿到更短的尾延迟和更高的吞吐。
  • 这组数据里,虚拟线程方案总耗时从 60285ms 压缩到 928ms,总体验证了“同步写法 + 大规模轻量并发”在 Java HTTP IO 场景中的优势。
  • 并发不能无限放大;实验里通过 Semaphore(96) 做限流,否则容易触发服务端 HTTP/2 并发流过载错误。