使用SparkTrials并行Hyperopt调参时MLFlow日志报错的解决方法
问题:Hyperopt SparkTrials结合MLFlow记录Keras自编码器超参数时出错
使用Hyperopt库对Keras自编码器进行超参数调优,并用MLFlow记录结果时,使用hyperopt.Trials()可以正常运行,但替换为hyperopt.SparkTrials()并行处理时,抛出属性错误,推测问题出在Spark对objective函数的并行序列化/执行逻辑上。
代码片段
AutoEncoder类
class AutoEncoder(mlflow.pyfunc.PythonModel): def __init__( self, kernel_initializer, num_layers, num_neurons, activation, regularization, reg_lambda, ): self.autoencoder = get_model( kernel_initializer=kernel_initializer, num_layers=num_layers, num_neurons=num_neurons, activation=activation, regularization=regularization, reg_lambda=reg_lambda, ) def load_context(self, context: None): pass def compile(self, optimizer, loss, metrics): if self.autoencoder is None: raise ValueError else: self.autoencoder.compile(optimizer=optimizer, loss=loss, metrics=metrics) def fit(self, X, y, epochs=None, batch_size=None, validation_data=None): history = self.autoencoder.fit( X, y, batch_size=batch_size, epochs=epochs, validation_data=validation_data, verbose=0, ) return history def predict(self, context, model_input): return self.autoencoder.predict(model_input)
目标函数(objective)
def objective(params): learning_rate = params["learning_rate"] optimizer = params["optimizer"] kernel_initializer = params["kernel_initializer"] num_layers = int(params["num_layers"]) num_neurons = int(params["num_neurons"]) activation = params["activation"] regularization = params["regularization"] reg_lambda = params["lambda"] autoencoder = AutoEncoder( kernel_initializer=kernel_initializer, num_layers=num_layers, num_neurons=num_neurons, activation=activation, regularization=regularization, reg_lambda=reg_lambda, ) autoencoder.compile(optimizer=optimizer, loss="mse", metrics=["accuracy"]) # 训练自编码器模型 history = autoencoder.fit( X_train, X_train, batch_size=32, epochs=1, validation_data=(X_val, X_val), ) with mlflow.start_run(nested=True): mlflow.log_params(params) # mlflow.log_model("autoencoder_model",python_model=autoencoder) mlflow.log_metric("val_loss", history.history["val_loss"][-1]) return history.history["val_loss"][-1]
错误信息
File "/databricks/spark/python/pyspark/serializers.py", line 189, in _read_with_length return self.loads(obj) File "/databricks/spark/python/pyspark/serializers.py", line 541, in loads return cloudpickle.loads(obj, encoding=encoding) File "/databricks/python/lib/python3.10/site-packages/mlflow/exceptions.py", line 117, in __init__ error_code = json.get("error_code", ErrorCode.Name(INTERNAL_ERROR)) AttributeError: 'str' object has no attribute 'get'
解决方法
1. 修复AutoEncoder类的load_context方法
原方法的参数类型标注错误,会导致序列化问题,修改为:
def load_context(self, context: mlflow.pyfunc.PythonModelContext): pass
2. 在Objective函数中显式设置MLFlow跟踪URI
Spark Worker节点无法自动继承Driver的MLFlow配置,需要在Objective内部显式指定:
def objective(params): # 根据你的环境设置正确的跟踪URI,Databricks环境直接用"databricks" mlflow.set_tracking_uri("databricks") # 剩余代码...
3. 优化Optimizer的传递方式
避免直接传递Optimizer对象(序列化易出问题),改为传递Optimizer名称,在Objective内部初始化并设置学习率:
def objective(params): mlflow.set_tracking_uri("databricks") learning_rate = params["learning_rate"] optimizer_name = params["optimizer"] # 根据名称初始化Optimizer if optimizer_name == "adam": optimizer = tf.keras.optimizers.Adam(learning_rate=learning_rate) elif optimizer_name == "sgd": optimizer = tf.keras.optimizers.SGD(learning_rate=learning_rate) elif optimizer_name == "rmsprop": optimizer = tf.keras.optimizers.RMSprop(learning_rate=learning_rate) else: raise ValueError(f"Unsupported optimizer: {optimizer_name}") # 剩余参数处理...
4. 确保get_model函数可被Worker访问
如果get_model是在Driver端定义的局部函数,Worker节点无法访问,需要将其定义为全局函数,或者放在可导入的模块中。
5. 简化MLFlow日志逻辑(可选)
在SparkTrials的并行环境中,嵌套Run的上下文可能出现问题,可以尝试去掉nested=True,确保每个Worker的Run上下文独立:
with mlflow.start_run(): mlflow.log_params(params) mlflow.log_metric("val_loss", history.history["val_loss"][-1])
内容的提问来源于stack exchange,提问作者WitchKingofAngmar
相关产品推荐
相关产品推荐

