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

如何将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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.11 08:32:16