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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.22 23:03:23