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

Python脚本处理大CSV时突发内存耗尽问题求助

大型CSV转换脚本内存溢出问题及优化方案

问题描述

我编写了Python脚本处理单文件约3GB的大型CSV,采用分块读取避免内存耗尽(本机32GB内存)。但处理第一个文件时内存仅占用约3GB,处理第二个文件时内存突然飙升至25GB,触发交换分区后进程被杀死。添加sleep(60)等待垃圾回收无效,跳过第二个文件后脚本可正常运行,怀疑该文件存在数据损坏;同时发现输出文件更新不及时,数据堆积在data_dict中导致内存溢出,写入逻辑存在问题。

核心代码

import time
import pandas as pd
import os

# 假设data_dict和csv_dict已提前定义
for file in files:
    time.sleep(60)
    print(file)
    read_names = True
    count = 0
    for df in pd.read_csv(file, encoding= 'unicode_escape', chunksize=1e4, names=['all']):
        start_index = 0
        count += 1
        if read_names:
            names = df.iloc[0,:].apply(lambda x: x.split(';')).values[0]
            names = names[1:]
            start_index = 2
            read_names = False
        for row in df.iloc[start_index:,:].iterrows():
            data = row[1]
            data_list = data['all'].split(';')
            date_time = data_list[0]
            values = data_list[1:]
            date, time = date_time.split(' ')
            dd, mm, yyyy = date.split('/')
            date = yyyy + '/' + mm + '/' + dd
            for name, value in zip(names, values):
                try:
                    data_dict[name].append([name, date, time, float(value)])
                except:
                    pass
        if count % 5 == 0:
            for name in names:
                start_date = data_dict[name][0][1]
                start_time = data_dict[name][0][2]
                end_date = data_dict[name][-1][1]
                end_time = data_dict[name][-1][2]
                start_dt = start_date + ' ' + start_time
                end_dt = end_date + ' ' + end_time
                dt_index = pd.date_range(start=start_dt, freq='1S', periods=len(data_dict[name]))
                df = pd.DataFrame(data_dict[name], index=dt_index)
                df = df[3].resample('1T').mean().round(10)
                with open(csv_dict[name], 'a') as ff:
                    for index, value in zip(df.index, df.values):
                        date, time = str(index).split(' ')
                        to_write = f"{name}, {date}, {time}, {value}\n"
                        ff.write(to_write)

数据格式说明

输入格式

time sensor1 sensor2 sensor3 sensor....
2022-07-01 00:00:00; 2.559;.234;0;0;0;.....
2022-07-01 00:00:01; 2.560;.331;0;0;0;.....
2022-07-01 00:00:02; 2.558;.258;0;0;0;.....

输出格式

sensor1, 2019-05-13, 05:58:00, 2.559 
sensor1, 2019-05-13, 05:59:00, 2.560 
sensor1, 2019-05-13, 06:00:00, 2.558 

优化方案

1. 清空全局数据字典,避免跨文件内存堆积

当前data_dict未在每个文件处理前清空,导致前一个文件的数据残留,第二个文件数据持续累加触发内存溢出。处理每个文件时必须重置或清空data_dict:

for file in files:
    time.sleep(60)
    print(file)
    data_dict.clear()  # 关键:清空上一个文件的残留数据
    read_names = True
    count = 0
    # 后续处理逻辑...

如果传感器列表每个文件一致,也可以在循环外初始化空字典,每个文件处理前清空即可。

2. 写入后立即释放内存,避免数据堆积

当前每5个块写入一次,但写入后未清空data_dict中对应传感器的数据,导致内存持续占用。写入完成后必须清空对应列表:

if count % 5 == 0:
    for name in names:
        # 现有写入逻辑...
        with open(csv_dict[name], 'a') as ff:
            for index, value in zip(df.index, df.values):
                date, time = str(index).split(' ')
                to_write = f"{name}, {date}, {time}, {value}\n"
                ff.write(to_write)
        data_dict[name].clear()  # 写入后立即清空,释放内存

3. 精准处理数据异常,定位损坏文件

第二个文件可能存在格式错误(如日期分割失败、数值转换失败),当前宽泛的try-except会掩盖错误并导致异常数据堆积。需精准捕获异常并记录错误行:

# 处理日期时添加异常判断
try:
    date, time = date_time.split(' ')
    dd, mm, yyyy = date.split('/')
    date = f"{yyyy}/{mm}/{dd}"
except ValueError as e:
    print(f"文件{file}行{row[0]}:日期格式错误[{date_time}],错误信息:{e}")
    continue  # 跳过错误行

# 处理数值转换时添加异常判断
for name, value in zip(names, values):
    try:
        data_dict[name].append([name, date, time, float(value.strip())])  # 先去除值的空格
    except ValueError as e:
        print(f"文件{file}行{row[0]}:传感器[{name}]的值[{value}]转换失败,错误信息:{e}")
    except KeyError as e:
        print(f"文件{file}行{row[0]}:传感器[{name}]未在data_dict中定义,错误信息:{e}")

4. 用Pandas批量处理替代逐行迭代,提升效率并降低内存占用

当前iterrows()逐行处理效率低且内存开销大,改用Pandas原生字符串处理函数批量处理分块数据:

for df in pd.read_csv(file, encoding='unicode_escape', chunksize=10000, names=['all']):
    if read_names:
        # 直接分割表头行,无需apply
        names = df.iloc[0,0].split(';')[1:]
        read_names = False
        # 跳过表头和空行(根据实际输入格式调整)
        df = df.iloc[2:]
    
    # 批量分割每行数据为多列
    df_split = df['all'].str.split(';', expand=True)
    df_split.columns = ['datetime'] + names
    
    # 批量处理日期时间格式
    df_split[['date', 'time']] = df_split['datetime'].str.split(' ', expand=True)
    df_split[['dd', 'mm', 'yyyy']] = df_split['date'].str.split('/', expand=True)
    df_split['date'] = df_split['yyyy'] + '/' + df_split['mm'] + '/' + df_split['dd']
    
    # 批量转换数值列,自动处理无效值
    for col in names:
        df_split[col] = pd.to_numeric(df_split[col].str.strip(), errors='coerce')
    
    # 按传感器整理数据并写入
    for name in names:
        # 过滤无效值
        sensor_data = df_split[['date', 'time', name]].dropna()
        # 转换为时间序列并重采样
        sensor_data['datetime'] = pd.to_datetime(sensor_data['date'] + ' ' + sensor_data['time'])
        sensor_resampled = sensor_data.set_index('datetime')[name].resample('1T').mean().round(10)
        
        # 写入文件
        with open(csv_dict[name], 'a') as ff:
            for idx, val in sensor_resampled.items():
                date_str = idx.strftime('%Y/%m/%d')
                time_str = idx.strftime('%H:%M:%S')
                ff.write(f"{name}, {date_str}, {time_str}, {val}\n")

5. 监控内存使用,定位异常文件

添加内存监控代码,确认每个文件处理前后的内存变化,精准定位导致内存飙升的文件:

import psutil

def get_memory_usage():
    """获取当前进程内存使用量(GB)"""
    process = psutil.Process()
    return round(process.memory_info().rss / (1024 ** 3), 2)

for file in files:
    print(f"开始处理文件:{file},当前内存使用:{get_memory_usage()}GB")
    # 文件处理逻辑...
    print(f"完成处理文件:{file},当前内存使用:{get_memory_usage()}GB")

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.25 18:36:25