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
相关产品推荐
相关产品推荐

