不能直接用 redistemplate 做流式响应,因其为阻塞式 api,会阻塞 netty 事件循环线程,破坏 webflux 非阻塞特性;必须使用 reactiveredistemplate(基于 lettuce),返回 mono/flux,支持背压与异步调度。

为什么不能直接用 RedisTemplate 做流式响应
因为 RedisTemplate 是阻塞式 API,它在调用 opsForValue().get() 或 scan() 时会同步等待 Redis 返回结果,这会卡住 Netty 的事件循环线程。一旦发生,整个 WebFlux 的非阻塞优势就没了,高并发下容易线程耗尽、超时堆积。
必须用响应式客户端——ReactiveRedisTemplate(底层基于 Lettuce),它返回的是 Mono 或 Flux,天然支持背压和异步调度。
-
ReactiveRedisTemplate的scan()方法返回Flux<string></string>,可直接用于流式分批拉取 key - 对大 value(如 JSON 字符串)做流式解析时,不能一次性
get()再拆,而应结合scan()+pipeline+ 分块订阅 - 若误配了
spring-boot-starter-data-redis(非 reactive 版),Spring Boot 会自动装配阻塞版,需手动排除
ReactiveRedisTemplate 流式扫描大键空间的写法
比如要从 Redis 扫描 10 万个 key 并逐个返回其 value,不能用 keys *(禁用!会阻塞 Redis),必须用游标式 scan:
@GetMapping(value = "/redis/keys", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public Flux<string> streamKeys() {
return redisTemplate.scan(ScanOptions.scanOptions()
.match("user:*")
.count(100)
.build())
.flatMap(key -> redisTemplate.opsForValue().get(key)
.map(value -> String.format("key=%s, value=%s", key, value))
.defaultIfEmpty(String.format("key=%s, value=null", key)))
.take(5000); // 防止无限流,加安全上限
}</string>
注意点:
-
count不是“每次返回多少条”,而是 hint,实际数量可能更少;值太大会增加单次 Redis 负担,建议 50–200 -
flatMap是关键:把每个 key 的异步get()转成并行流,但默认并发度是 256,生产环境建议加.concurrent(8)控制连接数 - 没加
take()或超时控制时,客户端断连后Flux可能不自动 cancel,需配合doOnCancel清理资源
流式写入 Redis 时怎么避免 OOM
向 Redis 批量写入大量数据(如导入日志),如果一次性构造百万级 Mono 再 collectList(),会吃光堆内存。正确做法是“边生成边发”,靠背压驱动:
Flux.range(1, 100_000)
.buffer(100) // 每 100 条打包成 list
.flatMap(batch -> Mono.fromRunnable(() -> {
// 构造 pipeline 命令,一次发 100 个 set
var pipeline = redisTemplate.getConnectionFactory().getConnection().pipelined();
batch.forEach(i -> pipeline.set(("log:" + i).getBytes(), ("data-" + i).getBytes()));
pipeline.exec();
}), 4) // 并发最多 4 个 pipeline
.then();
要点:
- 别用
Flux.concatMap—— 它是串行,吞吐低;flatMap并发可控才是流式写入的核心 - Redis 的 pipeline 不是原子的,失败需重试逻辑,
exec()返回List<object></object>,要检查 null 或异常 - Lettuce 默认连接池最大 8 个连接,
flatMap并发数 > 连接数会导致排队,可通过ReactiveRedisConnectionFactory调整maxIdle/maxAcquire
客户端断连后 Redis 流还在跑?得手动 cancel
WebFlux 的 Flux 默认不会感知 HTTP 连接关闭。用户关掉浏览器或网络中断,服务端仍可能继续 scan / get,浪费 Redis 资源和 CPU。
必须显式监听取消信号:
return redisTemplate.scan(options)
.doOnCancel(() -> log.info("Client disconnected, scan cancelled"))
.doOnTerminate(() -> log.info("Stream finished or cancelled"))
.onErrorResume(e -> {
log.error("Redis stream error", e);
return Flux.empty();
});
更稳妥的做法是:在 Controller 方法里注入 ServerWebExchange,用 exchange.getResponse().isCommitted() 判断是否已写出,但不如 doOnCancel 直接可靠。
真正容易被忽略的是:Lettuce 的 scan 游标本身不带 cancel 支持,所以 cancel 后当前批次可能仍会完成,但后续游标不再发起请求——这是框架层限制,不是 bug。











