如何使用并控制Quarkus响应式Redis客户端的流水线功能?
Quarkus ReactiveRedisDataSource 流水线功能详解
一、示例代码的流水线判定
你给出的这段代码:
@Inject ReactiveRedisDataSource reactiveDS; // some code below ... reactiveDS.list(String.class).lpush("myList", "randomString") .chain(() -> reactiveDS.hash(String.class).hset("myKey", "field1", "randomValue"))
不会触发流水线处理。因为chain()是串行执行逻辑:前一个命令执行完成后才会发起下一个请求,两个命令是先后发送到Redis服务器的,不存在批量入队后一次性发送的流水线行为。
二、正确实现流水线的方式
要实现Redis命令的流水线,需要基于同一个Redis连接批量提交命令,再统一等待响应。验证可行的实现方式如下:
rds.withConnection(redis -> { List<Uni<Void>> unis = new ArrayList<>(); for (int i = 0; i < 5000; i++) { unis.add(redis.value(Integer.class).set("key-" + i, i)); } return Uni.join().all(unis).andCollectFailures() .replaceWithVoid(); }).await().indefinitely();
核心逻辑说明:
- 使用
withConnection()获取独占的Redis连接,所有流水线命令都基于该连接提交 - 批量创建多个命令对应的
Uni对象并收集到列表中 - 通过
Uni.join().all(unis)合并多个Uni,等待所有命令执行完成 andCollectFailures()可收集所有执行失败的命令,避免单个失败导致整体终止
三、流水线的控制与响应获取
1. 控制发送前的入队命令数量
你可以直接调整批量提交的命令列表大小来控制入队数量(比如示例中的5000)。如果需要精细化流量控制,还可以结合Mutiny的流式操作分批次提交:
// 分批次提交,每批1000个命令 IntStream.range(0, 5000) .boxed() .collect(Collectors.groupingBy(i -> i / 1000)) .values() .forEach(batch -> { rds.withConnection(redis -> { List<Uni<Void>> batchUnis = batch.stream() .map(i -> redis.value(Integer.class).set("key-" + i, i)) .collect(Collectors.toList()); return Uni.join().all(batchUnis).andCollectFailures().replaceWithVoid(); }).await().indefinitely(); });
2. 获取数组响应
如果需要获取每个命令的响应结果,不要调用replaceWithVoid(),直接返回合并后的Uni<List<T>>即可:
List<String> responses = rds.withConnection(redis -> { List<Uni<String>> unis = new ArrayList<>(); unis.add(redis.value(String.class).get("key1")); unis.add(redis.value(String.class).get("key2")); return Uni.join().all(unis).andCollectFailures(); }).await().indefinitely(); // responses列表按顺序保存每个命令的响应结果
内容的提问来源于stack exchange,提问作者A.N.T.
相关产品推荐
相关产品推荐

