机器学习AutoML【免费下载链接】hyperoptDistributed Asynchronous Hyperparameter Optimization in Python项目地址https://gitcode.com/gh_mirrors/hy/hyperopt点击查看免费下载Hyperopt 的SparkTrials类可以将超参数调优任务分发到 Apache Spark 集群上让成千上百次模型训练以并行任务的形式运行在各个 Spark executor 上实现大规模分布式调优。本文将基于 docs/templates/scaleout/spark.md 的官方指南结合 hyperopt/spark.py 源码与 hyperopt/tests/integration/test_spark.py 测试完整讲解SparkTrials的工作原理、三个核心配置参数以及一个可运行的 scikit-learn 调优实战示例帮助你掌握在 Spark 集群上并行执行超参数搜索的完整方案。SparkTrials 是什么超参数调优与模型选择通常需要训练成百上千个模型。SparkTrials是 Hyperopt 中专门用于 Spark 集群分布式执行的新增 Trials 实现它会将一批批训练任务并行下发每个 Spark executor 上运行一个任务从而获得巨大的横向扩展能力。该 API 最初在 Databricks 内部开发随后贡献给了 Hyperopt 开源项目。在 hyperopt/spark.py 的类注释中明确说明SparkTrials是hyperopt.Trials的分布式执行实现要求fmin运行在 Spark 集群上每个 trial一组超参数取值由单个 Spark task 处理即每个模型在一台 worker 机器上完成拟合与评估trial 之间异步运行。使用方式非常简单——只需把SparkTrials对象传给 Hyperopt 的fmin()import hyperopt best_hyperparameters hyperopt.fmin( fn training_function, space search_space, algo hyperopt.tpe.suggest, max_evals 64, trials hyperopt.SparkTrials())底层工作原理从源码结构看SparkTrials的分布式执行由 hyperopt/spark.py 中的_SparkFMinState类管理它维护一个主 dispatcher 线程外加每个 trial 一个 task 线程每个 task 线程运行一个仅含单任务的 Spark job。整体流程如下Hyperopt 主逻辑运行在 Spark driver 上fmin在 driver 端运行负责根据已完成的 trial 结果计算新的超参数组合建议。dispatcher 线程通过_poll_new_tasks()hyperopt/spark.py轮询处于JOB_STATE_NEW状态的 trial并在不超过parallelism上限的前提下逐个启动异步任务。worker 就绪时为每个超参数组合启动一个单任务 Spark job_run_trial_async()hyperopt/spark.py为每个 trial 创建独立线程线程内通过spark.sparkContext.parallelize([0], 1).mapPartitions(...).collect()将任务发送到集群执行。任务在单个 executor 上运行用户代码executor 端通过base.Domain(local_eval_function, local_space)构造评估域并调用domain.evaluate(params, ctrlNone, attach_attachmentsFalse)hyperopt/spark.py完成模型训练与评估。任务完成后将结果含 loss返回 driver执行异常时traceback 会被格式化后通过异常对象属性传回 driverhyperopt/spark.py成功时通过_finish_trial_run()把结果写回 trial 对象hyperopt/spark.py随后这些新结果会被 Hyperopt 用于计算更好的超参数建议。异步 Trials 机制SparkTrials继承自base.Trials但将类属性asynchronous置为Truehyperopt/spark.py而内存版 hyperopt/base.py 中Trials.asynchronous False。正是这一标志让fmin知道当产生新搜索点时不要原地执行目标函数而是把任务交给集群异步执行、由其他进程/任务完成评估后再回填结果。fmin 的分发入口在 hyperopt/fmin.py 中可以看到当传入的trials对象带有fmin方法且allow_trials_fminTrue时fmin会直接将控制权转交给SparkTrials.fmin()hyperopt/spark.py。SparkTrials.fmin内部先启动 dispatcher 线程再调用通用fmin逻辑设置allow_trials_fminFalse防止递归并在finally块中等待所有线程结束后打印汇总日志例如Total Trials: 8: 6 succeeded, 2 failed, 0 cancelled.该日志格式被 hyperopt/tests/integration/test_spark.py 的check_run_status断言验证。适用边界单机模型还是分布式模型由于SparkTrials在单个 Spark worker 上完成每个模型的拟合与评估它仅适用于单机 ML 模型与工作流例如 scikit-learn 或单机版 TensorFlow。官方文档明确建议适合SparkTrialsscikit-learn、单机 TensorFlow 等单机模型。不适合SparkTrialsApache Spark MLlib、Horovod 等分布式 ML 算法——这类算法应使用 Hyperopt 默认的Trials类在 driver 上运行主逻辑即可否则会把本应分布式训练的算法硬塞进单个 executor。SparkTrials 配置参数详解SparkTrials可通过 3 个可选参数进行配置。下面结合源码逐项说明hyperopt/spark.py 的构造函数签名parallelism并发调优的上限parallelism指定最多同时评估多少个 trial即最大并发 Spark 任务数。默认值为 SparkSparkContext.defaultParallelism若未显式指定或传入非正值则取max(spark_default_parallelism, 1)见 hyperopt/spark.py 的_decide_parallelism。与max_evals的权衡关系若parallelism max_evalsHyperopt 会一次性独立选好全部超参数组合并并行评估等效于 Random Search随机搜索。若parallelism 1则完全串行可以充分利用 TPETree of Parzen Estimators等自适应算法的迭代探索能力——每个新组合都基于此前结果选择。将parallelism设置在 1 与max_evals之间即可在扩展性更快拿到结果与自适应性有时能得到更好模型之间做取舍。硬性上限当前并发任务数存在 128 的硬上限对应源码中的MAX_CONCURRENT_JOBS_ALLOWED 128hyperopt/spark.py。SparkTrials还会检查集群配置允许的并发任务数若parallelism超出该上限会自动降级为集群允许的最大值。测试 hyperopt/tests/integration/test_spark.py 专门验证了“None/负值回退到默认值”以及“超过 128 被截断”两种行为。timeout调优时长预算timeout是允许fmin()运行的最大秒数默认None不限时。它提供预算控制机制一旦超时fmin会尽可能终止正在运行的 trial然后退出并返回当前已完成的全部结果。在 dispatcher 循环中hyperopt/spark.py当elapsed_time self.trials.timeout时会置_fmin_cancelled标志并调用_cancel_running_trials()。取消行为取决于 PySpark 版本能力hyperopt/spark.py若支持 job group 取消collectWithJobGroup或 pinned-threads 模式会调用sparkContext.cancelJobGroup(...)取消该 job group 下所有运行中的 Spark job并将NEW/RUNNING状态的 trial 标记为CANCEL。若不支持较老版本则只能停止启动新任务运行中的 Spark job 会阻塞直到自然结束并打印警告日志。测试 hyperopt/tests/integration/test_spark.py 分别覆盖了“可取消”与“不可取消”两条路径。spark_session指定 SparkSessionspark_session是供SparkTrials使用的SparkSession实例。若不传SparkTrials会调用SparkSession.builder.getOrCreate()复用现有会话或创建新会话hyperopt/spark.py。注意构造SparkTrials时若环境中无法导入 PySpark_have_spark False会直接抛出异常提示先确保 PySpark 可用hyperopt/spark.py。源码中可用的附加参数除文档列出的 3 个参数外源码构造函数还支持两个进阶参数loss_threshold当最优 loss 降到该阈值以下时提前终止搜索需为正浮点数通过validate_loss_threshold校验。resource_profileResourceProfile对象用于在 Spark 训练任务中指定资源如 GPU。仅当 Spark 版本支持 stage-level scheduling即parallelize([1]).withResources可用且spark.master非 local 模式时生效否则打印警告并忽略。SparkTrials的完整 API 说明也可以通过help(SparkTrials)在交互环境中查看。实战用 SparkTrials 调优 scikit-learn 稀疏逻辑回归以下完整示例改编自 scikit-learn 官方文档中基于 MNIST 的稀疏逻辑回归示例展示如何用SparkTrials调优一个 scikit-learn 模型。原示例来自 scikit-learn 文档的 sparse logistic regression with MNIST 页面此处保留了完整的可运行代码from sklearn.datasets import fetch_openml from sklearn.linear_model import LogisticRegression from sklearn.model_selection import train_test_split from sklearn.preprocessing import StandardScaler from sklearn.utils import check_random_state from hyperopt import fmin, hp, tpe from hyperopt import SparkTrials, STATUS_OK # Load MNIST data, and preprocess it by standarizing features. X, y fetch_openml(mnist_784, version1, return_X_yTrue) random_state check_random_state(0) permutation random_state.permutation(X.shape[0]) X X[permutation] y y[permutation] X X.reshape((X.shape[0], -1)) X_train, X_test, y_train, y_test train_test_split( X, y, train_size5000, test_size10000) scaler StandardScaler() X_train scaler.fit_transform(X_train) X_test scaler.transform(X_test) # First, set up the scikit-learn workflow, wrapped within a function. def train(params): This is our main training function which we pass to Hyperopt. It takes in hyperparameter settings, fits a model based on those settings, evaluates the model, and returns the loss. :param params: map specifying the hyperparameter settings to test :return: loss for the fitted model # We will tune 2 hyperparameters: # regularization and the penalty type (L1 vs L2). regParam float(params[regParam]) penalty params[penalty] # Turn up tolerance for faster convergence clf LogisticRegression(C1.0 / regParam, multi_classmultinomial, penaltypenalty, solversaga, tol0.1) clf.fit(X_train, y_train) score clf.score(X_test, y_test) return {loss: -score, status: STATUS_OK} # Next, define a search space for Hyperopt. search_space { penalty: hp.choice(penalty, [l1, l2]), regParam: hp.loguniform(regParam, -10.0, 0), } # Select a search algorithm for Hyperopt to use. algotpe.suggest # Tree of Parzen Estimators, a Bayesian method # We can run Hyperopt locally (only on the driver machine) # by calling fmin without an explicit trials argument. best_hyperparameters fmin( fntrain, spacesearch_space, algoalgo, max_evals32) best_hyperparameters # We can distribute tuning across our Spark cluster # by calling fmin with a SparkTrials instance. spark_trials SparkTrials() best_hyperparameters fmin( fntrain, spacesearch_space, algoalgo, trialsspark_trials, max_evals32) best_hyperparameters示例要点解析目标函数契约train接收一组超参regParam、penalty返回包含loss与status的字典。返回{loss: -score, status: STATUS_OK}是因为 Hyperopt 默认最小化 loss取负准确率即可实现“准确率越高越好”的优化目标。搜索空间hp.choice(penalty, [l1, l2])在 L1/L2 惩罚间选择hp.loguniform(regParam, -10.0, 0)在[e^-10, e^0]区间对数均匀采样正则化强度C 1.0 / regParam完成到 scikit-learn 参数空间的换算。本地与集群两种运行方式第一个fmin不传trials仅在 driver 机器上串行执行 32 次评估第二个fmin传入SparkTrials()将 32 次评估分布到 Spark 集群并行执行。两者的目标函数、搜索空间与搜索算法完全一致可以直观对比串行与分布式的效果。执行异常的处理executor 上的训练异常不会导致整个fmin崩溃而是被捕获并格式化为 traceback 写回 trial 的misc.error字段trial 状态记为JOB_STATE_ERROR。测试 hyperopt/tests/integration/test_spark.py 验证了失败 trial 会保留RuntimeError与完整 Traceback 信息且count_successful_trials()/count_failed_trials()统计准确hyperopt/tests/integration/test_spark.py 同时确认失败任务不会被 Spark 自动重试每个 trial 恰好执行一次。使用注意事项结合文档与源码使用SparkTrials时有几点需要留意环境要求必须运行在 Spark 集群上且环境中可import pyspark否则构造SparkTrials会抛异常。不支持的 fmin 特性SparkTrials.fmin明确断言不支持pass_expr_memo_ctrl与catch_eval_exceptions也不支持points_to_evaluate、trials_save_filehyperopt/spark.pytrial_attachments同样不支持hyperopt/spark.py。超时取消的能力依赖超时终止运行中任务的能力取决于 PySpark 是否支持按 job group 取消若不支持超时只阻止新任务启动运行中的任务会继续跑完。提前停止early_stop_fn在 dispatcher 线程中被检查hyperopt/spark.py一旦触发会取消运行中的 trialfmin_cancelled_reason会记录为early stopping condition测试 hyperopt/tests/integration/test_spark.py 验证了该行为。结果的持久性与内存版Trials不同SparkTrials的结果保存在 driver 进程内的 trial 对象中若要跨进程/跨实验共享结果可参考 docs/templates/scaleout/mongodb.md 中基于 MongoDB 的MongoTrials持久化方案。小结SparkTrials让 Hyperopt 的超参数搜索能够利用 Spark 集群的横向扩展能力driver 负责搜索逻辑每个 executor 并行训练并评估一个模型通过parallelism在“随机搜索式的大规模并行”与“TPE 式的自适应串行”之间灵活调节配合timeout实现调优时长预算控制。对于 scikit-learn、单机 TensorFlow 等单机模型的超参数调优它是将调优任务规模化到 Spark 集群的首选方案而对于 MLlib、Horovod 等本身即分布式的算法则应回到默认Trials类。相关源码与测试可继续在 hyperopt/spark.py 与 hyperopt/tests/integration/test_spark.py 中深入研读。赞分享机器学习AutoML【免费下载链接】hyperoptDistributed Asynchronous Hyperparameter Optimization in Python项目地址https://gitcode.com/gh_mirrors/hy/hyperopt点击查看免费下载相关推荐XGBoost4J-Spark-GPU 实践指南在 Apache Spark 集群上端到端 GPU 加速分布式 XGBoost 训练XGBoost4J Spark GPU 实践指南在 Apache Spark 集群上端到端 GPU 加速分布式 XGBoost 训练 本文以官方教程 doc/人工智能机器学习Horovod Ray Tune 分布式超参数搜索实战指南从单机调参到多机集群Horovod Ray Tune 分布式超参数搜索实战指南从单机调参到多机集群 导读 本指南基于 Horovod 官方文档《Distributed Hy深度学习机器学习分布式训练Ray Tune 可扩展超参数调优实战指南从核心概念到分布式搜索Ray Tune 可扩展超参数调优实战指南从核心概念到分布式搜索 Ray TuneTune是 Ray 生态中专用于 可扩展超参数搜索 的 Python 库人工智能分布式训练强化学习任务调度模型推理服务后端上一篇分布式调度平台多集群协同5大跨区域任务负载均衡策略终极指南下一篇Relax CMS多数据库支持MongoDB之外的存储选择创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考