【问题标题】:How to prevent multiple threads from picking up same task from queue如何防止多个线程从队列中获取相同的任务
【发布时间】:2021-01-17 16:11:30
【问题描述】:

我想并行运行多个线程。每个线程从任务队列中选择一个任务并执行该任务。

from threading import Thread
from Queue import Queue
import time

class link(object):
    
    def __init__(self, i):
        self.name = str(i)

def run_jobs_in_parallel(consumer_func, jobs, results, thread_count, 
                         async_run=False):
    def consume_from_queue(jobs, results):
        while not jobs.empty():
            job = jobs.get()
            try:
                results.append(consumer_func(job))
            except Exception as e:
                print str(e)
                results.append(False)
            finally:
                jobs.task_done()
    #start worker threads
    if jobs.qsize() < thread_count:
        thread_count = jobs.qsize() 
    for tc in range(1,thread_count+1):
        worker = Thread(
            target=consume_from_queue, 
            name="worker_{0}".format(str(tc)), 
            args=(jobs,results,))
        worker.start()
    if not async_run:
        jobs.join()

def create_link(link):
    print str(link.name)
    time.sleep(10)
    return True
    
def consumer_func(link):
    return create_link(link)
    # create_link takes a while to execute

jobs = Queue()
results = list()
for i in range(0,10):
    jobs.put(link(i))
                
run_jobs_in_parallel(consumer_func, jobs, results, 25, async_run=False)

现在发生的事情是,假设我们在作业队列中有 10 个链接对象,当线程并行运行时,多个线程正在执行相同的任务。我怎样才能防止这种情况发生? 注意 - 上面的示例代码没有上面描述的问题,但我有完全相同的代码,除了 create_link 方法做了一些复杂的事情。

【问题讨论】:

  • 请编辑您问题中的代码,使其成为minimal reproducible example - 任何人都应该能够将其粘贴到文件中并不添加任何内容运行它以查看问题你看到了。您必须提供一些 minimal 示例数据,以便代码显示您遇到的问题。您可能需要在适当的位置添加诸如打印之类的内容来显示正在发生的问题,但这就是 minimal reproducible example 的意义所在。

标签: python multithreading python-multithreading


【解决方案1】:

我认为你需要的是一个锁对象(docs,tutorial+examples)。如果您创建此类对象的实例,您可以“锁定”代码的某些部分,确保一次只有一个线程执行该部分。

我猜你想锁定job = jobs.get()这一行。

首先,您必须在所有线程都可以访问它的范围内创建锁。 (您不希望每个线程都拥有一个锁,而是所有线程都需要一个锁。这意味着在获取锁之前在线程中创建锁是行不通的)

import threading    
lock = threading.Lock() 

然后你可以在你的线上使用它:

lock.acquire()
job = jobs.get()
lock.release()

with lock:
  job = jobs.get()

第一个到达acquire() 的线程将锁定锁。尝试acquire() 锁的其他线程将暂停,直到通过调用release() 再次解锁锁。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2020-04-28
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多