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 streamstoo many concurrent streams 是 HTTP/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 | 是否成功 | 获取专辑 | 获取歌曲 | 说明 |
|---|---|---|---|---|
| 64 | 是 | 271ms | 838ms | 保守值,稳定 |
| 96 | 是 | 210ms | 718ms | 更快,仍稳定 |
| 128 | 是 | 179ms | 586ms | 协议上限 |
| 129 | 否 | 144ms | 198ms | 请求失败,结果无效 |
可以看到:
- 在请求成功的前提下,并发度从
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的极限值运行,为了稳定性应预留余量,例如使用64或96这类更保守的并发值。