MY kafka Listener detail as below:
@KafkaListener(
topics="test",
groupId="groupid",
containerFactory="kafkaBatchListenerContainerFactory" )
public void onMessage(List<ConsumerRecord<String,String>> messages, Acknowledgment acknowledgment) throws Exception {
log.info("Batch size is {}",messages.size());
acknowledgment.acknowledge();
}
Bean Configuration detail:
@Bean
public ConcurrentKafkaListenerContainerFactory<String, String> kafkaBatchListenerContainerFactory() throws IOException {
ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>();
factory.setConsumerFactory(consumerFactory());
factory.setConcurrency(3); //Equal to partition
factory.getContainerProperties().setAckMode("MANUAL");
factory.getContainerProperties().setIdleBetweenPolls((long) 1 * 1000L);
factory.getContainerProperties().setPollTimeout(kafkaConfigProperties.getPollTimeout() * 1000L);
factory.setMissingTopicsFatal(false);
factory.setBatchListener(true);
factory.getContainerProperties().setAsyncAcks(true);
factory.setCommonErrorHandler(new DefaultErrorHandler(new FixedBackOff(1000L, 4)));
}
Consumer Factory:
@Bean
public ConsumerFactory<String, String> consumerFactory() throws IOException {
Map<String, Object> config = new HashMap<>();
config.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9080");
config.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, ErrorHandlingDeserializer.class);
config.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, ErrorHandlingDeserializer.class);
config.put(ErrorHandlingDeserializer.KEY_DESERIALIZER_CLASS, StringDeserializer.class);
config.put(ErrorHandlingDeserializer.VALUE_DESERIALIZER_CLASS, JsonDeserializer.class.getName());
config.put(ConsumerConfig.GROUP_ID_CONFIG, "groupid");
config.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false);
config.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
config.put(ConsumerConfig.REQUEST_TIMEOUT_MS_CONFIG, 60 * 1000);
config.put(ConsumerConfig.FETCH_MAX_WAIT_MS_CONFIG, 2 * 1000);
config.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, 30 * 1000);
config.put(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG, 900 * 1000);
config.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 100);
config.put(ConsumerConfig.ISOLATION_LEVEL_CONFIG, read_committed);
return new DefaultKafkaConsumerFactory<>(config, new StringDeserializer(), new StringDeserializer());
}
When stop and start the server every time, able to receive batch of records. However during live transaction, received one record instead of batch. Let me know any suggestion.
My logs as below :
Batch size is 1
Batch size is 1
Batch size is 1
Batch size is 1