如何使用Redisson Stream实现非阻塞读取?
用Redisson实现Redis Stream的非阻塞队列读取
方案1:用带阻塞参数的read方法替代轮询
Redis原生Stream支持XREAD BLOCK阻塞读取命令,Redisson的RStream提供了对应重载方法,能直接实现无轮询的阻塞读取,不用自己加Thread.sleep:
RStream<Object, Object> stream = redissonClient.getStream("my-stream"); // 设置阻塞时长,0表示一直阻塞直到有消息返回 StreamReadParams params = StreamReadParams.defaultParams().block(Duration.ofMillis(0)); // 从最新消息位置开始读取,也可以指定具体的StreamMessageId Map<StreamMessageId, Map<Object, Object>> messages = stream.read(params, StreamOffset.lastEntry("my-stream")); // 处理读取到的消息 for (var entry : messages.entrySet()) { StreamMessageId msgId = entry.getKey(); Map<Object, Object> content = entry.getValue(); // 你的业务处理逻辑 }
如果需要持续监听,把这段代码放到循环里即可——因为read方法本身会阻塞,没有消息时会挂起在Redis端,直到有新消息才返回,完全避免无效轮询。
方案2:异步监听器模式(无需主动循环)
Redisson还支持异步监听,后台自动处理阻塞读取,有新消息时直接回调处理逻辑:
RStream<Object, Object> stream = redissonClient.getStream("my-stream"); // 添加监听器,监听从当前最新位置开始的新消息 int listenerId = stream.addListener(StreamOffset.lastEntry("my-stream"), new StreamListener<Object, Object>() { @Override public void onMessage(StreamMessageId msgId, Map<Object, Object> msgContent) { // 收到新消息后的处理逻辑 System.out.println("新消息内容:" + msgContent); } }); // 后续不需要监听时可以移除 // stream.removeListener(listenerId);
这种方式完全不用自己写循环,Redisson会在后台维护连接处理阻塞逻辑,消息到达时自动触发回调。
多消费者场景:用消费组实现阻塞读取
如果是多消费者协作的队列场景,推荐用Redis Stream的消费组功能,Redisson同样支持阻塞读取:
RStream<Object, Object> stream = redissonClient.getStream("my-stream"); // 先创建消费组(仅需执行一次,重复调用会报错,可加判断) stream.createGroup("my-consumer-group", StreamOffset.lastEntry("my-stream")); // 从消费组中阻塞读取消息,指定消费者ID StreamReadParams params = StreamReadParams.defaultParams().block(Duration.ofMillis(0)).count(10); Map<StreamMessageId, Map<Object, Object>> messages = stream.readGroup( "my-consumer-group", "consumer-001", params, StreamOffset.lastEntry("my-stream") ); // 处理完消息后确认,避免重复消费 for (StreamMessageId msgId : messages.keySet()) { stream.ack("my-consumer-group", msgId); }
内容的提问来源于stack exchange,提问作者Qifan.y
相关产品推荐
相关产品推荐

