迭代过程中内存占用持续升高引发内存错误的解决方法
内存泄漏问题解决方案
问题描述
每30秒获取一批JSON数据:第一批7.5万条记录处理正常,第二批9万条无问题,但第三批8万条时出现内存错误,迭代过程中内存占用持续增加。
代码示例
from multiprocessing.pool import ThreadPool from threading import Lock def divide_chunks(iterable, n): if type(iterable) is range and iterable.step != 1: # algorithm doesn't work with steps other than 1: iterable = list(iterable) l = len(iterable) n = min(l, n) k, m = divmod(l, n) return [iterable[i * k + min(i, m):(i + 1) * k + min(i + 1, m)] for i in range(n)] FILE_NO = 1 lock = Lock() def process_data(i, msgs): # arguments must be in this order global FILE_NO #process data for chunks parsed_records = [] for msg in msgs: #just deleting unnecessary keys and few key data manipulated #assigning processed msgs to record_data parsed_records.append(record_data) # Get next file number with lock: file_no = FILE_NO FILE_NO += 1 name = f"sample_{file_no}.json" with open(name, "w") as outfile: outfile.write(parsed_records) return True if __name__ == "__main__": # only imports and function/class defs before this line. # The number of chunks you want msgs split into # (this will be the number of files created for each invocation of process_data) N_CHUNKS = 10 POOL_SIZE = N_CHUNKS with ThreadPool(POOL_SIZE) as pool: while True: # Process next list of messages: msgs = [...] chunks = divide_chunks(msgs, N_CHUNKS) msgs.clear() results = pool.starmap(process_data, enumerate(chunks))
解决方案
1. 修复文件写入的致命错误
代码中outfile.write(parsed_records)是错误的,write方法仅接受字符串,直接写入列表会抛出TypeError,导致任务异常终止、资源无法正常释放。必须用json模块序列化数据:
import json # 替换原写入代码 with open(name, "w") as outfile: json.dump(parsed_records, outfile)
如果需要格式化输出,可添加indent参数:json.dump(parsed_records, outfile, indent=2)
2. 显式清理大内存对象
在process_data函数写入文件后,立即删除大列表parsed_records,帮助垃圾回收(GC)更快释放内存:
def process_data(i, msgs): # ... 其他代码 ... with open(name, "w") as outfile: json.dump(parsed_records, outfile) # 显式清理大对象 del parsed_records return True
3. 优化原数据的内存释放
原代码中msgs.clear()仅清空列表内容,但列表对象本身仍占用内存。直接删除msgs能更快释放内存:
while True: msgs = [...] chunks = divide_chunks(msgs, N_CHUNKS) del msgs # 替代msgs.clear(),直接销毁原列表 results = pool.starmap(process_data, enumerate(chunks)) # 清理临时变量 del chunks, results
4. 避免原对象引用泄漏
如果record_data是直接修改原msg对象(如删除键),会导致原msg被parsed_records引用,无法被GC回收。应创建新字典存储处理后的数据:
# 示例:仅保留需要的键 keep_keys = ["key1", "key2", "key3"] for msg in msgs: record_data = {k: v for k, v in msg.items() if k in keep_keys} parsed_records.append(record_data)
5. 降低线程池并发数
当前POOL_SIZE = N_CHUNKS会导致所有chunk同时处理,内存占用峰值过高。可将线程池大小设为CPU核心数或更小值(如5),减少同时运行的线程数:
POOL_SIZE = 5 # 替代原POOL_SIZE = N_CHUNKS
6. 手动触发垃圾回收
在循环末尾手动触发GC,强制回收未释放的内存(仅在内存增长明显时使用):
import gc while True: # ... 处理代码 ... gc.collect()
内容的提问来源于stack exchange,提问作者ditil
相关产品推荐
相关产品推荐

