从大文件生成DataFrame时内存不足问题的优化建议咨询
问题
需要将一个大文件转换为DataFrame,文件每行一条数据,当前做法是先读取为二维数组,遇到-99999时切换数组行,但文件过大时创建DataFrame出现内存错误:
MemoryError: Unable to allocate 11.3GiB for an Array with shape (2735, 555873)
此时系统已占用至少等量内存,却仍需额外分配11.3GiB。尝试对浮点数取整后内存消耗未改变,求优化内存占用的技巧。
当前实现代码:
temperature_matrix = [] # 所有温度数据,作为DataFrame的数据源 time_array = [] # 时间步数组,作为DataFrame的索引 elements = [] # 元素ID数组,作为DataFrame的列名 with open(self.file_location, "r") as f: timestep = -1 # 时间步计数 for line in f: line = line.strip() line_data = line.split() if line_data[0] == "-99999": # 检测到新时间步,切换行 time_array.append(float(line_data[1])) timestep = timestep + 1 temperature_matrix.append([]) continue if timestep == 0: # 第一个时间步时存储元素ID elements.append(int(line_data[0])) temperature_matrix[timestep].append(float(line_data[1])) df = pd.DataFrame(data=temperature_matrix, index=time_array, columns=elements)
输入文件示例:
-99999 0.0 # 时间步0 125 25.774447 # 时间步0时元素125的温度 126 35.774447 127 45.774447 128 55.774447 ... -99999 60.0 # 时间步60.0 125 25.774447 # 时间步60时元素125的温度 126 35.774447 127 45.774447 128 55.774447 ...
优化内存占用的技巧
1. 降低浮点数据类型精度
pandas默认用float64存储浮点数,占8字节,如果你不需要这么高的精度,换成float32(4字节)能直接把内存占用减半,甚至可以用float16(2字节,需确认精度满足需求)。创建DataFrame时直接指定类型:
df = pd.DataFrame(data=temperature_matrix, index=time_array, columns=elements, dtype='float32')
或者先把二维数组转成numpy指定类型再创建DataFrame:
import numpy as np temperature_matrix = np.array(temperature_matrix, dtype='float32') df = pd.DataFrame(data=temperature_matrix, index=time_array, columns=elements)
2. 消除中间数组的冗余存储
当前代码先把数据存在列表的列表里,转成DataFrame时会再复制一份,等于占用双倍内存。可以直接预分配numpy数组存储数据,避免动态扩容的开销:
import numpy as np # 先统计时间步和元素数量,方便预分配数组 timestep_count = 0 element_count = 0 with open(self.file_location, "r") as f: for line in f: line = line.strip() if not line: continue line_data = line.split() if line_data[0] == "-99999": timestep_count +=1 if timestep_count ==1: # 统计第一个时间步的元素总数 element_count = 0 for sub_line in f: sub_line = sub_line.strip() if not sub_line: continue sub_data = sub_line.split() if sub_data[0] == "-99999": break element_count +=1 break # 预分配固定大小的numpy数组 temperature_matrix = np.empty((timestep_count, element_count), dtype='float32') time_array = np.empty(timestep_count, dtype='float32') elements = [] with open(self.file_location, "r") as f: timestep_idx = -1 element_idx = 0 for line in f: line = line.strip() if not line: continue line_data = line.split() if line_data[0] == "-99999": timestep_idx +=1 time_array[timestep_idx] = float(line_data[1]) element_idx = 0 continue if timestep_idx ==0: elements.append(int(line_data[0])) temperature_matrix[timestep_idx, element_idx] = float(line_data[1]) element_idx +=1 df = pd.DataFrame(data=temperature_matrix, index=time_array, columns=elements)
3. 分块处理或用大数据工具
如果文件大到单内存装不下,别硬凑,用分块存储或者专门的大数据框架:
分块存到HDF5
把每个时间步的数据单独转成小DataFrame,追加到HDF5文件里,最后再读取合并:
import pandas as pd store = pd.HDFStore('temperature_data.h5') current_timestep_data = {} current_time = None with open(self.file_location, "r") as f: for line in f: line = line.strip() if not line: continue line_data = line.split() if line_data[0] == "-99999": if current_time is not None: # 存储当前时间步数据 df_chunk = pd.DataFrame([current_timestep_data], index=[current_time]) store.append('temperature', df_chunk) current_time = float(line_data[1]) current_timestep_data = {} continue element_id = int(line_data[0]) temp = float(line_data[1]) current_timestep_data[element_id] = temp # 处理最后一个时间步 if current_time is not None: df_chunk = pd.DataFrame([current_timestep_data], index=[current_time]) store.append('temperature', df_chunk) store.close() # 后续读取完整数据 df = pd.read_hdf('temperature_data.h5', 'temperature')
用Dask处理超大数据集
Dask可以把数据拆成小块并行处理,不用一次性加载到内存:
import dask.dataframe as dd from dask.diagnostics import ProgressBar # 解析每行数据的函数 def parse_line(line): line = line.strip() if not line: return None parts = line.split() if parts[0] == '-99999': return {'time': float(parts[1]), 'element_id': None, 'temperature': None} else: return {'time': None, 'element_id': int(parts[0]), 'temperature': float(parts[1])} # 读取文件并处理 ddf = dd.read_text(self.file_location).map(parse_line).dropna() # 填充每个元素对应的时间步 ddf['time'] = ddf['time'].ffill() # 转成宽表(和原DataFrame结构一致) ddf = ddf.pivot(index='time', columns='element_id', values='temperature') # 计算并转成pandas DataFrame(内存允许时),或继续用Dask操作 with ProgressBar(): df = ddf.compute()
4. 精简列名存储
如果元素ID是连续整数,可以不用把ID存成列名,用默认整数索引代替,后续需要时再映射回ID,能减少列名的内存开销(优化幅度相对小,优先级低于前面的方法)。
内容的提问来源于stack exchange,提问作者user1766050
相关产品推荐
相关产品推荐

