如何将MNIST的ndarray转换为Spark DataFrame以进行MLflow预测?
MNIST模型Spark DataFrame预测报错解决
问题描述
已基于MNIST数据集训练好TensorFlow模型,需将MNIST测试集的ndarray(形状为(10000,28,28))转换为Spark DataFrame用于MLflow预测,执行output.show(10)时抛出转换错误。
尝试的代码及报错
预测初始化代码
import mlflow from pyspark.sql.functions import struct, col logged_model = 'runs:/myid/myModel' loaded_model = mlflow.pyfunc.spark_udf(spark, model_uri=logged_model, result_type='double') (x_train, y_train), (x_test, y_test) = mnist.load_data()
DataFrame转换代码
mydf = spark.sparkContext.parallelize(x_test).map(lambda x: [x.tolist()]).toDF(["Doc2Vec"]) output = mydf.withColumn('predictions', loaded_model(struct(*map(col, mydf.columns))))
报错信息
An exception was thrown from a UDF: 'ValueError: Failed to convert a NumPy array to a Tensor (Unsupported object type numpy.ndarray).'.
训练模型代码
import tensorflow as tf import uuid BUFFER_SIZE = 10000 BATCH_SIZE = 64 def make_datasets(): (mnist_images, mnist_labels), _ = \ tf.keras.datasets.mnist.load_data(path=str(uuid.uuid4())+'mnist.npz') dataset = tf.data.Dataset.from_tensor_slices(( tf.cast(mnist_images[..., tf.newaxis] / 255.0, tf.float32), tf.cast(mnist_labels, tf.int64)) ) dataset = dataset.repeat().shuffle(BUFFER_SIZE).batch(BATCH_SIZE) return dataset def build_and_compile_cnn_model(): model = tf.keras.Sequential([ tf.keras.layers.Conv2D(32, 3, activation='relu', input_shape=(28, 28, 1)), tf.keras.layers.MaxPooling2D(), tf.keras.layers.Flatten(), tf.keras.layers.Dense(64, activation='relu'), tf.keras.layers.Dense(10, activation='softmax'), ]) model.compile( loss=tf.keras.losses.sparse_categorical_crossentropy, optimizer=tf.keras.optimizers.SGD(learning_rate=0.001), metrics=['accuracy'], ) return model train_datasets = make_datasets() options = tf.data.Options() options.experimental_distribute.auto_shard_policy = tf.data.experimental.AutoShardPolicy.DATA train_datasets = train_datasets.with_options(options) multi_worker_model = build_and_compile_cnn_model() multi_worker_model.fit(x=train_datasets, epochs=3, steps_per_epoch=5)
解决方案
问题根源是输入数据的预处理逻辑与训练时不一致、DataFrame格式不符合模型输入要求、MLflow UDF结果类型设置错误,具体修正步骤如下:
1. 复刻训练时的数据预处理
训练时对MNIST图片做了两个关键处理,预测时必须完全对齐:
- 添加通道维度:将(28,28)的单通道图片转为(28,28,1)
- 归一化:将像素值从0-255缩放到0-1区间
2. 修正Spark DataFrame转换逻辑
调整测试集转换方式,确保输入格式匹配模型要求:
import numpy as np # 预处理测试集:添加通道维度 + 归一化 x_test_processed = np.expand_dims(x_test / 255.0, axis=-1) # 将预处理后的数组转为Spark DataFrame,这里选择扁平化数组(适配MLflow PyFunc UDF的输入要求) mydf = spark.createDataFrame( [(x.flatten().tolist(),) for x in x_test_processed], ["image_features"] )
3. 调整MLflow UDF配置与调用
模型输出是10维的softmax概率数组,因此result_type需设为array<double>,且无需用struct包裹列:
# 加载模型时指定正确的结果类型 loaded_model = mlflow.pyfunc.spark_udf(spark, model_uri=logged_model, result_type='array<double>') # 生成预测结果,若需要类别标签可额外处理 from pyspark.sql.functions import argmax output = mydf.withColumn('prediction_probs', loaded_model(col("image_features"))) # 从概率数组中提取预测类别(最大概率对应的索引) output = output.withColumn('predicted_label', argmax(col("prediction_probs")).cast("int")) output.show(10)
4. 补充模型保存步骤(若未执行)
确保训练后的模型以MLflow兼容格式保存:
import mlflow.tensorflow with mlflow.start_run(run_id="myid"): mlflow.tensorflow.log_model(multi_worker_model, "myModel")
内容的提问来源于stack exchange,提问作者olaf
相关产品推荐
相关产品推荐

