Webflux实现首元素立即返回,后续元素延迟输出的方法
解决方案
要实现首元素无延迟、后续元素添加限流延迟的需求,有两种简洁的实现方式:
方式一:通过索引判断延迟
利用index()操作符为每个元素添加索引,针对第一个元素(索引为0)直接返回,后续元素则添加延迟:
var connection = Mono.fromFuture( () -> dataCollector.getDataFromCassandra((CassandraExposureRetrieveRequest) command.getBody())) .repeat() .index() .flatMap(tuple -> { if (tuple.getT1() == 0) { return Mono.just(tuple.getT2()); } else { return Mono.just(tuple.getT2()).delayElement(Duration.ofMillis(delayExposure)); } }) .share();
方式二:拆分首元素与后续元素
单独获取无延迟的首元素,再将后续重复调用的逻辑添加延迟后合并:
// 首元素直接返回,无延迟 Mono<T> firstElement = Mono.fromFuture( () -> dataCollector.getDataFromCassandra((CassandraExposureRetrieveRequest) command.getBody())); // 后续元素每次获取前添加延迟,持续重复 Flux<T> subsequentElements = Mono.fromFuture( () -> dataCollector.getDataFromCassandra((CassandraExposureRetrieveRequest) command.getBody())) .delayElement(Duration.ofMillis(delayExposure)) .repeat(); // 合并首元素和后续流 var connection = Flux.concat(firstElement, subsequentElements) .share();
两种方式都能满足需求:首元素会立即触发数据库查询并返回,后续元素则会在每次查询前等待指定延迟,达到限流效果。
内容的提问来源于stack exchange,提问作者Icarium
相关产品推荐
相关产品推荐

