【问题标题】:Python threading.Thread can be stopped only with private method self.__Thread_stop()Python threading.Thread 只能使用私有方法 self.__Thread_stop() 停止
【发布时间】:2011-10-06 21:20:30
【问题描述】:

我有一个函数,它接受一个大数组 x,y 对作为输入,它使用 numpy 和 scipy 进行一些复杂的曲线拟合,然后返回一个值。为了尝试加快速度,我尝试使用两个线程将数据提供给使用 Queue.Queue 。一旦数据完成。我试图让线程终止,然后结束调用进程并将控制权返回给 shell。

我试图理解为什么我必须求助于 threading.Thread 中的私有方法来停止我的线程并将控制权返回给命令行。

self.join() 不会结束程序。重新获得控制权的唯一方法是使用私有停止方法。

        def stop(self):
            print "STOP CALLED"
            self.finished.set()
            print "SET DONE"
            # self.join(timeout=None) does not work
            self._Thread__stop()

这是我的代码的近似值:

    class CalcThread(threading.Thread):
        def __init__(self,in_queue,out_queue,function):
            threading.Thread.__init__(self)
            self.in_queue = in_queue
            self.out_queue = out_queue
            self.function = function
            self.finished = threading.Event()

        def stop(self):
            print "STOP CALLED"
            self.finished.set()
            print "SET DONE"
            self._Thread__stop()

        def run(self):
            while not self.finished.isSet():
                params_for_function = self.in_queue.get()
                try:
                    tm = self.function(paramsforfunction)
                    self.in_queue.task_done()
                    self.out_queue.put(tm)
                except ValueError as v:
                    #modify params and reinsert into queue
                    window = params_for_function["window"]
                    params_for_function["window"] = window + 1
                    self.in_queue.put(params_for_function)

    def big_calculation(well_id,window,data_arrays):
            # do some analysis to calculate tm
            return tm

    if __name__ == "__main__":
        NUM_THREADS = 2
        workers = []
        in_queue = Queue()
        out_queue = Queue()

        for i in range(NUM_THREADS):
            w = CalcThread(in_queue,out_queue,big_calculation)
            w.start()
            workers.append(w)

        if options.analyze_all:
              for i in well_ids:
                  in_queue.put(dict(well_id=i,window=10,data_arrays=my_data_dict))

        in_queue.join()
        print "ALL THREADS SEEM TO BE DONE"
        # gather data and report it from out_queue
        for i in well_ids:
            p = out_queue.get()
            print p
            out_queue.task_done()
            # I had to do this to get the out_queue to proceed
            if out_queue.qsize() == 0:
                out_queue.join()
                break
# Calling this stop method does not seem to return control to the command line unless I use threading.Thread private method

        for aworker in workers:
            aworker.stop()

【问题讨论】:

  • sys.exit()(仅杀死线程)
  • self.daemon = True 仅在调用 start() 之前设置有效,否则引发 RuntimeError
  • sys.exit() 不会杀死线程,但会在当前线程中引发 SystemExit 异常。

标签: python multithreading queue


【解决方案1】:

一般来说,杀死修改共享资源的线程是个坏主意。

除非在执行计算时释放 GIL,否则多线程中的 CPU 密集型任务在 Python 中比无用更糟糕。许多numpy 函数确实发布了 GIL。

ThreadPoolExecutor example from the docs

import concurrent.futures # on Python 2.x: pip install futures 

calc_args = []
if options.analyze_all:
    calc_args.extend(dict(well_id=i,...) for i in well_ids)

with concurrent.futures.ThreadPoolExecutor(max_workers=NUM_THREADS) as executor:
    future_to_args = dict((executor.submit(big_calculation, args), args)
                           for args in calc_args)

    while future_to_args:
        for future in concurrent.futures.as_completed(dict(**future_to_args)):
            args = future_to_args.pop(future)
            if future.exception() is not None:
                print('%r generated an exception: %s' % (args,
                                                         future.exception()))
                if isinstance(future.exception(), ValueError):
                    #modify params and resubmit
                    args["window"] += 1
                    future_to_args[executor.submit(big_calculation, args)] = args

            else:
                print('f%r returned %r' % (args, future.result()))

