【问题标题】:Python multiproccessing memory increasePython多处理内存增加
【发布时间】:2014-11-06 07:47:50
【问题描述】:

我有一个应该永远运行的程序。 这是我正在做的事情:

from myfuncs import do, process

class Worker(multiprocessing.Process):

    def __init__(self, lock):
        multiprocesing.Process.__init__(self)
        self.lock = lock
        self.queue = Redis(..) # this is a redis based queue
        self.res_queue = Redis(...)

     def run():
         while True:
             job = self.queue.get(block=True)
             job.results = process(job)
             with self.lock:
                 post_process(self.res_queue, job)


def main():
    lock = multiprocessing.Semaphore(1)
    ps = [Worker(lock) for _ in xrange(4)]
    [p.start() for p in ps]
    [p.join() for p in ps]

self.queue 和 self.res_queue 是两个像 python stdlib Queue 一样工作的对象,但它们 使用 Redis 数据库作为后端。

函数进程对job携带的数据(主要是html解析)和返回的数据做一些处理 一本字典。

函数 post_process 通过检查某些条件将作业写入另一个 redis 队列(一次只有一个进程可以检查导致锁定的条件)。它返回真/假。

程序每天使用的内存在增加。 有人能弄清楚发生了什么吗?

当作业超出运行方法的范围时,内存应该是空闲的吗?

【问题讨论】:

  • 是什么让您确定保留的是作业对象?你是用tracemalloc,还是在调试器中扫描gc堆,还是只是猜测?
  • 现在只是猜测

标签: python memory multiprocessing


【解决方案1】:

当作业超出运行方法的范围时,内存应该是空闲的吗?

首先,作用域是整个run 方法,它永远循环,因此永远不会发生。 (此外,当您退出 run 方法时,进程将关闭并且它的内存无论如何都会被释放......)

但即使它确实超出了范围,也并不意味着您似乎认为它意味着什么。 Python 不像 C++,其中的变量存储在堆栈上。所有对象都存在于堆中,并且它们一直存在,直到不再有对它们的引用。超出范围的变量意味着该变量不再引用它曾经引用的任何对象。如果该变量是对该对象的唯一引用,那么它将被释放*,但如果您在其他地方进行了其他引用,则在这些其他引用消失之前无法释放该对象。

同时,超出范围并没有什么神奇之处。变量停止引用对象的任何方式都具有相同的效果——无论是变量超出范围,您在其上调用del,还是为它分配一个新值。因此,每次通过循环时,当您执行job = 时,您将删除先前对job 的引用,即使没有超出范围。 (但请记住,在高峰期您​​将有 两个 个工作,而不是一个,因为新工作在旧工作发布之前就已从队列中移除。如果这是一个问题,您总是可以这样做job = None 在队列中阻塞之前。)

所以,假设问题实际上是 job 对象(或它拥有的东西),问题是您没有向我们展示的一些代码在某处保留了对它的引用。

在不知道自己在做什么的情况下,很难提出修复建议。它可能只是“不要在那里存储”。或者它可能是“存储弱引用而不是对象本身”。或“添加 LRU 算法”。或者“添加一些流量控制,这样如果你备份太多,你就不会继续工作,直到内存用完”。


* 在 CPython 中,这会立即发生,因为垃圾收集器是基于引用计数的。另一方面,在 Jython 和 IronPython 中,垃圾收集器仅依赖于底层 VM 的垃圾收集器,因此在 JVM 或 CLR 注意到它不再被引用之前不会释放对象,这通常不是立即的,并且是不确定的.

