【问题标题】:multithreading check membership in Queue and stop the threads多线程检查队列中的成员资格并停止线程
【发布时间】:2015-05-08 12:08:21
【问题描述】:

我想使用 2 个线程遍历一个列表。一个来自前导,另一个来自尾随,并在每次迭代时将元素放在Queue 中。但是在将值放入Queue 之前,我需要检查Queue 中是否存在该值(当其中一个线程将该值放入Queue 时),所以当这种情况发生时我需要停止线程并返回每个线程的遍历值列表。

这是我迄今为止尝试过的:

from Queue import Queue
from threading import Thread, Event

class ThreadWithReturnValue(Thread):
    def __init__(self, group=None, target=None, name=None,
                 args=(), kwargs={}, Verbose=None):
        Thread.__init__(self, group, target, name, args, kwargs, Verbose)
        self._return = None
    def run(self):
        if self._Thread__target is not None:
            self._return = self._Thread__target(*self._Thread__args,
                                                **self._Thread__kwargs)
    def join(self):
        Thread.join(self)
        return self._return

main_path = Queue()

def is_in_queue(x, q):
   with q.mutex:
      return x in q.queue

def a(main_path,g,l=[]):
  for i in g:
    l.append(i)
    print 'a'
    if is_in_queue(i,main_path):
      return l
    main_path.put(i)

def b(main_path,g,l=[]):
  for i in g:
    l.append(i)
    print 'b'
    if is_in_queue(i,main_path):
      return l
    main_path.put(i)

g=['a','b','c','d','e','f','g','h','i','j','k','l']

t1 = ThreadWithReturnValue(target=a, args=(main_path,g))
t2 = ThreadWithReturnValue(target=b, args=(main_path,g[::-1]))
t2.start()
t1.start()
# Wait for all produced items to be consumed
print main_path.join()

我使用了ThreadWithReturnValue,它将创建一个返回值的自定义线程。

对于成员资格检查,我使用了以下功能:

def is_in_queue(x, q):
   with q.mutex:
      return x in q.queue

现在,如果我先启动t1,然后启动t2,我会得到12 个a,然后是一个b,那么它什么也做不了,我需要手动终止python!

但如果我先运行t2 然后t1 我会得到以下结果:

b
b
b
b
 ab

ab
b

b
b
 b
a
a

所以我的问题是,为什么 python 在这种情况下会有所不同?以及如何终止线程并使它们相互通信?

【问题讨论】:

  • 看这里pymotw.com/2/multiprocessing/communication.html ...您对管理共享状态更感兴趣
  • @OWADVL 听起来很有用,我会看到的!谢谢!
  • 您是否有从列表两端进行迭代的实际需求,或者您只是将其作为划分任务的一种方式?

标签: python multithreading python-2.7


【解决方案1】:

在我们遇到更大的问题之前,你没有使用Queue.join 对。

这个函数的全部意义在于,将一堆项目添加到队列中的生产者可以等到消费者完成所有这些项目的工作。这可以通过让消费者在完成使用get 完成的每个项目的工作后致电task_done 来实现。一旦有与put 调用一样多的task_done 调用,队列就完成了。你没有在任何地方做get,更不用说task_done,所以队列永远不可能完成。所以,这就是为什么你在两个线程完成后永远阻塞的原因。


这里的第一个问题是您的线程在实际同步之外几乎没有做任何工作。如果他们所做的唯一一件事就是为了排队而战,那么他们中只有一个能够同时运行。

当然这在玩具问题中很常见,但你必须考虑你真正的问题:

  • 如果您正在执行大量 I/O 工作(侦听套接字、等待用户输入等),线程工作得很好。
  • 如果您正在执行大量 CPU 工作(计算素数),由于 GIL,线程在 Python 中无法工作,但进程可以。
  • 如果您实际上主要处理同步单独的任务,那么任何一个都不会很好地工作(并且流程会更糟)。从线程的角度考虑可能仍然更简单,但它会是最慢的做事方式。您可能想研究协程; Greg Ewing 有一个 great demonstration 说明如何使用 yield from 使用协程来构建调度程序或多角色模拟等内容。

