Spring Kafka配置批量消费每次仅拉取1条消息问题求解
ConsumerConfig.MAX_POLL_RECORDS_CONFIG配置的是消费者单次poll操作最多可拉取的消息数上限,而非强制拉取的下限,不会等攒够5条消息才返回结果。
你当前仅配置了最大拉取条数、开启了批量监听器,没有配置消息攒批的等待逻辑,加上测试时生产者发送速率设置为每秒10条,消费端poll请求的速度远快于消息生产速度,每次请求打到Broker时分区内只有1条未消费消息,自然每次监听器只能拿到1条数据。
你需要从两个层面补充攒批配置,确保每次拉取尽可能凑够5条消息再投递到监听器:
1. 补充消费者端攒批参数
修改consumerFactory的配置项,追加fetch相关参数,让Broker主动攒够一定量数据再返回给消费者:
@Bean public ConsumerFactory<String, String> consumerFactory(){ Map<String, Object> config = new HashMap<>(); config.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092"); config.put(ConsumerConfig.GROUP_ID_CONFIG, "batch"); config.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); config.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG,StringDeserializer.class); config.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, "5"); // 新增以下配置 // 单次fetch请求最少返回的字节数,单条消息100字节,5条共500字节 config.put(ConsumerConfig.FETCH_MIN_BYTES_CONFIG, "500"); // fetch请求最长等待时间,单位毫秒,最多等3秒,没攒够500字节也会返回现有数据 config.put(ConsumerConfig.FETCH_MAX_WAIT_MS_CONFIG, "3000"); return new DefaultKafkaConsumerFactory<>(config); }
2. 补充监听容器的批量控制参数
修改ConcurrentKafkaListenerContainerFactory的配置,调整poll超时时间,高版本Spring Kafka可直接设置监听器接收的批量大小:
@Bean public ConcurrentKafkaListenerContainerFactory concurrentKafkaListenerContainerFactory() { ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory()); factory.setBatchListener(true); // 新增以下配置 // 单次poll操作的超时时间,和FETCH_MAX_WAIT_MS对齐 factory.getContainerProperties().setPollTimeout(3000); // 批量消费模式下按批次提交偏移量,避免重复消费 factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.BATCH); // Spring Kafka 2.8+ 版本支持直接设置投递到监听器的批量大小,容器层面自动攒够5条再调用监听器 factory.setBatchSize(5); return factory; }
如果你使用的是2.8以下的Spring Kafka版本,没有
setBatchSize方法,靠前面的Broker侧攒批参数也可以实现批量拉取效果。
如果要快速验证批量效果,可以先把生产者测试命令的--throughput参数设为-1(即最大吞吐量一次性发完所有消息),此时分区内存在足够的消息积压,就算不配置等待时间,每次poll也会拉满5条消息。
生产环境建议根据业务可接受的消费延迟调整FETCH_MAX_WAIT_MS参数,一般设置100~1000毫秒即可,避免等待时间过长导致消息延迟升高。
调整配置后重启消费端,即可看到监听器每次接收到的消息列表长度为5(最后一批消息不足5条时会返回实际剩余的消息条数)。
内容的提问来源于stack exchange,提问作者Venkata Krishna Jonnabhatla

