Spring Boot Data Redis Stream批量消费疑问:Imperative容器如何接收批量消息
解答
1. 澄清batchSize参数的作用
你之前确实误解了batchSize的作用:它控制的是每次从Redis Stream拉取的最大消息数量,但默认的StreamListener接口会把拉取到的消息逐条分发给onMessage方法,所以你看到的是单条消息接收。
2. 实现批量接收的正确方式
要通过StreamMessageListenerContainer接收批量消息,需要替换为BatchMessageListener接口,它的onMessage方法直接接收List<MapRecord>类型的批量消息。
修改后的Batch监听器代码
public class RedisStreamBatchListener implements BatchMessageListener<String, MapRecord<String, Object, Object>> { @Value("${redis-stream-name2}") private String redisStreamName; @Value("${redis-stream-consumer2}") private String consumerGroup; private static final Logger LOG = LoggerFactory.getLogger(RedisStreamBatchListener.class); @Autowired private RedisTemplate redisTemplate; @Override public void onMessage(List<MapRecord<String, Object, Object>> messages) { LOG.debug("Received batch messages from Redis, total count: {}", messages.size()); for (MapRecord<String, Object, Object> message : messages) { String stream = message.getStream(); RecordId id = message.getId(); String data = (String) message.getValue().get("testdata"); try { EventEFRModel eventEFRModel = CommonUtils.readStringToObject(data, EventEFRModel.class); LOG.debug("Consumed data: {}", eventEFRModel); acknowledgeAndDeleteRecord(stream, stream, id); } catch (JsonProcessingException e) { e.printStackTrace(); } } } }
修改后的订阅代码
只需要调整监听器注册的部分,传入BatchMessageListener实例即可:
public Subscription subscribeToEvents(RedisConnectionFactory redisConnectionFactory) throws InterruptedException { StreamMessageListenerContainer<String, MapRecord<String, Object, Object>> listenerContainer = StreamMessageListenerContainer.create(Objects.requireNonNull(redisConnectionFactory), StreamMessageListenerContainer.StreamMessageListenerContainerOptions.builder() .hashKeySerializer(new Jackson2JsonRedisSerializer<>(String.class)) .hashValueSerializer(new Jackson2JsonRedisSerializer<>(Object.class)) .pollTimeout(Duration.ofSeconds(1)) .batchSize(3) // 现在该参数会生效,每次拉取最多3条消息 .build()); StreamMessageListenerContainer.StreamReadRequest<String> streamReadRequest = StreamMessageListenerContainer.StreamReadRequest .builder(StreamOffset.create(redisStreamName, ReadOffset.lastConsumed())) .consumer(Consumer.from(consumerGroup, redisStreamName)) .cancelOnError(ex -> false) .autoAcknowledge(false) .build(); // 注册BatchMessageListener,而非普通StreamListener Subscription subscription = listenerContainer.register(streamReadRequest, new RedisStreamBatchListener()); listenerContainer.start(); LOGGER.info("Registered subscriber for stream key: {} with consumer group: {} & consumer name: {}", redisStreamName, consumerGroup); return subscription; }
3. 两种批量方式对比
StreamMessageListenerContainer + BatchMessageListener:框架自动异步轮询,无需手动添加调度器,适合长期持续监听的场景,batchSize控制每次拉取的消息上限。RedisTemplate + 调度器:需要自己实现定时拉取逻辑,灵活性更高,但需要处理调度、异常重试等额外逻辑。
内容的提问来源于stack exchange,提问作者Pritam Kumar
相关产品推荐
相关产品推荐

