如何在8GB内存下加载百万级DataFrame并进行Python线性回归?
大内存DataFrame线性回归解决方案
针对8GB内存无法加载百万行DataFrame的问题,可通过以下几种方案解决:
1. 增量学习(分块训练线性回归)
线性回归的参数可通过逐块数据迭代更新,无需一次性加载全量数据。使用sklearn.linear_model.SGDRegressor(支持partial_fit方法)实现增量训练:
from pymongo import MongoClient import pandas as pd import numpy as np from sklearn.linear_model import SGDRegressor from sklearn.preprocessing import StandardScaler def get_data_chunk(chunk_size=10000): client = MongoClient(host='127.0.0.1', port=27017) database = client['database'] collection = database['AI'] # 分块读取数据,每次返回指定数量的条目 total_docs = collection.count_documents({}) for i in range(0, total_docs, chunk_size): chunk = list(collection.find({}).skip(i).limit(chunk_size)) df_chunk = pd.DataFrame(chunk) # 预处理:截取目标列、删除缺失值 df_chunk = df_chunk.iloc[:, 2:] df_chunk.dropna(inplace=True) if len(df_chunk) == 0: continue X = df_chunk.iloc[:, :-1].values y = df_chunk['price'].values.reshape(-1, 1) yield X, y # 初始化标准化器与增量回归模型 scaler_X = StandardScaler() scaler_y = StandardScaler() regr = SGDRegressor(loss='squared_error', random_state=42) first_chunk = True for X, y in get_data_chunk(chunk_size=10000): # 增量标准化数据 X_scaled = scaler_X.partial_fit(X).transform(X) y_scaled = scaler_y.partial_fit(y).transform(y).ravel() if first_chunk: # 首次训练需指定特征维度 regr.partial_fit(X_scaled, y_scaled) first_chunk = False else: regr.partial_fit(X_scaled, y_scaled) # 单独取一批数据评估模型 test_X, test_y = next(get_data_chunk(chunk_size=50000)) test_X_scaled = scaler_X.transform(test_X) test_y_scaled = scaler_y.transform(test_y).ravel() print(f"模型得分:{regr.score(test_X_scaled, test_y_scaled)}")
2. 优化数据类型减少内存占用
通过调整DataFrame字段类型,大幅降低内存消耗,尝试全量加载数据:
from pymongo import MongoClient import pandas as pd import numpy as np from sklearn.linear_model import LinearRegression from sklearn.model_selection import train_test_split def get_optimized_data(): client = MongoClient(host='127.0.0.1', port=27017) database = client['database'] collection = database['AI'] df = pd.DataFrame(list(collection.find({}))) # 数值型字段降级为更小的数据类型 for col in df.select_dtypes(include=['int64', 'float64']).columns: df[col] = pd.to_numeric(df[col], downcast='float') # 高重复率字符串字段转为category类型 for col in df.select_dtypes(include=['object']).columns: if df[col].nunique() / len(df) < 0.1: df[col] = df[col].astype('category') # 预处理 df = df.iloc[:, 2:] df.dropna(inplace=True) return df df = get_optimized_data() X = df.iloc[:, :-1].values y = df['price'].values.reshape(-1, 1) X_train, X_test, y_train, y_test = train_test_split(X, y, test_size=0.25) regr = LinearRegression() regr.fit(X_train, y_train) print(f"模型得分:{regr.score(X_test, y_test)}")
3. MongoDB端预处理减少数据量
在数据库层面先过滤缺失值、仅提取所需字段,减少传输到Python的数据量:
from pymongo import MongoClient import pandas as pd import numpy as np from sklearn.linear_model import LinearRegression from sklearn.model_selection import train_test_split def get_filtered_data(): client = MongoClient(host='127.0.0.1', port=27017) database = client['database'] collection = database['AI'] # 仅提取第2列及以后的字段(替换为实际字段名更高效) all_fields = list(collection.find_one().keys()) target_fields = all_fields[2:] projection = {field: 1 for field in target_fields} projection['_id'] = 0 # 过滤所有目标字段非空的文档 query = {field: {'$exists': True, '$ne': None} for field in target_fields} df = pd.DataFrame(list(collection.find(query, projection))) return df df = get_filtered_data() X = df.iloc[:, :-1].values y = df['price'].values.reshape(-1, 1) X_train, X_test, y_train, y_test = train_test_split(X, y, test_size=0.25) regr = LinearRegression() regr.fit(X_train, y_train) print(f"模型得分:{regr.score(X_test, y_test)}")
4. 使用Dask处理超内存数据集
Dask可将数据分块存储,在磁盘与内存间切换,支持类pandas API与scikit-learn兼容的模型:
from pymongo import MongoClient import pandas as pd import dask.dataframe as dd from dask_ml.linear_model import LinearRegression from dask_ml.model_selection import train_test_split def get_dask_data(): client = MongoClient(host='127.0.0.1', port=27017) database = client['database'] collection = database['AI'] # 将MongoDB数据转为Dask DataFrame,指定分块大小 df_pandas = pd.DataFrame(list(collection.find({}))) df = dd.from_pandas(df_pandas, chunksize=10000) # 预处理 df = df.iloc[:, 2:] df = df.dropna() return df df = get_dask_data() X = df.iloc[:, :-1] y = df['price'] X_train, X_test, y_train, y_test = train_test_split(X, y, test_size=0.25) regr = LinearRegression() regr.fit(X_train, y_train) # Dask结果需调用compute()获取实际值 print(f"模型得分:{regr.score(X_test, y_test).compute()}")
内容的提问来源于stack exchange,提问作者mahsa nemati
相关产品推荐
相关产品推荐

