You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

迭代过程中内存占用持续升高引发内存错误的解决方法

内存泄漏问题解决方案

问题描述

每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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.03 09:46:05