Groovy脚本生产者消费者模型:批量消息并行化处理求助
单生产者单消费者场景下批量消息并行处理问题
问题背景
生产者从数据库拉取数据推入阻塞队列,消费者每次取出一批消息后,需并行执行fileExists逻辑;同时想确认是否需要多消费者来完成任务。
初始串行处理代码
ThreadX = Thread.start('producer') { // 从数据库获取数据 while(row){ queue.put(message) } queue.put("KILL") } ThreadY = Thread.start('Consumer') { while(true){ sleep(200) def jsonSlurper = new JsonSlurper() def var = jsonSlurper.parseText(queue.take().toString()) if(var.getAt(0).equals("KILL")) return var.each { fileExists(it) } // 需要并行化此部分 } } boolean fileExists(key){ if(key) { // 业务逻辑 sleep 1000 } }
尝试并行化后仅处理第一批消息的代码
ExecutorService exeSvc = Executors.newFixedThreadPool(5) ThreadY = Thread.start('Consumer') { while(true){ sleep(200) def jsonSlurper = new JsonSlurper() def var = jsonSlurper.parseText(queue.take().toString()) if(var.getAt(0).equals("KILL")) return var.each { exeSvc.execute({-> fileExists(it) sleep(200) }) } } }
问题原因
仅处理第一批消息的核心原因是:
- 消费者收到
KILL信号后直接返回,未等待线程池内已提交的任务执行完毕,JVM主线程退出会直接终止后台线程池任务。 - 不排除消息解析异常或
KILL信号误判的可能,需确认每批消息的格式正确性。
修复方案
方案1:确保线程池任务全部执行完毕再退出
ExecutorService exeSvc = Executors.newFixedThreadPool(5) ThreadY = Thread.start('Consumer') { try { while(true){ // 非必需的sleep(200)建议移除,会增加处理延迟 def jsonSlurper = new JsonSlurper() def var = jsonSlurper.parseText(queue.take().toString()) if(var.getAt(0).equals("KILL")){ break } // 批量提交任务到线程池 var.each { key -> exeSvc.execute({ fileExists(key) }) } } } finally { // 关闭线程池,拒绝新任务 exeSvc.shutdown() // 等待所有已提交任务完成,超时时间可根据业务调整 exeSvc.awaitTermination(1, TimeUnit.HOURS) } } boolean fileExists(key){ if(key) { // 业务逻辑 sleep 1000 } }
方案2:用CompletableFuture简化并行处理(Groovy 2.5+支持)
ThreadY = Thread.start('Consumer') { def futures = [] def threadPool = Executors.newFixedThreadPool(5) try { while(true){ def jsonSlurper = new JsonSlurper() def var = jsonSlurper.parseText(queue.take().toString()) if(var.getAt(0).equals("KILL")){ break } // 并行提交任务 var.each { key -> futures.add(CompletableFuture.runAsync({ fileExists(key) }, threadPool)) } } } finally { // 等待所有任务执行完成 CompletableFuture.allOf(*futures).join() threadPool.shutdown() } }
关于多消费者的建议
- 如果单消费者的并行处理未充分利用资源(比如
fileExists是IO密集型操作,线程池线程多数时间处于等待状态),可以增加消费者线程:每个消费者独立从队列取批量消息,再各自用线程池并行处理。 - 如果单消费者的并行处理已占满资源(比如CPU密集型操作,线程池线程数等于CPU核心数),多消费者会增加上下文切换开销,无需添加。
- 多消费者注意事项:
- 阻塞队列本身线程安全,多消费者
take()操作无冲突。 KILL信号需适配多消费者:可以推送与消费者数量相同的KILL,每个消费者收到后退出;或用计数器标记所有生产者任务完成,消费者检测到后统一退出。
- 阻塞队列本身线程安全,多消费者
内容的提问来源于stack exchange,提问作者deewreck
相关产品推荐
相关产品推荐

