ConcurrentKafkaListenerContainerFactory配置AbstractConsumerSeekAware后监听器未触发排查
问题描述
我尝试在ConcurrentKafkaListenerContainerFactory上设置通用的MessageListener,代码如下:
@Component @Slf4j public class GossiperMessageListener extends AbstractConsumerSeekAware implements MessageListener<String, String> { @Override public void onPartitionsAssigned( @NotNull Map<TopicPartition, Long> assignments, @NotNull ConsumerSeekCallback callback) { val groupId = KafkaUtils.getConsumerGroupId(); log.info("XXXXX groupId: {} assigned partitions: {}", groupId, assignments); } @Override public void onMessage(ConsumerRecord<String, String> data) { log.info("Received message with Key: {}", data.key()); } }
@Configuration public class KafkaConsumerConfig { public static final String AT_LEAST_ONCE_KAFKA_CONTAINER_FACTORY = "atLeastOnceKafkaListenerContainerFactory"; public Map<String, Object> consumerConfigBase() { val props = new HashMap<String, Object>(); props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:19092"); props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, ErrorHandlingDeserializer.class); props.put(ErrorHandlingDeserializer.VALUE_DESERIALIZER_CLASS, StringDeserializer.class); props.put( ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 1); props.put( ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG, 300000); return props; } @Bean public AbstractKafkaListenerContainerFactory< ConcurrentMessageListenerContainer<String, String>, String, String> atLeastOnceKafkaListenerContainerFactory(GossiperMessageListener listener) { val factory = new ConcurrentKafkaListenerContainerFactory<String, String>(); factory.getContainerProperties().setAckMode(AckMode.RECORD); factory.setConsumerFactory(new DefaultKafkaConsumerFactory<>(consumerConfigBase())); factory.getContainerProperties().setMessageListener(listener); return factory; } }
我能看到分区已分配且可接收Kafka消息,但自定义的GossiperMessageListener并未被调用,日志如下:
2024-06-17T21:54:54.439+05:30 INFO 38597 --- [ntainer#0-0-C-1] o.s.k.l.KafkaMessageListenerContainer : TOPIC_A.group_id: partitions assigned: [TOPIC_A-0] ...
实际项目中有多个标注@KafkaListener的服务类监听不同主题,希望通过ConcurrentKafkaListenerContainerFactory设置ConsumerSeekAware,无需修改所有带@KafkaListener的类。使用版本:spring-kafka 3.0.12。
解决方案
问题根源
当使用@KafkaListener注解时,Spring Kafka会为每个注解方法自动生成对应的MessageListener代理实例,容器工厂中手动设置的全局messageListener会被这个代理覆盖,因此自定义的GossiperMessageListener不会生效。
方案一:通过基类实现全局ConsumerSeekAware
创建一个继承AbstractConsumerSeekAware的基类,实现通用的分区分配逻辑,然后让所有业务监听器继承该基类:
@Slf4j public abstract class BaseKafkaListener extends AbstractConsumerSeekAware { @Override public void onPartitionsAssigned( @NotNull Map<TopicPartition, Long> assignments, @NotNull ConsumerSeekCallback callback) { String groupId = KafkaUtils.getConsumerGroupId(); log.info("XXXXX groupId: {} assigned partitions: {}", groupId, assignments); // 若需要执行seek操作,可通过callback调用,例如:callback.seekToBeginning(assignments.keySet()); } }
业务监听器继承基类即可:
@KafkaListener(topics = "TOPIC_A", containerFactory = "atLeastOnceKafkaListenerContainerFactory") @Slf4j public class BusinessListener extends BaseKafkaListener { @KafkaHandler public void handleMessage(String message) { log.info("Received business message: {}", message); } }
方案二:全局添加ConsumerAwareRebalanceListener
如果不想修改现有业务监听器代码,可以通过ConsumerFactory添加全局的重平衡监听器,实现分区分配监听:
@Configuration public class KafkaConsumerConfig { // ... 保留原有consumerConfigBase()方法 @Bean public AbstractKafkaListenerContainerFactory< ConcurrentMessageListenerContainer<String, String>, String, String> atLeastOnceKafkaListenerContainerFactory() { val factory = new ConcurrentKafkaListenerContainerFactory<String, String>(); factory.getContainerProperties().setAckMode(AckMode.RECORD); DefaultKafkaConsumerFactory<String, String> consumerFactory = new DefaultKafkaConsumerFactory<>(consumerConfigBase()); // 添加全局重平衡监听器 consumerFactory.addListener(new ConsumerAwareRebalanceListener() { @Override public void onPartitionsAssigned(Consumer<?, ?> consumer, Collection<TopicPartition> partitions) { String groupId = consumer.groupMetadata().groupId(); log.info("XXXXX groupId: {} assigned partitions: {}", groupId, partitions); // 直接操作consumer执行seek,例如:consumer.seekToBeginning(partitions); } }); factory.setConsumerFactory(consumerFactory); return factory; } }
关键注意点
- 若仅需要实现
ConsumerSeekAware的分区监听或seek能力,推荐使用方案二,无需修改现有业务代码,侵入性更低。 - 方案一适合需要在业务监听器中复用更多通用逻辑的场景。
内容的提问来源于stack exchange,提问作者tuk
相关产品推荐
相关产品推荐

