Python2.7 Queue多进程任务架构内存持续上涨问题问询
我来帮你分析下为啥设置了队列maxsize,主进程内存还是持续上涨,结合你的代码和内存分析结果,核心原因和解决办法如下:
1. 核心问题:Worker的结果列表一直在膨胀
先看你task_worker.py里的process_task函数:
def process_task(task, result): r = list() # 这个列表在循环外面初始化,整个worker生命周期都不会重置 while True: t = task.get() for i in xrange(10): a = list() a.append(uuid.uuid1()) for j in xrange(10): a.append(random.random()) r.append(a) result.put(r) # 每次扔给队列的都是不断变大的r!
这里的r是在while True循环外面定义的——意思是每个worker进程处理任务时,会一直往同一个列表里加元素:处理第1个任务,r有10个元素;第2个任务,r变成20个;第N个任务,r就有10*N个元素。
哪怕你给队列设了maxsize=50,队列里的每个元素体积却在持续变大,总内存自然会不断上升。你用meliae查到的那个超大列表list(4432453288 14520B 1600refs),就是这个不断膨胀的r,它的父帧是get_result,说明主进程已经拿到了这个对象,但worker还在生成更大的列表,内存根本没法稳定。
2. 次要问题:主进程垃圾回收不及时
主进程的get_result线程拿到结果列表后,只是打印就覆盖了变量r,理论上旧列表会被垃圾回收,但如果对象太大或者GC触发不及时,内存会暂时没法释放,不过这只是雪上加霜,核心还是worker的结果对象本身在变大。
修复步骤
第一步:重置Worker的结果列表(必做)
把r的初始化移到while True循环内部,让每个任务都生成独立的、固定大小的结果列表:
def process_task(task, result): while True: t = task.get() r = list() # 每个任务重新创建列表,避免持续追加 for i in xrange(10): a = list() a.append(uuid.uuid1()) for j in xrange(10): a.append(random.random()) r.append(a) print "pid: %d, result %d is waiting to put into queue" % (os.getpid(), t) result.put(r) print "pid: %d, result %d has been put into queue" % (os.getpid(), t)
这样每个任务的结果列表都是固定10个元素,队列满了之后内存就会稳定下来,不会再持续上涨。
第二步:优化主进程的内存回收(可选)
在get_result里显式清空对象并手动触发GC,帮助及时释放内存:
def get_result(result): print threading.currentThread() i = 0 while True: r = result.get() print r # 清空列表并解除引用 r.clear() r = None i = i + 1 if (i % 100 == 0): scanner.dump_all_objects('dump%s.txt' % time.time()) # 手动触发垃圾回收 import gc gc.collect()
额外提醒:确认队列阻塞逻辑
设置maxsize后,当队列满时worker的result.put()会阻塞,直到主进程取走数据。如果worker处理速度远快于主进程消费速度,队列会很快满,worker会暂停处理,这也能限制内存上涨——但前提是每个结果对象的大小固定,否则这个机制也没用。
按照上面的方法修改后,主进程的内存占用应该会稳定在一个合理范围,不会再持续攀升了。
内容的提问来源于stack exchange,提问作者蔡启申

