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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.21 19:12:21