【讨论】:

  • 我很确定过程和 post_process 不会保留任何参考。 Process 接受作业对象获取包含字符串的属性解析该字符串并将解析结果作为字符串返回(实际上返回 zlib.compress(json.dumps(result))。我在这里发现了类似的问题:stackoverflow.com/questions/21485319/…。也在这里:python.dzone.com/articles/diagnosing-memory-leaks-python 在解释中说:长时间运行的 Python 作业在运行时会占用大量内存,直到进程终止才会将内存返回给操作系统
  • @gosom:确实,CPython 几乎从不释放操作系统的内存,所以如果您的内存使用量达到峰值,那么在您退出之前,该峰值将是您的内存使用量。但是,至少在 64 位域中,这通常不是问题。如果那个额外的内存永远不会被触及,并且任何其他进程都需要它,它只会被换出并且永远不会换回,所以你实际上只是在浪费 1MB 的页表空间,而不是 12GB 的活动内存.因此,如果您认为这是正在发生的事情,请确保问题正在影响性能或稳定性,然后再浪费太多调试时间......
  • 无论如何,如果您确实正在保留垃圾,那么仅保留未使用的堆页面可能是一个更大的问题。如果你需要调试它,Python 有一些工具可以帮助你做到这一点,比如 gc 模块和(3.4+)tracemalloc,还有很多第三方 Python 模块和外部工具也可以提供帮助。除非您的工作对象随着时间的推移逐渐变大,或者multiprocessing 本身存在泄漏,否则您在某处保留了一些东西。
  • 最后一件事:您在什么平台上使用什么版本的 Python(整个 X.Y.Z,而不仅仅是 X.Y)?因为我隐约记得在 2.7.x 和 3.3.y 或类似的东西中修复的非 OS X POSIX 系统上的多处理本身有一个半严重的泄漏......所以它可能是值得的升级到 3.4 或 2.7.8 或任何合适的版本,看看它是否有所作为。
  • 我在 Linux 上使用 python 2.7.6
【解决方案2】:

如果您找不到泄漏源,您可以通过让每个工作人员只处理有限数量的任务来解决它。一旦它们达到任务限制,您就可以让它们退出,并用新的工作进程替换它们。内置的 multiprocessing.Pool 对象通过 maxtasksperchild 关键字参数支持这一点。你可以做类似的事情:

import multiprocessing
import threading

class WorkerPool(object):
    def __init__(self, workers=multiprocessing.cpu_count(),
                 maxtasksperchild=None, lock=multiprocessing.Semaphore(1)):
        self._lock = multiprocessing.Semaphore(1)
        self._max_tasks = maxtasksperchild
        self._workers = workers
        self._pool = []
        self._repopulate_pool()
        self._pool_monitor = threading.Thread(self._monitor_pool)
        self._pool_monitor.daemon = True
        self._pool_monitor.start()

    def _monitor_pool(self):
        """ This runs in its own thread and monitors the pool. """
        while True:
            self._maintain_pool()
            time.sleep(0.1)

    def _maintain_pool(self):
        """ If any workers have exited, start a new one in its place. """
        if self._join_exited_workers():
            self._repopulate_pool()

    def _join_exited_workers(self):
        """ Find exited workers and join them. """
        cleaned = False
        for i in reversed(range(len(self._pool))):
            worker = self._pool[i]
            if worker.exitcode is not None:
                # worker exited
                worker.join()
                cleaned = True
                del self._pool[i]
        return cleaned

    def _repopulate_pool(self):
        """ Start new workers if any have exited. """
        for i in range(self._workers - len(self._pool)):
            w = Worker(self._lock, self._max_tasks)
            self._pool.append(w)
            w.start()    


class Worker(multiprocessing.Process):

    def __init__(self, lock, max_tasks):
        multiprocesing.Process.__init__(self)
        self.lock = lock
        self.queue = Redis(..) # this is a redis based queue
        self.res_queue = Redis(...)
        self.max_tasks = max_tasks

     def run():
         runs = 0
         while self.max_tasks and runs < self.max_tasks:
             job = self.queue.get(block=True)
             job.results = process(job)
             with self.lock:
                 post_process(self.res_queue, job)
            if self.max_tasks:
                 runs += 1


def main():
    pool = WorkerPool(workers=4, maxtasksperchild=1000)
    # The program will block here since none of the workers are daemons.
    # It's not clear how/when you want to shut things down, but the Pool
    # can be enhanced to support that pretty easily.

请注意,上面的池监控代码与 multiprocessing.Pool 中用于相同目的的代码几乎完全相同。

【讨论】:

  • 很好的解释。我发现一次或两次有用的一种变体是使进程在达到某个峰值内存使用量时而不是在执行一定数量的任务后进行回收。 (事实上​​,这是我编写自己的池而不是仅仅使用futuresmultiprocessing...的罕见原因之一......)
  • @dano 谢谢,我尝试了类似的方法并且它有效。这是一种解决方法。当我有时间时,我会尝试弄清楚为什么会发生这种情况
  • @gosom:总是值得一探究竟......但如果它是 2.7.6 多处理中的一个错误,并且您无法升级到 2.7.8,或者它是隐含在您的您无法更改的设计等,无论如何这最终可能成为您的永久答案。
猜你喜欢
  • 1970-01-01
  • 2015-01-04
  • 2019-12-21
  • 2013-01-22
  • 1970-01-01
  • 1970-01-01
  • 2015-03-15
  • 1970-01-01
  • 2017-05-07
相关资源
最近更新 更多