Spring Kafka无消息5分钟时,如何调用API停止当前Listener并启动新Listener?
Spring Kafka 无消息时切换监听器实现方案
核心思路
通过定时任务监控活跃监听器的最后消费时间,当连续5分钟无消息消费时,利用KafkaListenerEndpointRegistry控制监听器的启停,完成切换。
具体实现步骤
1. 定义带ID的Kafka监听器
为两个监听器指定唯一ID,方便后续通过ID控制启停;初始只启动第一个监听器,第二个设为autoStartup=false:
import org.apache.kafka.clients.consumer.ConsumerRecord; import org.springframework.kafka.annotation.KafkaListener; import org.springframework.stereotype.Component; @Component public class KafkaListeners { @KafkaListener(id = "primary-listener", topics = "topic-primary", autoStartup = "true") public void listenPrimary(ConsumerRecord<String, String> record) { // 业务消费逻辑 LastConsumeRecorder.update("primary-listener"); } @KafkaListener(id = "secondary-listener", topics = "topic-secondary", autoStartup = "false") public void listenSecondary(ConsumerRecord<String, String> record) { // 第二个监听器的业务逻辑 LastConsumeRecorder.update("secondary-listener"); } }
2. 记录最后消费时间
用线程安全的容器记录每个监听器的最后消费时间,确保多线程环境下的准确性:
import java.util.Map; import java.util.concurrent.ConcurrentHashMap; public class LastConsumeRecorder { private static final Map<String, Long> LAST_CONSUME_TIMES = new ConcurrentHashMap<>(); public static void update(String listenerId) { LAST_CONSUME_TIMES.put(listenerId, System.currentTimeMillis()); } public static long get(String listenerId) { return LAST_CONSUME_TIMES.getOrDefault(listenerId, 0L); } }
3. 定时检查并切换监听器
使用Spring定时任务周期检查活跃监听器的闲置状态,达到阈值时执行切换:
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Component; import org.springframework.kafka.config.KafkaListenerEndpointRegistry; @Component public class ListenerSwitcher { @Autowired private KafkaListenerEndpointRegistry listenerRegistry; // 闲置阈值:5分钟(毫秒) private static final long INACTIVE_THRESHOLD = 5 * 60 * 1000; // 检查周期:每分钟一次 private static final long CHECK_INTERVAL = 60 * 1000; @Scheduled(fixedRate = CHECK_INTERVAL) public void checkAndSwitch() { String activeId = "primary-listener"; String targetId = "secondary-listener"; long lastConsumeTime = LastConsumeRecorder.get(activeId); long now = System.currentTimeMillis(); if (now - lastConsumeTime >= INACTIVE_THRESHOLD) { // 停止当前活跃监听器 listenerRegistry.getListenerContainer(activeId).stop(); // 启动目标监听器 listenerRegistry.getListenerContainer(targetId).start(); // 更新记录,避免重复触发 LastConsumeRecorder.update(targetId); } } }
4. 开启定时任务支持
在Spring配置类上添加@EnableScheduling注解,启用定时任务功能:
import org.springframework.scheduling.annotation.EnableScheduling; import org.springframework.context.annotation.Configuration; @Configuration @EnableScheduling public class SchedulerConfig { }
注意事项
- 若需要双向切换(比如第二个监听器闲置后切回第一个),可扩展定时任务逻辑,同时监控两个监听器的状态。
- 监听器启停操作是异步的,可通过
container.isRunning()方法校验状态。 - 生产环境建议添加日志记录,便于排查切换逻辑的执行情况。
内容的提问来源于stack exchange,提问作者Arjun
相关产品推荐
相关产品推荐