接下来,正如我在您之前的问题中提到的那样,使线程(或进程)在共享状态下有效地工作需要尽可能短的持有锁。

因此,如果您必须在锁定下搜索整个队列,那最好是恒定时间搜索,而不是线性时间搜索。这就是为什么我建议使用类似OrderedSet 配方而不是list 的原因,就像stdlib 中的Queue.Queue 一样。那么这个函数:

def is_in_queue(x, q):
   with q.mutex:
      return x in q.queue

... 只是阻塞队列的一小部分时间——刚好足够在表中查找哈希值,而不是足够长以将队列中的每个元素与x 进行比较。


最后,我试图解释你的另一个问题的竞争条件,但让我再试一次。

您需要锁定代码中的每个完整“事务”,而不仅仅是单个操作。

例如,如果你这样做:

with queue locked:
    see if x is in the queue
if x was not in the queue:
    with queue locked:
        add x to the queue

...那么,当您检查时,x 总是可能不在队列中,但是在您解锁和重新锁定之间的时间里,有人添加了它。这正是两个线程都可能提前停止的原因。

要解决此问题,您需要锁定整个事物:

with queue locked:
    if x is not in the queue:
        add x to the queue

当然,这与我之前所说的尽可能短地锁定队列的说法背道而驰。确实,简而言之,这就是使多线程变得困难的原因。编写安全代码很容易,只要可以想象到有必要就锁定所有内容,但是您的代码最终只使用一个内核,而所有其他线程都被阻塞等待锁定。而且很容易编写快速的代码,尽可能短暂地锁定所有内容,但是这样是不安全的,你会得到垃圾值,甚至到处崩溃。弄清楚什么需要成为一个事务,以及如何最小化这些事务中的工作,以及如何处理多个锁,您可能需要在不使它们死锁的情况下使其工作......这并不容易。

【讨论】:

  • 非常感谢@abarnert 花时间和完整的解释!你澄清了我的一些错误理解,实际上我在多处理方面很糟糕,我想我需要更多的学习! :)
  • @Kasra:“共享内存多线程很难”是陈词滥调是有原因的。每个人都不擅长,因为我们的直觉是错误的(至少在从语言级别到 CPU 微码级别设计以优化单处理且无法更改的系统上)。尽可能使用更高级别的抽象(消息传递而不是共享内存、STM 等),如果不是……嗯,你会逐渐感觉到什么时候应该忽略你的直觉并严格执行(和/或测试) ) 会发生什么,但您仍然会犯难以调试的错误……
  • 是的,我的问题有很多场景!我正在慢慢研究它,我想深入学习它!正如你所说,“共享内存多线程很难”是陈词滥调,对我来说也是如此,但我喜欢困难的事情,我会努力的!
【解决方案2】:

我认为可以改进的几点:

  1. 由于GIL,您可能希望使用multiprocessing(而不是threading)模块。一般来说,CPython 线程不会导致 CPU 密集型工作加速。 (根据您问题的具体内容,multiprocessing 也可能不会,但threading 几乎肯定不会。)
  2. is_inqueue 这样的函数可能会导致高争用。

锁定时间似乎与需要遍历的项目数呈线性关系:

def is_in_queue(x, q):
    with q.mutex:
        return x in q.queue

因此,您可以执行以下操作。

multiprocessing 与共享的dict 一起使用:

 from multiprocessing import Process, Manager

 manager = Manager()
 d = manager.dict()

 # Fn definitions and such

 p1 = Process(target=p1, args=(d,))
 p2 = Process(target=p2, args=(d,))

在每个函数中,像这样检查项目:

def p1(d):

    # Stuff

    if 'foo' in d:
        return 

【讨论】:

    猜你喜欢
    • 2017-05-18
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2013-04-18
    • 1970-01-01
    • 2013-10-10
    • 1970-01-01
    相关资源
    最近更新 更多