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

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());
        }
    });
}

原因排查

  1. 主线程被阻塞:ApplicationReadyEvent的监听方法onApplicationEvent运行在Spring容器初始化的主线程中,方法内的while(true)死循环+Thread.sleep(60000)直接占用了主线程。而Spring Cloud Stream的@StreamListener消费线程的初始化、调度依赖主线程完成后续启动流程,主线程被卡死导致StreamListener无法正常启动,自然停止消费。
  2. 单线程事件模型冲突: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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.10 08:50:27