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

已知内存泄漏来源仍无法修复生成器内存泄漏问题

批量读取Binance聚合交易数据时内存持续上涨的解决方案

问题背景

批量读取Binance BTCUSDT期货每日聚合交易的zip文件,使用生成器逐文件处理,但内存占用持续增长。尝试过del(batch)和gc.collect()都无法解决,内存增长情况如下:

Before _read_csv
Used: 10.31 GB
After _read_csv
Used: 10.31 GB
1
Before _read_csv
Used: 10.31 GB
After _read_csv
Used: 10.32 GB
2
Before _read_csv
Used: 10.32 GB
After _read_csv
Used: 10.33 GB
3
Before _read_csv
Used: 10.33 GB
After _read_csv
Used: 10.35 GB
4

最小复现代码:

from zipfile import ZipFile
import numpy as np
import pandas as pd
import psutil
import re

def print_memory_usage(msg):
    # Get memory usage information
    memory_info = psutil.virtual_memory()
    
    # Print memory usage details
    print(msg)
    print(f"Used: {memory_info.used / (1024 ** 3):.2f} GB")

def batch_generator(files: list):

    def _create_batch(file):
        return _read_file(file)

    for batch in files:
        df = _create_batch(batch)
        yield df

def _read_file(file):
    with ZipFile(file) as zipfile:
        csv_filename = re.split(r'/', file)[-1][:-4] + ".csv"
        with zipfile.open(csv_filename) as f:
            try:
                return read_aggtrades(f)
            except Exception as e:
                print(e)
                raise Exception(f"Error occurred reading file: {csv_filename}")

def read_aggtrades(file) -> pd.DataFrame:

    # EX: 578304464,17085,0.01449000,684164672,684164672,14,True,True
    columns = ['a', 'price', 'q', 'first_trade_id', 'last_trade_id', 't', 'was_the_buyer_maker']
    usecols = ['a', 'price', 'q', 't']
    dtype = {'a': np.int64, 'price': str, 'q': str, 't': np.int64}

    def peek_line(f):
        pos = f.tell()
        line = f.readline()
        f.seek(pos)

        # Convert bytes to str (line can be bytes or str)
        if type(line) == bytes:
            return line.decode()

        return line

    def _read_csv(f1):
        # 99.9% files don't have headers, some do. Discard it if we encounter it by reading header line.
        first_line = peek_line(f1)
        if first_line.startswith('agg_trade_id'):
            f1.readline()
        return pd.read_csv(f1,
                         sep=',',
                         header=None,
                         names=columns,
                         usecols=usecols,
                         dtype=dtype)
    print_memory_usage("Before _read_csv")
    df = _read_csv(file)
    print_memory_usage("After _read_csv")

    return df

file_list = [
    "/home/owner/Desktop/BTCUSDT/BTCUSDT-aggTrades-2019-12-31.zip", 
    "/home/owner/Desktop/BTCUSDT/BTCUSDT-aggTrades-2020-01-01.zip", 
    "/home/owner/Desktop/BTCUSDT/BTCUSDT-aggTrades-2020-01-02.zip", 
    "/home/owner/Desktop/BTCUSDT/BTCUSDT-aggTrades-2020-01-03.zip"]
generator = batch_generator(file_list)

i = 0
for batch in generator:
    i += 1
    print(i)

解决方案

1. 优化数据类型,减少内存占用

原代码中price和q使用字符串类型,内存占用远高于数值类型。直接改为浮点类型(精度满足需求时):

# 修改read_aggtrades里的dtype
dtype = {'a': np.int64, 'price': np.float64, 'q': np.float64, 't': np.int64}

若需要更高精度,可使用pd.DecimalDtype,但浮点类型的内存效率仍远优于字符串。

2. 修正内存统计方式,准确监控Python进程内存

原代码统计的是系统整体内存,包含系统缓存,无法反映Python进程的真实占用。改为监控当前进程的RSS内存:

def print_memory_usage(msg):
    process = psutil.Process()
    mem_rss = process.memory_info().rss / (1024 ** 3)
    print(msg)
    print(f"Python Process RSS: {mem_rss:.2f} GB")

3. 优化循环中的内存释放逻辑

确保处理完每个batch后彻底释放引用并强制垃圾回收:

import gc  # 导入gc模块

i = 0
for batch in generator:
    i += 1
    print(i)
    # 此处添加数据处理逻辑
    del batch  # 删除DataFrame引用
    gc.collect()  # 强制触发垃圾回收

4. 简化函数嵌套,避免闭包引用残留

原代码在read_aggtrades中嵌套_read_csv函数,闭包可能持有额外内存引用。去掉嵌套直接实现逻辑:

def read_aggtrades(file) -> pd.DataFrame:
    columns = ['a', 'price', 'q', 'first_trade_id', 'last_trade_id', 't', 'was_the_buyer_maker']
    usecols = ['a', 'price', 'q', 't']
    dtype = {'a': np.int64, 'price': np.float64, 'q': np.float64, 't': np.int64}

    def peek_line(f):
        pos = f.tell()
        line = f.readline()
        f.seek(pos)
        if isinstance(line, bytes):
            return line.decode()
        return line

    print_memory_usage("Before reading CSV")
    first_line = peek_line(file)
    if first_line.startswith('agg_trade_id'):
        file.readline()
    df = pd.read_csv(file,
                     sep=',',
                     header=None,
                     names=columns,
                     usecols=usecols,
                     dtype=dtype)
    print_memory_usage("After reading CSV")

    return df

5. 分块读取单个文件(可选)

若单个文件数据量极大,可通过chunksize分块读取,降低单次内存占用:

df_iter = pd.read_csv(file,
                     sep=',',
                     header=None,
                     names=columns,
                     usecols=usecols,
                     dtype=dtype,
                     chunksize=100000)  # 按10万行分块
df = pd.concat(df_iter)

若无需完整DataFrame,可直接遍历df_iter处理每个chunk,无需合并。

验证效果

应用上述修改后,运行脚本观察Python进程的内存变化,内存占用应稳定在合理范围,不会持续上涨。

内容的提问来源于stack exchange,提问作者user11629

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.27 05:44:59