【问题标题】:Python Queues memory leaks when called inside thread在线程内调用时,Python Queues 内存泄漏
【发布时间】:2013-10-18 23:44:38
【问题描述】:

我有 python TCP 客户端,需要循环发送媒体(.mpg)文件到“C”TCP 服务器。

我有以下代码,在单独的线程中,我正在读取 10K 文件块并将其发送并在循环中重新执行,我认为这是因为我实现了线程模块或 tcp send . 我正在使用 Queues 在我的 GUI (Tkinter) 上打印日志,但过了一段时间它内存不足。

更新 1 - 根据要求添加了更多代码

线程类“Sendmpgthread”用于创建线程发送数据

.
. 
def __init__ ( self, otherparams,MainGUI):
    .
    .
    self.MainGUI = MainGUI
    self.lock = threading.Lock()
    Thread.__init__(self)

#This is the one causing leak, this is called inside loop
def pushlog(self,msg):
    self.MainGUI.queuelog.put(msg)

def send(self, mysocket, block):
    size = len(block)
    pos = 0;
    while size > 0:
        try:
            curpos = mysocket.send(block[pos:])
        except socket.timeout, msg:
            if self.over:
                 self.pushlog(Exit Send)
                return False
        except socket.error, msg:
            print 'Exception'     
            return False  
        pos = pos + curpos
        size = size - curpos
    return True

def run(self):
    media_file = None
    mysocket = None 

    try:
        mysocket = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
        mysocket.connect((self.ip, string.atoi(self.port)))
        media_file = open(self.file, 'rb') 

        while not self.over:
            chunk = media_file.read(10000)
            if not chunk:   # EOF Reset it
                print 'resetting stream'
                media_file.seek(0, 0)
                continue
            if not self.send(mysocket, chunk): # If some error or thread is killed 
                break;

            #disabling this solves the issue
            self.pushlog('print how much data sent')       

    except socket.error, msg:
        print 'print exception'
    except Exception, msg:
        print 'print exception'

    try:
        if media_file is not None:
            media_file.close()
            media_file = None            
        if mysocket is not None:
            mysocket.close()
            mysocket = None
    finally:
            print 'some cleaning'   

def kill(self):
    self.over = True

我发现这是因为 Queue 的错误实现,因为评论该部分解决了问题

更新 2 - 从 Thread 类调用的 MainGUI 类

class MainGUI(Frame):
    def __init__(self, other args):
       #some code
       .
       .
        #from the above thread class used to send data
        self.send_mpg_status = Sendmpgthread(params)
        self.send_mpg_status.start()     
        self.after(100, self.updatelog)
        self.queuelog = Queue.Queue()

    def updatelog(self):
       try:
           msg = self.queuelog.get_nowait() 

           while msg is not None:
               self.printlog(msg)
               msg = self.queuelog.get_nowait() 
        except Queue.Empty:
           pass

        if self.send_mpg_status: # only continue when sending   
            self.after(100, self.updatelog)

    def printlog(self,msg):
        #print in GUI

【问题讨论】:

  • 为什么你'需要循环运行它'?需求不应该规定这样的微小细节。除非您处于非阻塞模式,否则一次发送将花费整个过程。
  • 服务器选择文件并播放它,它需要循环播放,我无法将文件发送到服务器并让它从那里播放,因为我们已经逆向工程了完整的 linux 内核和媒体从客户端控制播放器
  • 你能注释掉调用self.printlog(msg),看看问题是在实际队列中还是在打印日志中?
  • 如果你不断地调用 self.edit.insert 那么编辑控件使用的内存会一直增长。但是没必要去猜测,直接注释掉看看会发生什么。内存泄漏通常并不意味着函数“不起作用”,而是某个地方正在“积累内存”
  • 控制台将只显示最后的 X 日志消息,您可以对编辑控件执行相同操作。我会为此添加一个答案,因为 cmets 不是最好的地方。您还应该更新提及 tkinter self.edit.insert 调用的问题

标签: python multithreading python-2.7 memory-leaks queue


【解决方案1】:

由于 printlog 正在添加到 tkinter 文本控件,因此该控件占用的内存将随着每条消息而增长(它必须存储所有日志消息才能显示它们)。

除非存储所有日志至关重要,否则一个常见的解决方案是限制显示的最大日志行数。

一个简单的实现是在控件达到最大消息数后从一开始就消除多余的行。向get the number of lines in the control 添加一个函数,然后在 printlog 中添加类似于:

while getnumlines(self.edit) > self.maxloglines:
    self.edit.delete('1.0', '1.end')

(以上代码未测试)

更新:一些一般准则

请记住,看起来像内存泄漏的情况并不总是意味着函数是wrong,或者内存不再可访问。很多时候,正在积累元素的容器缺少清理代码。

此类问题的基本通用方法:

  • 就代码的哪一部分可能导致问题形成意见
  • 通过注释掉该代码来检查它(或继续注释代码直到找到候选人)
  • 在负责的代码中查找容器,添加代码以打印其大小
  • 决定哪些元素可以从容器中安全移除,以及何时移除
  • 测试结果

