【发布时间】:2022-11-10 12:11:15
【问题描述】:
我有一个生产者消费者问题,我有单个生产者推入阻塞队列和单个消费者从队列中消费。一旦使用了一条消息,我就会对该批消息进行多项操作。如何并行处理每批消息的逻辑处理。下面是代码sn-p。还建议我是否应该考虑多个消费者来完成这项任务。
ThreadX = Thread.start('producer') {
//data retrieve from DB
while(row){
queue.put(message)
}
queue.put("KILL")
}
ThreadY = Thread.start('Consumer') {
while(true){
sleep(200)
// print(Thread.currentThread().name)
def jsonSlurper = new JsonSlurper()
def var = jsonSlurper.parseText(queue.take().toString())
if(var.getAt(0).equals("KILL"))
return
var.each { fileExists(it) } // **need parallelize this part**
}
boolean fileExists(key){
if(key) {
//some logic
sleep 1000
}
}
}
更新:尝试了以下代码,但它以某种方式只处理消费者消费的第一批 10 条消息
ExecutorService exeSvc = Executors.newFixedThreadPool(5)
ThreadY = Thread.start('Consumer') {
while(true){
sleep(200)
// print(Thread.currentThread().name)
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)
}) }
}
}
请帮忙
【问题讨论】:
标签: multithreading groovy parallel-processing