Spring Boot+Kafka应用停止消费队列,多线程并发问题求助
问题描述
基于Spring Boot + Kafka开发应用,实现ApplicationReadyEvent监听类ApplicationMainClass后,Kafka的@StreamListener停止消费队列,需排查原因并实现两个线程同时运行。
主应用类代码
@Service public class ApplicationMainClass implements ApplicationListener<ApplicationReadyEvent> { @Autowired PlayerDaoRepository playerDaoRepository; @Autowired DataColectorServiceImp dataColectorServiceImp; @Autowired BattleDaoRepository battleDaoRepository; @Autowired BattleService battleService; private static final Logger log = LogManager.getLogger(ApplicationMainClass.class); @Override public void onApplicationEvent(ApplicationReadyEvent applicationReadyEvent) { List<Playerdao> listPlayersActive; List<BattleDao> battle; List<BattleDao> battleDaoAux; while (true) { log.info("Comienza la ejecución"); listPlayersActive = playerDaoRepository.findByActive(true); for (Playerdao player : listPlayersActive) { try { String battleString = dataColectorServiceImp.apiConexion(player.getUri()); if (battleString.equals("")) { continue; } battle = player.getBatallasPlayed(); battleDaoAux = battleService.getBattle(battleString); player.setLastGamePlayed(!battle.isEmpty() ? battle.get(battle.size()-1).getBattletlime() : LocalDateTime.MIN.toString()); battleDaoAux = player.kafkaHandler(battleDaoAux); battleService.postBattle(battleDaoAux, player.getTag()); player.setBatallasPlayed(player.listBuilder(player.getBatallasPlayed(), battleDaoAux)); battleDaoRepository.saveAll(battle); playerDaoRepository.save(player); } catch (Exception e) { log.error("", e); } } try { log.info("Termina la ejecucion"); Thread.sleep(60000); } catch (InterruptedException e) { } } } }
StreamListener代码
@StreamListener public KStream<IdBattle, textBattle> newBattle(@Input(BinderProcessor.battles) KStream<IdBattle, textBattle> battleKStream){ updateDatabase(battleKStream); return null; } private void updateDatabase(KStream<IdBattle, textBattle> battleKStream) { battleKStream.foreach((IdBattle,textBattle)->{ if(textBattle==null){ playerDaoRepository.deleteById(IdBattle.getIdBattle()); }else{ Battle battle= playerDaoService.textTreatment(textBattle.getText(),Battle); battle=battleDaoService.setBattleTime(textBattle.getText(),battle); Event event = playerDaoService.textTreatmentEvent(textBattle.getText(),Event); battle.setMap(event.getMap()); battleDaoService.updateDatabase(battle,IdBattle.getIdBattle()); } }); }
原因排查
- 主线程被阻塞:
ApplicationReadyEvent的监听方法onApplicationEvent运行在Spring容器初始化的主线程中,方法内的while(true)死循环+Thread.sleep(60000)直接占用了主线程。而Spring Cloud Stream的@StreamListener消费线程的初始化、调度依赖主线程完成后续启动流程,主线程被卡死导致StreamListener无法正常启动,自然停止消费。 - 单线程事件模型冲突:Spring事件监听线程是单线程的,你的循环逻辑完全占用该线程,导致Spring Cloud Stream的初始化步骤无法执行。
解决方法
方法1:将循环逻辑放到独立线程执行
修改ApplicationMainClass,把死循环逻辑提交到Spring线程池或自定义线程中,避免阻塞主线程:
@Service public class ApplicationMainClass implements ApplicationListener<ApplicationReadyEvent> { @Autowired private PlayerDaoRepository playerDaoRepository; @Autowired private DataColectorServiceImp dataColectorServiceImp; @Autowired private BattleDaoRepository battleDaoRepository; @Autowired private BattleService battleService; private static final Logger log = LogManager.getLogger(ApplicationMainClass.class); @Autowired private TaskExecutor taskExecutor; @Override public void onApplicationEvent(ApplicationReadyEvent applicationReadyEvent) { // 把循环逻辑交给独立线程执行 taskExecutor.execute(() -> { List<Playerdao> listPlayersActive; List<BattleDao> battle; List<BattleDao> battleDaoAux; while (true) { log.info("Comienza la ejecución"); listPlayersActive = playerDaoRepository.findByActive(true); for (Playerdao player : listPlayersActive) { try { String battleString = dataColectorServiceImp.apiConexion(player.getUri()); if (battleString.equals("")) { continue; } battle = player.getBatallasPlayed(); battleDaoAux = battleService.getBattle(battleString); player.setLastGamePlayed(!battle.isEmpty() ? battle.get(battle.size()-1).getBattletlime() : LocalDateTime.MIN.toString()); battleDaoAux = player.kafkaHandler(battleDaoAux); battleService.postBattle(battleDaoAux, player.getTag()); player.setBatallasPlayed(player.listBuilder(player.getBatallasPlayed(), battleDaoAux)); battleDaoRepository.saveAll(battle); playerDaoRepository.save(player); } catch (Exception e) { log.error("", e); } } try { log.info("Termina la ejecucion"); Thread.sleep(60000); } catch (InterruptedException e) { Thread.currentThread().interrupt(); // 恢复中断状态 log.error("线程被中断", e); break; } } }); } }
添加线程池配置类(如果Spring未提供默认TaskExecutor):
@Configuration public class ThreadPoolConfig { @Bean public TaskExecutor taskExecutor() { ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); executor.setCorePoolSize(1); executor.setMaxPoolSize(1); executor.setQueueCapacity(10); executor.setThreadNamePrefix("Player-Task-"); executor.initialize(); return executor; } }
方法2:改用@Scheduled注解替代死循环
你的逻辑是每隔60秒执行一次,用Spring定时任务注解更符合规范,无需手动管理线程:
@Service @EnableScheduling // 开启定时任务支持 public class ApplicationMainClass { @Autowired private PlayerDaoRepository playerDaoRepository; @Autowired private DataColectorServiceImp dataColectorServiceImp; @Autowired private BattleDaoRepository battleDaoRepository; @Autowired private BattleService battleService; private static final Logger log = LogManager.getLogger(ApplicationMainClass.class); @Scheduled(fixedDelay = 60000) // 每次执行完成后间隔60秒再执行 public void executePlayerTask() { log.info("Comienza la ejecución"); List<Playerdao> listPlayersActive = playerDaoRepository.findByActive(true); for (Playerdao player : listPlayersActive) { try { String battleString = dataColectorServiceImp.apiConexion(player.getUri()); if (battleString.equals("")) { continue; } List<BattleDao> battle = player.getBatallasPlayed(); List<BattleDao> battleDaoAux = battleService.getBattle(battleString); player.setLastGamePlayed(!battle.isEmpty() ? battle.get(battle.size()-1).getBattletlime() : LocalDateTime.MIN.toString()); battleDaoAux = player.kafkaHandler(battleDaoAux); battleService.postBattle(battleDaoAux, player.getTag()); player.setBatallasPlayed(player.listBuilder(player.getBatallasPlayed(), battleDaoAux)); battleDaoRepository.saveAll(battle); playerDaoRepository.save(player); } catch (Exception e) { log.error("", e); } } log.info("Termina la ejecucion"); } }
额外注意事项
- 确保
@StreamListener所在类添加了@Service或@Component注解,被Spring扫描到。 - 检查Kafka Stream配置(如topic、consumer group)是否正确生效。
- 线程中断处理:自定义线程时需正确恢复中断状态,避免后续逻辑异常。
内容的提问来源于stack exchange,提问作者miguel30452
相关产品推荐
相关产品推荐

