【问题标题】:How do I limit the number of active threads in python?如何限制python中的活动线程数?
【发布时间】:2010-12-19 17:35:33
【问题描述】:

我是 python 新手,在threading 方面取得了一些进展——我正在做一些音乐文件转换,并希望能够利用我机器上的多个内核(每个内核一个活动转换线程)。

class EncodeThread(threading.Thread):
    # this is hacked together a bit, but should give you an idea
    def run(self):
        decode = subprocess.Popen(["flac","--decode","--stdout",self.src],
                            stdout=subprocess.PIPE)
        encode = subprocess.Popen(["lame","--quiet","-",self.dest],
                                stdin=decode.stdout)
        encode.communicate()

# some other code puts these threads with various src/dest pairs in a list

for proc in threads: # `threads` is my list of `threading.Thread` objects
    proc.start()

一切正常,所有文件都已编码,太棒了! ...但是,所有进程都会立即产生,但我只想一次运行两个(每个核心一个)。一个完成后,我希望它移动到列表中的下一个,直到完成,然后继续执行该程序。

我该怎么做?

(我查看了线程池和队列函数,但找不到简单的答案。)

编辑: 也许我应该补充一点,我的每个线程都使用subprocess.Popen 运行一个单独的命令行decoder (flac),通过管道传输到输入的标准输出命令行编码器 (lame/mp3)。

【问题讨论】:

  • 何必呢?让你的线程相互竞争有什么问题?让每个核心都充满工作会更快。
  • 好吧,我想我没有这样想过......拥有超过 2,000 个文件的音乐库,我认为(同时)产生 2,000 个解码过程(flac)到 2,000 个编码过程(跛脚)同时将是次优的。我错了吗?
  • @thornomad:是的,你错了。因为你有 2 个核心而将自己限制在 2 个进程中是错误的。一个过程不会使核心充满工作。即使是由三部分组成的流程管道也可能有足够的 I/O,以至于内核没有被完全占用。
  • 如果您的每个进程需要超过 1/2000 的物理内存,您可能会因为无法分配足够的内存或者更糟糕的机器将交换到死而导致进程死亡...

标签: python multithreading


【解决方案1】:

我想添加一些东西,作为其他人的参考,希望做类似的事情,但他们可能编写了与 OP 不同的东西。这个问题是我在搜索时遇到的第一个问题,选择的答案为我指明了正确的方向。只是想回馈一些东西。

import threading
import time
maximumNumberOfThreads = 2
threadLimiter = threading.BoundedSemaphore(maximumNumberOfThreads)

def simulateThread(a,b):
    threadLimiter.acquire()
    try:
        #do some stuff
        c = a + b
        print('a + b = ',c)
        time.sleep(3)
    except NameError: # Or some other type of error
        # in case of exception, release
        print('some error')
        threadLimiter.release()
    finally:
        # if everything completes without error, release
        threadLimiter.release()
        
        
threads = []
sample = [1,2,3,4,5,6,7,8,9]
for i in range(len(sample)):
    thread = threading.Thread(target=(simulateThread),args=(sample[i],2))
    thread.daemon = True
    threads.append(thread)
    thread.start()
    
for thread in threads:
    thread.join()

这基本上遵循您将在本网站上找到的内容: https://www.kite.com/python/docs/threading.BoundedSemaphore

