You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.26 17:12:52