多数据集多模型训练:多进程方案是否更高效?
Great question—this is a super common scenario when scaling ML experiments, and the answer depends on a few key factors. Let’s break this down step by step.
Core Decision: Parallel vs. Sequential Training
- Go parallel if: You have multiple datasets/models, and individual training runs don’t max out your CPU. Parallelizing independent tasks will cut down total runtime significantly.
- Stick to sequential if: Your single-model training already uses all available CPU cores (e.g.,
RandomForest(n_jobs=-1)). Adding manual multiprocessing here causes CPU oversubscription—your system will waste time switching between competing processes, slowing everything down.
Handling Built-in Multiprocessing in scikit-learn/XGBoost
Models like RandomForest and XGBoost use n_jobs (or nthread for older XGBoost versions) to parallelize training internally. This creates a conflict if you layer on multiprocessing.Pool:
- For example, an 8-core CPU running 4 parallel processes, each with
n_jobs=8, tries to use 32 cores total—way more than your hardware supports. - Fix this by:
- Setting
n_jobs=1for all models andRandomizedSearchCVwhen using manual multiprocessing. This lets each parallel task use exactly one core. - Or, balance resources: If you have 8 cores, run 2 parallel processes each with
n_jobs=4(total 8 cores used). Test this to find the sweet spot for your hardware.
- Setting
Feedback & Optimizations for Your Code
Your implementation is a solid start, but here are some tweaks to make it cleaner and more efficient:
- Simplify task queuing: Instead of looping through model functions and calling
map_asynceach time, package all (model, dataset) pairs into a single task list. This avoids redundant Pool setup/teardown. - Avoid nested list flattening: Your final
flattenedlogic can be simplified withitertools.productandstarmap(which handles multiple arguments for parallel tasks). - Explicitly set
n_jobs=1: As mentioned earlier, this prevents nested parallelism and resource fights.
Optimized Code Example
Here’s a revised version that fixes these issues and avoids resource conflicts:
import numpy as np import pandas as pd import xgboost as xgb from xgboost.sklearn import XGBClassifier from sklearn.ensemble import RandomForestClassifier from sklearn.model_selection import RandomizedSearchCV import multiprocessing as mp from sklearn.datasets import load_iris import itertools def run_task(model_fn, dataset): """Wrapper to execute one model-dataset training pair""" return model_fn(dataset) def apply_parallel_training(model_fns, datasets): """Parallelize training across all model-dataset combinations""" total_cores = mp.cpu_count() pool_size = total_cores - 1 # Leave 1 core for system processes with mp.Pool(pool_size) as pool: # Create all possible (model, dataset) pairs tasks = itertools.product(model_fns, datasets) # Use starmap to pass multiple arguments to run_task results = pool.starmap(run_task, tasks) return results def random_forest(x): """Random Forest with n_jobs=1 to avoid parallelism conflicts""" return train_some_models( x=x, clf=RandomForestClassifier(n_jobs=1), param_dist={ "n_estimators": [500], "max_depth": [5, 7], "max_features": ["auto", "log2"], "bootstrap": [True, False], "criterion": ["gini", "entropy"] }, n_iter=5, model_name="random_forest" ) def gradient_boosting_tree(x): """XGBoost with n_jobs=1 to avoid parallelism conflicts""" return train_some_models( x=x, clf=XGBClassifier(n_jobs=1), param_dist={ "n_estimators": [200], "learning_rate": [0.005, 0.01, 0.05, 0.1], "booster": ["gbtree"], "max_depth": [5, 7] }, n_iter=5, model_name="gradient_boosting_tree" ) def train_some_models(x, clf, param_dist, n_iter, model_name): """Train model and return prediction results as DataFrame""" Y = x[["target", "train_test_label"]].copy() X = x.drop(columns=["target"]) # Split train/test and clean up labels X_train = X[X.train_test_label == "train"].drop(columns=["train_test_label"]) X_test = X[X.train_test_label == "test"].drop(columns=["train_test_label"]) y_train = Y[Y.train_test_label == "train"]["target"].values.ravel() y_test = Y[Y.train_test_label == "test"]["target"].values.ravel() # Set n_jobs=1 here too to avoid nested parallelism in RandomizedSearchCV random_clf = RandomizedSearchCV( clf, param_distributions=param_dist, n_iter=n_iter, scoring='f1', n_jobs=1 ) random_clf.fit(X_train, y_train) y_class = random_clf.predict(X_test) return pd.DataFrame({ "y_class": y_class, "model_name": model_name, "y_test": y_test }) if __name__ == '__main__': # Prepare binary classification Iris dataset iris = load_iris() iris = pd.DataFrame( data=np.c_[iris['data'], iris['target']], columns=iris['feature_names'] + ['target'] ) iris["target"] = np.where(iris["target"] == 2, 1, 0) # Create Dataset 1 iris_train_1 = iris.loc[0:70,:].assign(train_test_label="train") iris_test_1 = iris.loc[70:100,:].assign(train_test_label="test") iris_1 = pd.concat([iris_train_1, iris_test_1], axis=0) # Create Dataset 2 iris_train_2 = iris.loc[100:130,:].assign(train_test_label="train") iris_test_2 = iris.loc[130:150,:].assign(train_test_label="test") iris_2 = pd.concat([iris_train_2, iris_test_2], axis=0) datasets = [iris_1, iris_2] model_fns = [random_forest, gradient_boosting_tree] # Run parallel training results = apply_parallel_training(model_fns, datasets) print(pd.concat(results))
Final Takeaways
- Parallelism is worth it when you have independent training tasks—just make sure to control CPU usage per task.
- Never mix nested parallelism: Always set
n_jobs=1for models and hyperparameter search when using manual multiprocessing. - Test small first: Run a quick comparison of sequential vs. parallel training on a subset of your data to confirm which approach is faster for your specific setup.
内容的提问来源于stack exchange,提问作者Aerocell
相关产品推荐
相关产品推荐

