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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 07:34:58