【讨论】:

    【解决方案2】:

    我看不出你的代码 sn-p 有什么明显错误。

    为了在 Python 2.7 下稍微减少内存使用量,我会使用 buffer(block, pos) 而不是 block[pos:]。另外我会使用mysocket.sendall(block) 而不是你的send 方法。

    如果上面的想法不能解决您的问题,那么错误很可能在您的代码中的其他地方。您能否发布完整的 Python 脚本的最短版本,它仍然会出现内存不足(http://sscce.org/)?这会增加您获得有用帮助的变化。

    【讨论】:

    • 我正在循环一个 mpg 文件,通过 tcp 读取和发送数据。我不确定是否在内存中只有 10k 块需要发送并刷新其他所有内容。
    • 请发布简化的完整 Python 源代码 (sscce.org) 以获得帮助。
    • 对我来说,您更新的代码似乎一切正常,但仍然不完整。您确定您发布的代码内存不足吗?请创建一个 sscce.org ,检查它是否内存不足并发布完整的源代码。
    • @andrewcooke 抱歉,我正在调试别人编写的代码,并且有 4000 行代码.. 无法判断要发布什么。顺便说一句,我发现这是因为队列,我用更多细节更新了代码,其中一个导致泄漏。它被循环调用,如果可能的话建议一些可以解决它的东西,我对队列很陌生。
    • @pts 我在 Update1 和 Update2 中都进行了更改。还评论了哪个部分导致泄漏
    【解决方案3】:

    内存不足错误表明数据正在生成但未使用或释放​​。浏览您的代码,我猜想这两个方面:

    • 消息正在pushlog 方法中被推送到Queue.Queue() 实例。它们被消耗了吗?
    • MainGuiprintlog 方法可能在某处写入文本。例如。它是否会持续写入某种 GUI 小部件而没有任何消息修剪?

    根据您发布的代码,我将尝试以下方法:

    1. updatelog 中添加print 语句。如果由于某种原因(例如 after() 调用失败)而没有继续调用它,那么 queuelog 将继续无限增长。
    2. 如果updatelog 不断被调用,则将注意力转移到printlog。注释这个函数的内容,看看是否仍然出现内存不足的错误。如果他们不这样做,那么 printlog 中的某些东西可能会保留记录的数据,您需要深入挖掘以找出是什么。

    除此之外,还可以稍微清理一下代码。 self.queuelog 直到线程启动后才会创建,这会产生竞争条件,线程可能会在创建队列之前尝试写入队列。 queuelog 的创建应该在线程启动之前移动到某个地方。

    updatelog 也可以重构以消除冗余:

    def updatelog(self):
           try:
               while True:
                   msg = self.queuelog.get_nowait() 
                   self.printlog(msg)
            except Queue.Empty:
               pass
    

    我假设 kill 函数是从 GUI 线程调用的。为避免线程竞争情况,self.over 应该是线程安全变量,例如 threading.Event 对象。

    def __init__(...):
        self.over = threading.Event()
    
    def kill(self):
        self.over.set()
    

    【讨论】:

      【解决方案4】:

      您的 TCP 发送循环中没有数据堆积。

      内存错误可能是由日志队列引起的,因为您还没有发布完整的代码尝试使用以下类进行日志记录:

      from threading import Thread, Event, Lock
      from time import sleep, time as now
      
      
      class LogRecord(object):
          __slots__ = ["txt", "params"]
          def __init__(self, txt, params):
              self.txt, self.params = txt, params
      
      class AsyncLog(Thread):
          DEBUGGING_EMULATE_SLOW_IO = True
      
          def __init__(self, queue_max_size=15, queue_min_size=5):
              Thread.__init__(self)
              self.queue_max_size, self.queue_min_size = queue_max_size, queue_min_size
              self._queuelock = Lock()
              self._queue = []            # protected by _queuelock
              self._discarded_count = 0   # protected by _queuelock
              self._pushed_event = Event()
              self.setDaemon(True)
              self.start()
      
          def log(self, message, **params):
              with self._queuelock:
                  self._queue.append(LogRecord(message, params))
                  if len(self._queue) > self.queue_max_size:
                      # empty the queue:
                      self._discarded_count += len(self._queue) - self.queue_min_size
                      del self._queue[self.queue_min_size:] # empty the queue instead of creating new list (= [])
                  self._pushed_event.set()
      
          def run(self):
              while 1: # no reason for exit condition here
                  logs, discarded_count = None, 0
                  with self._queuelock:
                      if len(self._queue) > 0:
                          # select buffered messages for printing, releasing lock ASAP
                          logs = self._queue[:]
                          del self._queue[:]
                          self._pushed_event.clear()
                          discarded_count = self._discarded_count
                          self._discarded_count = 0
                  if not logs:
                      self._pushed_event.wait()
                      self._pushed_event.clear()
                      continue
                  else:
                      # print logs
                      if discarded_count:
                          print ".. {0} log records missing ..".format(discarded_count)
                      for log_record in logs:
                          self.write_line(log_record)
                      if self.DEBUGGING_EMULATE_SLOW_IO:
                          sleep(0.5)
      
          def write_line(self, log_record):
              print log_record.txt, " ".join(["{0}={1}".format(name, value) for name, value in log_record.params.items()])
      
      
      
      if __name__ == "__main__":
          class MainGUI:
              def __init__(self):
                  self._async_log = AsyncLog()
                  self.log = self._async_log.log # stored as bound method
      
              def do_this_test(self):
                  print "I am about to log 100 times per sec, while text output frequency is 2Hz (twice per second)"
      
                  def log_100_records_in_one_second(itteration_index):
                      for i in xrange(100):
                          self.log("something happened", timestamp=now(), session=3.1415, itteration=itteration_index)
                          sleep(0.01)
      
                  for iter_index in range(3):
                      log_100_records_in_one_second(iter_index)
      
          test = MainGUI()
          test.do_this_test()
      

      我注意到您在发送循环中的任何地方都没有 sleep(),这意味着数据被尽可能快地读取并尽可能快地发送。请注意,在播放媒体文件时,这是不可取的行为 - 容器时间戳决定数据速率。

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 2014-10-29
        • 2013-12-18
        • 1970-01-01
        • 2017-02-04
        • 2016-12-09
        • 2020-08-08
        • 2012-02-17
        • 2021-01-10
        相关资源
        最近更新 更多