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

多数据集多模型训练:多进程方案是否更高效?

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=1 for all models and RandomizedSearchCV when 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.
Feedback & Optimizations for Your Code

Your implementation is a solid start, but here are some tweaks to make it cleaner and more efficient:

  1. Simplify task queuing: Instead of looping through model functions and calling map_async each time, package all (model, dataset) pairs into a single task list. This avoids redundant Pool setup/teardown.
  2. Avoid nested list flattening: Your final flattened logic can be simplified with itertools.product and starmap (which handles multiple arguments for parallel tasks).
  3. 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=1 for 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 04:58:21