【问题标题】:How to manage python threads results?如何管理python线程结果?
【发布时间】:2010-07-13 17:15:36
【问题描述】:

我正在使用此代码:

def startThreads(arrayofkeywords):
    global i
    i = 0
    while len(arrayofkeywords):
        try:
            if i<maxThreads:
                keyword = arrayofkeywords.pop(0)
                i = i+1
                thread = doStuffWith(keyword)
                thread.start()
        except KeyboardInterrupt:
            sys.exit()
    thread.join()

对于python中的线程,我几乎已经完成了所有工作,但我不知道如何管理每个线程的结果,在每个线程上我都有一个字符串数组作为结果,我怎样才能安全地将所有这些数组合并为一个?因为,如果我尝试写入全局数组,两个线程可能会同时写入。

【问题讨论】:

    标签: python multithreading arrays


    【解决方案1】:

    首先,您实际上需要保存所有那些thread 对象以调用join()。正如所写,您只保存最后一个,并且只有在没有异常的情况下才保存。

    进行多线程编程的一种简单方法是为每个线程提供运行所需的所有数据,然后让它不写入该工作集之外的任何内容。如果所有线程都遵循该准则,则它们的写入不会相互干扰。然后,一旦一个线程完成,让仅主线程将结果聚合到一个全局数组中。这被称为“fork/join 并行性”。

    如果您将 Thread 对象子类化,您可以给它空间来存储该返回值,而不会干扰其他线程。然后你可以这样做:

    class MyThread(threading.Thread):
        def __init__(self, ...):
            self.result = []
            ...
    
    def main():
        # doStuffWith() returns a MyThread instance
        threads = [ doStuffWith(k).start() for k in arrayofkeywords[:maxThreads] ]
        for t in threads:
            t.join()
            ret = t.result
            # process return value here
    

    编辑:

    看了一圈,好像上面的方法isn't the preferred way to do threads in Python。上面更多的是线程的Java-esque模式。相反,您可以执行以下操作:

    def handler(outList)
        ...
        # Modify existing object (important!)
        outList.append(1)
        ...
    
    def doStuffWith(keyword):
        ...
        result = []
        thread = Thread(target=handler, args=(result,))
        return (thread, result)
    
    def main():
        threads = [ doStuffWith(k) for k in arrayofkeywords[:maxThreads] ]
        for t in threads:
            t[0].start()
        for t in threads:
            t[0].join()
            ret = t[1]
            # process return value here
    

    【讨论】:

    • 谢谢,你的回答是我最容易理解的,马上试试。
    【解决方案2】:

    使用Queue.Queue 实例,它本质上是线程安全的。每个线程在完成后可以.put 将其结果传递给该全局实例,并且主线程(当它知道所有工作线程都已完成时,通过.joining 他们,例如@unholysampler 的答案)可以循环.getting从中得到每个结果,并将每个结果用于.extend“整体结果”列表,直到队列被清空。

    编辑:您的代码还有其他大问题——如果最大线程数小于关键字的数量,它将永远不会终止(您试图在每个线程中启动一个线程)关键字 - 永远不会少 - 但如果你已经开始了最大数量,你会永远循环到没有进一步的目的)。

    考虑改为使用 线程池,有点像 this recipe 中的那个,除了你将排队关键字而不是排队可调用对象 - 因为你想在其中运行可调用对象每个线程中的线程都是相同的,只是改变了参数。当然,该可调用对象将被更改为从传入任务队列(使用.get)和.put 完成后将结果列表剥离到传出结果队列。

    要终止 N 个线程,你可以在所有关键字之后,.putN 个“哨兵”(例如None,假设没有关键字可以是None):如果“关键字”它只是一个线程的可调用对象将退出拉取的是None

    通常情况下,Queue.Queue 提供了在 Python 中组织线程(和多处理!)架构的最佳方式,无论它们是像我指出的食谱中那样通用的,还是像我建议您使用的那样更专业最后两段的大小写。

    【讨论】:

      【解决方案3】:

      您需要保留指向您创建的每个线程的指针。照原样,您的代码仅确保最后创建的线程完成。这并不意味着您在它之前开始的所有那些也都完成了。

      def startThreads(arrayofkeywords):
          global i
          i = 0
          threads = []
          while len(arrayofkeywords):
              try:
                  if i<maxThreads:
                      keyword = arrayofkeywords.pop(0)
                      i = i+1
                      thread = doStuffWith(keyword)
                      thread.start()
                      threads.append(thread)
              except KeyboardInterrupt:
                  sys.exit()
          for t in threads:
              t.join()
          //process results stored in each thread
      

      这也解决了写访问的问题,因为每个线程都会在本地存储它的数据。然后在所有这些都完成之后,您就可以进行合并每个线程本地数据的工作了。

      【讨论】:

      • 我将如何访问每个线程的本地数据?因为线程一直在创建/完成,并不像总是相同的 10 个线程。
      • 根据本地数据,我指的是 Karmastan 的解决方案。根据您在问题中发布的内容,您似乎创建了 N 个线程,然后开始,然后加入它们。鉴于在线程完成后访问本地数据的模式可以正常工作。如果您希望事情变得更加动态,那么您将需要查看讨论线程池并将结果存储在数据队列中的答案。
      【解决方案4】:

      我知道这个问题有点老了,但是最好的办法就是不要像其他同事提出的那样伤害自己太多:)

      请阅读Pool 上的参考资料。这样你就可以分叉加入你的工作:

      def doStuffWith(keyword):
          return keyword + ' processed in thread'
      
      def startThreads(arrayofkeywords):
          pool = Pool(processes=maxThreads)
          result = pool.map(doStuffWith, arrayofkeywords)
          print result
      

      【讨论】:

      • 值得注意的是Pool使用multiprocessing而不是threading,它为每个进程使用单独的内存空间,因此比线程慢。使用from multiprocessing.pool import ThreadPool 等效于Pool
      【解决方案5】:

      如果您使用信号量来保护临界区,则写入全局数组是可以的。当您想附加到全局数组时,您“获取”锁,然后在完成后“释放”。这样,每次只有一个线程附加到数组中。

      查看http://docs.python.org/library/threading.html 并搜索信号量以获取更多信息。

      sem = threading.Semaphore()
      ...
      sem.acquire()
      # do dangerous stuff
      sem.release()
      

      【讨论】:

        【解决方案6】:

        尝试一些信号量的方法,比如获取和释放.. http://docs.python.org/library/threading.html

        【讨论】:

          猜你喜欢
          • 2022-01-10
          • 1970-01-01
          • 1970-01-01
          • 1970-01-01
          • 2016-10-28
          • 1970-01-01
          • 1970-01-01
          • 2021-04-02
          • 2014-11-24
          相关资源
          最近更新 更多