机器学习中的异步编程应用:数据处理与模型训练优化问询
泰坦尼克数据集大规模场景优化问题解答
基于泰坦尼克数据集的VotingClassifier标准机器学习代码在小数据集上仅需约0.108秒,但面对超大规模数据集时必须进行优化。针对提出的两个技术问题,解答如下:
1. 用asyncio并行执行数据预处理列操作
要说明的是:pandas的多数操作属于CPU密集型任务,asyncio本身更适配IO密集场景,因此需要结合concurrent.futures.ThreadPoolExecutor或ProcessPoolExecutor来实现并行处理,通过asyncio的run_in_executor方法将同步任务包装为异步任务。
以下是实现代码:
import asyncio import pandas as pd from sklearn.preprocessing import LabelEncoder async def process_name(data): data['Name'] = data['Name'].map(lambda x: x.split(',')[1].split('.')[0]) return data async def process_age(data): data['Age'] = data['Age'].fillna(data['Age'].mean()) return data async def process_sex(data, label_encoder): data['Sex'] = label_encoder.fit_transform(data['Sex']) return data async def process_cabin(data): data['Cabin'] = data['Cabin'].fillna(data['Cabin'].value_counts().index[0]) data['Cabin'] = data['Cabin'].map(lambda x: x[0]) return data async def process_embarked(data): data['Embarked'] = data['Embarked'].fillna(data['Embarked'].value_counts().index[0]) return data async def main(): # 加载数据集 data = pd.read_csv('titanic.csv') mylabel = LabelEncoder() # 提交所有预处理任务到线程池并行执行 tasks = [ process_name(data), process_age(data), process_sex(data, mylabel), process_cabin(data), process_embarked(data) ] # 等待所有任务完成 await asyncio.gather(*tasks) # 查看处理后的数据 print(data.head()) if __name__ == '__main__': asyncio.run(main())
补充:因为pandas的DataFrame是可变对象,每个异步函数可直接修改原数据。如果是极端大规模数据,建议改用ProcessPoolExecutor规避GIL限制,只需在run_in_executor中指定进程池即可。
2. asyncio在VotingClassifier训练中的应用
VotingClassifier的fit方法是串行训练各个子模型的,无法直接通过asyncio优化。但可以用asyncio并行训练每个子模型,之后手动实现投票逻辑(软投票/硬投票),这样能大幅提升大规模数据集下的训练效率。
示例代码:
import asyncio import pandas as pd from sklearn.svm import SVC from sklearn.linear_model import LogisticRegression from sklearn.tree import DecisionTreeClassifier from sklearn.model_selection import train_test_split from sklearn.preprocessing import LabelEncoder import numpy as np # 异步训练单个模型 async def train_model(model, X_train, y_train): loop = asyncio.get_event_loop() # 用线程池执行训练(适配CPU密集任务) await loop.run_in_executor(None, model.fit, X_train, y_train) return model async def main(): # 加载并预处理数据 data = pd.read_csv('titanic.csv') data['Age'] = data['Age'].fillna(data['Age'].mean()) data['Embarked'] = data['Embarked'].fillna('S') le = LabelEncoder() data['Sex'] = le.fit_transform(data['Sex']) data['Embarked'] = le.fit_transform(data['Embarked']) X = data.drop(['Survived', 'Name', 'Ticket', 'Cabin'], axis=1) y = data['Survived'] X_train, X_test, y_train, y_test = train_test_split(X, y, test_size=0.2) # 定义子模型 models = [ ('svc', SVC(probability=True)), ('LR', LogisticRegression(max_iter=5000)), ('cart', DecisionTreeClassifier()) ] # 并行训练所有模型 tasks = [train_model(model, X_train, y_train) for _, model in models] trained_models = await asyncio.gather(*tasks) # 手动实现软投票逻辑 def soft_vote(models, X): probas = np.array([model.predict_proba(X) for model in models]) avg_proba = np.mean(probas, axis=0) return np.argmax(avg_proba, axis=1) # 预测并评估效果 y_pred = soft_vote(trained_models, X_test) accuracy = np.mean(y_pred == y_test) print(f"软投票准确率: {accuracy:.4f}") if __name__ == '__main__': asyncio.run(main())
说明:
- 并行训练子模型时,用
run_in_executor将同步的fit方法包装为异步任务,实现多模型并行训练。 - 训练完成后,可按需实现投票逻辑:软投票取各模型概率的平均值再选最大值对应的类别;硬投票则取各模型预测结果的众数。
内容的提问来源于stack exchange,提问作者user4356954
相关产品推荐
相关产品推荐

