JDK Virtual Thread + Spring RestClient + HTTP/2 最大并发流测试

场景

使用第三方开放 API 做 ETL,使用虚拟线程做 Extract 过程。

过程

主要代码:

    		// 使用虚拟线程做 HTTP IO 密集型请求
        val t1 = System.currentTimeMillis()
        // use 会自动关闭线程池,所以如果使用 use 则不要抽成类变量
        val albumDetailDtos = Executors.newVirtualThreadPerTaskExecutor().use { executor ->
            val futures = newAlbums.map { album ->
                executor.submit<SirenAlbumDetailDto?> {
                    try {
                        api.getAlbumDetail(album.cid)
                    } catch (e:Exception) {
                        logger.error("获取专辑详情失败: {}", album.cid, e)
                        null
                    }
                }
            }
 
            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?> {
                    try {
                        val dto = api.getSongDetail(song.cid)
                        val album = cidAlbumMap[dto.albumCid]
                        dto.toEntity(album)
                    } catch (e: Exception) {
                        logger.error("获取歌曲详情失败: {}", song.cid, e)
                        null
                    }
                }
            }
 
            futures.mapNotNull { it.get() }
        }
        val t6 = System.currentTimeMillis()
        
				...
 
        logger.info("获取新专辑用时:${t2 - t1}ms")
        logger.info("获取新歌曲用时:${t6 - t5}ms")
        logger.info("本次更新,专辑 ${persistedAlbums.size} 张, 歌曲 ${persistedSongs.size} 首.")

结果出现大量并发错误:

2026-07-17T21:34:26.040+08:00 ERROR 31146 --- [terra-echo] [    virtual-382] i.g.l.t.c.xxxSynchronizer   : 获取专辑详情失败: 8947
org.springframework.web.client.ResourceAccessException: I/O error on GET request for "https://monster-siren.hypergryph.com/api/song/125010": too many concurrent streams

too many concurrent streamsHTTP/2 的并发流(stream)限制 ,先用 nghttp 查看服务器连接信息:

# looko @ MacBookPro in ~ [21:46:44]
$ nghttp -nv https://monster-siren.hypergryph.com/api/song/125010
[  0.038] Connected
The negotiated protocol: h2
[  0.153] recv SETTINGS frame <length=18, flags=0x00, stream_id=0>
          (niv=3)
          [SETTINGS_MAX_CONCURRENT_STREAMS(0x03):128]
          [SETTINGS_INITIAL_WINDOW_SIZE(0x04):524288]
          [SETTINGS_MAX_FRAME_SIZE(0x05):16777215]
[  0.153] recv WINDOW_UPDATE frame <length=4, flags=0x00, stream_id=0>
          (window_size_increment=2147418112)
[  0.153] send SETTINGS frame <length=18, flags=0x00, stream_id=0>
          (niv=3)
          [SETTINGS_MAX_CONCURRENT_STREAMS(0x03):100]
          [SETTINGS_INITIAL_WINDOW_SIZE(0x04):65535]
          [SETTINGS_NO_RFC7540_PRIORITIES(0x09):1]
[  0.153] send SETTINGS frame <length=0, flags=0x01, stream_id=0>
          ; ACK
          (niv=0)
[  0.153] send HEADERS frame <length=67, flags=0x05, stream_id=1>
          ; END_STREAM | END_HEADERS
          (padlen=0)
          ; Open new stream
          :method: GET
          :path: /api/song/125010
          :scheme: https
          :authority: monster-siren.hypergryph.com
          priority: u=3
          accept: */*
          accept-encoding: gzip, deflate
          user-agent: nghttp2/1.69.0
[  0.190] recv SETTINGS frame <length=0, flags=0x01, stream_id=0>
          ; ACK
          (niv=0)
[  0.196] recv (stream_id=1) :status: 200
[  0.196] recv (stream_id=1) date: Fri, 17 Jul 2026 13:53:46 GMT
[  0.196] recv (stream_id=1) content-type: application/json; charset=utf-8
[  0.196] recv (stream_id=1) content-length: 325
[  0.196] recv (stream_id=1) x-xss-protection: 1; mode=block
[  0.196] recv (stream_id=1) x-content-type-options: nosniff
[  0.196] recv (stream_id=1) x-download-options: noopen
[  0.196] recv (stream_id=1) x-readtime: 1
[  0.196] recv HEADERS frame <length=136, flags=0x04, stream_id=1>
          ; END_HEADERS
          (padlen=0)
          ; First response header
[  0.196] recv DATA frame <length=325, flags=0x00, stream_id=1>
[  0.196] recv DATA frame <length=0, flags=0x01, stream_id=1>
          ; END_STREAM
[  0.196] send GOAWAY frame <length=8, flags=0x00, stream_id=0>
          (last_stream_id=0, error_code=NO_ERROR(0x00), opaque_data(0)=[])

可知:
SETTINGS_MAX_CONCURRENT_STREAMS = 128

解决方案:用 Semaphore 限制并发:

        // 限制并发,避免触发服务器的 HTTP2 并发流过载错误
        val semaphore = Semaphore(64)
				...
        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, 准备同步...
获取新专辑用时:271ms
获取新歌曲用时:838ms
本次更新,专辑 285 张, 歌曲 877 首.

能请求成功了,接下来做逼近实验:

val semaphore = Semaphore(96) 时 :

新专辑数量: 285, 准备同步...
新歌曲数量: 877, 准备同步...
获取新专辑用时:210ms
获取新歌曲用时:718ms
本次更新,专辑 285 张, 歌曲 877 首.

val semaphore = Semaphore(128) 时 :

新专辑数量: 285, 准备同步...
新歌曲数量: 877, 准备同步...
获取新专辑用时:179ms
获取新歌曲用时:586ms
本次更新,专辑 285 张, 歌曲 877 首.

val semaphore = Semaphore(129) 时,请求有错误,部分请求失败,数据没能完整更新 :

获取歌曲详情失败: 779437
org.springframework.web.client.ResourceAccessException: I/O error on GET request for "https://monster-siren.hypergryph.com/api/song/779437": too many concurrent streams
新专辑数量: 285, 准备同步...
新歌曲数量: 877, 准备同步...
获取新专辑用时:144ms
获取新歌曲用时:198ms
本次更新,专辑 235 张, 歌曲 480 首.

不同并发度结果对比

Semaphore是否成功获取专辑获取歌曲说明
64271ms838ms保守值,稳定
96210ms718ms更快,仍稳定
128179ms586ms协议上限
129144ms198ms请求失败,结果无效

可以看到:

  • 在请求成功的前提下,并发度从 64 提升到 96、再到 128,吞吐持续提升。
  • 129 虽然表面耗时更短,但这是部分请求失败后的不完整结果,不能作为有效性能数据。
  • 这一组数据说明:并发度的上升确实能继续压缩 HTTP IO 耗时,但前提是不能突破协议与服务端的承载上限。

结论

  • 服务器通过 HTTP/2 SETTINGS_MAX_CONCURRENT_STREAMS 将单连接最大并发流限制为 128
  • Semaphore(128) 可以跑满协议上限,Semaphore(129) 会触发 too many concurrent streams,实验结果与 nghttp 查询一致。
  • 虚拟线程解决的是线程调度成本,不会自动控制资源并发度。
  • 在 HTTP IO 场景里,并发度仍然需要按协议约束和服务端承载能力显式限流。
  • 实际使用不宜贴着 128 的极限值运行,为了稳定性应预留余量,例如使用 6496 这类更保守的并发值。