print("ALL work SEEMs TO BE DONE")

如果没有共享状态,您可以将 ThreadPoolExecutor 替换为 ProcessPoolExecutor。将代码放入您的 main() 函数中。

【讨论】:

  • 哇,这让我大开眼界。非常感谢您向我介绍 concurrent.futures。它适用于 python 2.7 和 numpy 和 scipy。没有 thread.Threading 的麻烦和所有并发执行的好处
【解决方案2】:

详细说明我的评论 - 如果您的线程的唯一目的是使用队列中的值并对它们执行功能,那么您最好做这样的事情恕我直言:

q = Queue()
results = []

def worker():
  while True:
    x, y = q.get()
    results.append(x ** y)
    q.task_done()

for _ in range(workerCount):
  t = Thread(target = worker)
  t.daemon = True
  t.start()

for tup in listOfXYs:
  q.put(tup)

q.join()

# Some more code here with the results list.

q.join() 将阻塞,直到它再次为空。工作线程将继续尝试检索值,但找不到任何值,因此一旦队列为空,它们将无限期地等待。当您的脚本稍后完成执行时,工作线程将死亡,因为它们被标记为守护线程。

【讨论】:

  • 您可以使用哨兵值,而不是为这些东西使用守护进程(对于这种情况,imo 不是一个好的设计,YMMV)。 IE。所有作业完成后,将nrThreads 标记值放入队列,然后再次加入队列或线程。线程只是检查get() 是否返回了哨兵(通常没有是一个好的选择)并在这种情况下停止。还可以更轻松地将代码包含在更大的设计中。
  • @Voo:工作线程本身正在将新值放入in_queue。如果主线程将哨兵放入in_queue,它们可能会提前发出终止信号。您将如何处理这种情况?
  • @unutbu - 我个人没有看到哨兵值的优势,但您可以(理论上)通过使用 LifoQueue 代替标准队列并使用预填充它来解决这个问题每个工作线程的哨兵值。这确实有可能(至少在 op 的情况下)你的一些工作人员提前死亡,但最终重新添加到 in_queue 的工作人员最终运行时间明显更长。被阻塞的queue.get() 中的守护线程几乎不消耗资源,而且根据我的经验,它不会消耗性能。
  • @g.d.d.c:我喜欢你使用 LifoQueue 的想法。我认为这可能是可行的。但是还有另一个问题:如何知道out_queue 何时为空。我不认为测试qsize 是安全的——一个线程可能即将put out_queue 中的一个新项目,而主线程正在测试qsize,此时它暂时为零。
  • 我显然没有很好地解释这个概念,但是是的,g.d.d.c 是对的。等待,将哨兵放入队列,再次等待(尽管您也可以在线程/进程/线程池/第二次等待;语义不完全相同但足够接近)。对于这类问题,这是一个非常有用的模式 - 可以通过多个输入队列、输出队列等变得更加复杂,但这是野兽的本质,我们也可以推广解决方案。
【解决方案3】:

我尝试了 g.d.d.c 的方法,它产生了一个有趣的结果。我可以让他精确的 x**y 计算在线程之间很好地分布。

当我在工作线程中调用我的函数时,while True 循环。只有在调用线程 start() 方法的 for 循环中放入 time.sleep(1) 才能在多个线程之间执行计算。

所以在我的代码中。没有 time.sleep(1) 程序给了我一个没有输出的干净退出,或者在某些情况下

“线程 Thread-2 中的异常(很可能在解释器关闭期间引发):线程 Thread-1 中的异常(很可能在解释器关闭期间引发):”

添加 time.sleep() 后,一切正常。

for aworker in range(5):
    t = Thread(target = worker)
    t.daemon = True
    t.start()
    # This sleep was essential or results for my specific function were None
    time.sleep(1)
    print "Started"

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2013-01-07
    • 2013-11-06
    • 1970-01-01
    • 2013-03-14
    • 2013-09-11
    • 1970-01-01
    • 2013-11-08
    • 2017-05-29
    相关资源
    最近更新 更多