【讨论】:

    【解决方案2】:

    如果要限制并行线程的数量,请使用semaphore

    threadLimiter = threading.BoundedSemaphore(maximumNumberOfThreads)
    
    class EncodeThread(threading.Thread):
    
        def run(self):
            threadLimiter.acquire()
            try:
                <your code here>
            finally:
                threadLimiter.release()
    

    一次启动所有线程。除了maximumNumberOfThreads 之外的所有线程都将在threadLimiter.acquire() 中等待,并且只有在另一个线程通过threadLimiter.release() 时,等待的线程才会继续。

    【讨论】:

    • 这正好回答了最初的问题。非常适合通过 Google 搜索最终来到这里的人。
    【解决方案3】:

    简答:不要使用线程。

    对于一个工作示例,您可以查看我最近在工作中组合在一起的一些东西。它是ssh 的一个小包装器,它运行可配置数量的Popen() 子进程。我已将其发布在:Bitbucket: classh (Cluster Admin's ssh Wrapper)

    如上所述,我不使用线程;我只是生成孩子,循环调用他们的.poll() 方法并检查超时(也可配置)并在我收集结果时补充池。我玩过不同的 sleep() 值,过去我写过一个版本(在将 subprocess 模块添加到 Python 之前),它使用了 signal 模块( SIGCHLD 和 SIGALRM)以及 os.fork()os.execve() 函数 --- 我的管道和文件描述符管道等)。

    在我的情况下,我会在收集结果时逐步打印结果......并记住所有结果以在最后进行总结(当所有作业都已完成或因超过超时而被终止时)。

    我在一个包含 25,000 个内部主机的列表上运行了它,其中包括 25,000 个内部主机(其中许多已关闭、退役、位于国际上、我的测试帐户无法访问等)。它在两个多小时内完成了工作,没有任何问题。 (其中大约 60 个由于系统处于退化/颠簸状态而超时——证明我的超时处理工作正常)。

    所以我知道这个模型很可靠。使用此代码运行 100 个当前的ssh 进程似乎不会造成任何明显的影响。 (这是一个中等老旧的 FreeBSD 机器)。我曾经在我的旧 512MB 笔记本电脑上运行具有 100 个并发进程的旧(预子进程)版本,也没有问题)。

    (顺便说一句:我计划清理它并为其添加功能;随意贡献或克隆您自己的分支;这就是 Bitbucket.org 的用途)。

    【讨论】:

    • 谢谢 - 今天我会更仔细地研究一下。我很快想出了一组非常简单的 while 循环,似乎只检查 p.communicate() 方法就可以工作。 (PS:我认为您在源代码的第 4 行缺少结束 '''。)
    【解决方案4】:

    “我的每个线程都使用subprocess.Popen 运行单独的命令行[进程]”。

    为什么有一堆线程管理一堆进程?这正是操作系统为您所做的。为什么要对操作系统已经管理的内容进行微观管理?

    与其用线程来监督进程,不如直接派生出进程。您的进程表可能无法处理 2000 个进程,但它可以轻松处理几十个(可能是几百个)。

    您希望更多 工作超出您的 CPU 可能处理的队列。真正的问题是内存问题——而不是进程或线程。如果所有进程的所有活动数据的总和超过物理内存,则必须交换数据,这会减慢您的速度。

    如果您的进程具有相当小的内存占用,您可以运行很多很多。如果您的进程占用大量内存,则不能运行很多。

    【讨论】:

    • 嘿。我现在看到了我被黑在一起的方法的愚蠢——它有点多余。那么,有没有办法通过子流程来管理“池”(正如其他人所建议的那样)。感谢您的输入。边走边学……是否只是使用subprocess.poll() 查看已完成和仍在运行的问题?再次感谢。
    • 正确。您可以使用一组简单的流程;删除已完成的。添加一个并将集合的大小保持在某个限制之下。这只是addremove 的集合。
    【解决方案5】:

    在我看来,您想要的是某种类型的池,并且在该池中您希望有 n 个线程,其中 n == 系统上的处理器数量。然后,您将拥有另一个线程,其唯一的工作是将作业送入队列,工作线程可以在它们空闲时拾取和处理(因此对于双代码机器,您将拥有三个线程,但主线程会做很少)。

    由于您是 Python 新手,但我假设您不知道 GIL 以及它对线程的副作用。如果您阅读我链接的文章,您很快就会明白为什么传统的多线程解决方案并不总是 Python 世界中最好的。相反,您应该考虑使用multiprocessing 模块(Python 2.6 中的新模块,在 2.5 中您可以使用use this backport)来实现相同的效果。它通过使用多个进程来回避 GIL 的问题,就好像它们是同一应用程序中的线程一样。关于如何共享数据(您在不同的内存空间中工作)存在一些限制,但实际上这并不是一件坏事:它们只是鼓励良好的做法,例如最小化线程(或本例中的进程)之间的接触点。

    在您的情况下,您可能对使用here 指定的池感兴趣。

    【讨论】:

    • 谢谢 - 我会看看多进程......我编辑了我的问题以获得更多细节......似乎 subprocess.Popen 确实有点中断并做自己的事情。
    • 多处理模块 BTW 是 2.6 的一个很好的补充(来自支持 2.4 和 2.5 的 pyprocessing 3rd 方模块)。但是,它不太适合运行外部程序。多处理模块的主要优点在于它在线程支持之后建模的方式。您可以创建 Queue()s 作为主要的(线程/进程)间通信机制,以消除对您自己的显式锁定的大部分需求。 (Queue()s 为任意对象的多个生产者和消费者提供一致的支持)。如果孩子们运行 Python 代码,那就太好了。
    【解决方案6】:

    我不是这方面的专家,但我读过一些关于“锁”的东西。 This article 可能会帮到你

    希望对你有帮助

    【讨论】:

      【解决方案7】:

      如果您使用默认的“cpython”版本,那么这对您没有帮助,因为一次只能执行一个线程;查找Global Interpreter Lock。相反,我建议查看 Python 2.6 中的 multiprocessing module ——它使并行编程变得轻而易举。您可以使用2*num_threads 进程创建一个Pool 对象,并为其分配大量任务。它将一次最多执行2*num_threads 个任务,直到全部完成。

      在工作中,我最近迁移了一堆 Python XML 工具(不同的 xpath grepper 和批量 xslt 转换器)来使用它,并且每个处理器有两个进程,结果非常好。

      【讨论】:

      • 如果您的子进程将在您的 Python 代码中执行函数,那么多处理模块非常棒。如果您正在调用外部程序,那么该模块不会比子进程模块提供优势......因为这些外部程序除了临时文件或管道等之外,没有任何方法可以将其结果返回给父级. 多处理模块的巨大 IPC 优势在您执行的外部程序中丢失了。 (例如,让每个进程都在一个多进程调用子进程中听起来很愚蠢)。
      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2018-08-04
      • 1970-01-01
      • 2016-08-25
      • 1970-01-01
      • 1970-01-01
      • 2011-07-11
      • 1970-01-01
      相关资源
      最近更新 更多