mlflow.spark

mlflow.spark 模块提供用于记录和加载 Spark MLlib 模型的 API。该模块以以下 flavors 导出 Spark MLlib 模型:

Spark MLlib (native) format

允许将模型作为 Spark Transformer 在 Spark 会话中加载以进行评分。带有此 flavor 的模型可以在 Python 中作为 PySpark PipelineModel 对象加载。这是主要的 flavor,并且始终会生成。

mlflow.pyfunc

支持在 Spark 之外进行部署,方法是实例化 SparkContext 并在评分之前将输入数据读取为 Spark DataFrame。也支持作为 Spark UDF 在 Spark 中部署。具有此 flavor 的模型可以作为 Python 函数加载以执行推理。该 flavor 始终会被生成。

mlflow.spark.autolog(disable=False, silent=False)[source]

注意

自动记录与以下包版本已知兼容: 3.3.0 <= pyspark <= 4.0.0。当与该范围之外的包版本一起使用时,自动记录可能无法成功。

启用(或禁用)并配置在读取 Spark 数据源时对数据源路径、版本(如适用)和格式的日志记录。此方法不是线程安全的,并假定已存在一个SparkSession,且已附加了mlflow-spark JAR。它应在 Spark 驱动程序上调用,而不是在执行器上(即不要在由 Spark 并行化的函数内调用此方法)。所使用的 mlflow-spark JAR 必须与 Spark 的 Scala 版本匹配。有关可用版本,请参阅 Maven Repository。此 API 需要 Spark 3.0 或更高版本。

数据源信息会被缓存在内存中并记录到随后所有的 MLflow 运行中,包括活动的 MLflow 运行(如果在读取数据时存在)。注意,通过该 API 目前不支持对 Spark ML (MLlib) 模型的自动记录。数据源自动记录为尽力而为,这意味着如果 Spark 处于高负载或 MLflow 的记录因任何原因失败(例如,MLflow 服务器不可用),则可能会丢弃记录。

对于 autologging 的任何意外问题,除了检查由你的 MLflow 代码生成的 stderr & stdout 外,还应检查 Spark driver 和 executor 日志 —— 数据源信息是从 Spark 提取的,因此与调试相关的日志可能会出现在 Spark 日志中。

注意

Spark 数据源自动记录仅支持在单线程中将日志记录到 MLflow 运行

Example
import mlflow.spark
import os
import shutil
from pyspark.sql import SparkSession

# Create and persist some dummy data
# Note: the 2.12 in 'org.mlflow:mlflow-spark_2.12:2.16.2' below indicates the Scala
# version, please match this with that of Spark. The 2.16.2 indicates the mlflow version.
# Note: On environments like Databricks with pre-created SparkSessions,
# ensure the org.mlflow:mlflow-spark_2.12:2.16.2 is attached as a library to
# your cluster
spark = (
    SparkSession.builder.config(
        "spark.jars.packages",
        "org.mlflow:mlflow-spark_2.12:2.16.2",
    )
    .master("local[*]")
    .getOrCreate()
)
df = spark.createDataFrame(
    [(4, "spark i j k"), (5, "l m n"), (6, "spark hadoop spark"), (7, "apache hadoop")],
    ["id", "text"],
)
import tempfile

tempdir = tempfile.mkdtemp()
df.write.csv(os.path.join(tempdir, "my-data-path"), header=True)
# Enable Spark datasource autologging.
mlflow.spark.autolog()
loaded_df = spark.read.csv(
    os.path.join(tempdir, "my-data-path"), header=True, inferSchema=True
)
# Call toPandas() to trigger a read of the Spark datasource. Datasource info
# (path and format) is logged to the current active run, or the
# next-created MLflow run if no run is currently active
with mlflow.start_run() as active_run:
    pandas_df = loaded_df.toPandas()
Parameters
  • disable – 如果 True,则禁用 Spark 数据源的自动记录集成。 如果 False,则启用 Spark 数据源的自动记录集成。

  • silent – 如果 True,在 Spark 数据源 autologging 期间抑制来自 MLflow 的所有事件日志和警告。 如果 False,在 Spark 数据源 autologging 期间显示所有事件和警告。

mlflow.spark.get_default_conda_env(is_spark_connect_model=False)[source]
Returns

默认的 Conda 环境,用于由对 save_model()log_model() 的调用产生的 MLflow Models。该 Conda 环境包含调用者系统上安装的当前版本的 PySpark。dev 版本的 PySpark 会在生成的 Conda 环境中被替换为稳定版本(例如,如果您运行的 PySpark 版本为 2.4.5.dev0,调用此方法会生成一个对 PySpark 版本 2.4.5 有依赖的 Conda 环境)。

mlflow.spark.get_default_pip_requirements(is_spark_connect_model=False)[source]
Returns

该列表列出了由此 flavor 生成的 MLflow Models 的默认 pip 依赖项。 对 save_model()log_model() 的调用会生成一个 pip 环境,该环境至少包含这些依赖项。

mlflow.spark.load_model(model_uri, dfs_tmpdir=None, dst_path=None)[source]

从路径加载 Spark MLlib 模型。

Parameters
  • model_uri

    以 URI 格式表示的 MLflow 模型的位置,例如:

    • /Users/me/path/to/local/model

    • relative/path/to/local/model

    • s3://my_bucket/path/to/model

    • runs:/<mlflow_run_id>/run-relative/path/to/model

    • models:/<model_name>/<model_version>

    • models:/<model_name>/<stage>

    有关支持的 URI 方案的更多信息,请参见 Referencing Artifacts

  • dfs_tmpdir – 临时目录路径,位于分布式(Hadoop)文件系统(DFS)或在本地模式下的本地文件系统。模型从该位置加载。默认值为 /tmp/mlflow

  • dst_path – 要将模型工件下载到的本地文件系统路径。该目录必须已存在。如果未指定,将创建一个本地输出路径。

Returns

pyspark.ml.pipeline.PipelineModel

Example
import mlflow

model = mlflow.spark.load_model("spark-model")
# Prepare test documents, which are unlabeled (id, text) tuples.
test = spark.createDataFrame(
    [(4, "spark i j k"), (5, "l m n"), (6, "spark hadoop spark"), (7, "apache hadoop")],
    ["id", "text"],
)
# Make predictions on test documents
prediction = model.transform(test)
mlflow.spark.log_model(spark_model, artifact_path, conda_env=None, code_paths=None, dfs_tmpdir=None, registered_model_name=None, signature: mlflow.models.signature.ModelSignature = None, input_example: Union[pandas.core.frame.DataFrame, numpy.ndarray, dict, list, csr_matrix, csc_matrix, str, bytes, tuple] = None, await_registration_for=300, pip_requirements=None, extra_pip_requirements=None, metadata=None)[source]

将 Spark MLlib 模型记录为当前运行的 MLflow artifact。此操作使用 MLlib 的持久化格式,并生成一个带有 Spark flavor 的 MLflow Model。

注意:如果没有活动的 run,它将实例化一个 run 来获取 run_id。

Parameters
  • spark_model

    要保存的 Spark 模型 - MLflow 只能保存实现了 MLReadable 和 MLWritable 的 pyspark.ml.Model 或 pyspark.ml.Transformer 的子类。

    注意

    所提供的 Spark 模型的 transform 方法必须生成一个名为 “prediction” 的列,该列被用作 MLflow pyfunc 模型的输出。 大多数 Spark 模型默认会生成名为 “prediction” 的输出列,其中包含预测标签。 要将概率列设为概率性分类模型的输出列,您需要将 “probabilityCol” 参数设置为 “prediction” 并将 “predictionCol” 参数设置为 “”。 (例如 model.setProbabilityCol(“prediction”).setPredictionCol(“”))

  • artifact_path – 相对于运行的工件路径。

  • conda_env

    要么是 Conda 环境的字典表示,要么是指向 conda 环境 yaml 文件的路径。如果提供,它描述了该模型应在其中运行的环境。至少,它应指定包含在 get_default_conda_env() 中的依赖项。如果 None,会向模型添加一个包含由 mlflow.models.infer_pip_requirements() 推断的 pip 依赖的 conda 环境。如果依赖推断失败,则回退使用 get_default_pip_requirements。来自 conda_env 的 pip 依赖将被写入 pip requirements.txt 文件,完整的 conda 环境将被写入 conda.yaml。 以下是一个 示例 的 conda 环境字典表示:

    {
        "name": "mlflow-env",
        "channels": ["conda-forge"],
        "dependencies": [
            "python=3.8.15",
            {
                "pip": [
                    "pyspark==x.y.z"
                ],
            },
        ],
    }
    

  • code_paths

    本地文件系统中指向 Python 文件依赖(或包含文件依赖的目录)路径的列表。这些文件在模型加载时会被预先添加到系统路径中。如果为某个模型声明了依赖文件且多个文件之间存在导入依赖关系,那么这些文件应从一个共同的根路径声明相对导入,以避免在加载模型时发生导入错误。

    有关 code_paths 功能、推荐的使用模式和限制的详细说明,请参阅 code_paths usage guide

  • dfs_tmpdir – 在分布式(Hadoop)文件系统(DFS)上的临时目录路径;如果在本地模式下运行,则为本地文件系统上的路径。模型会先写入该位置,然后复制到模型的 artifact 目录中。由于在集群上运行时 Spark ML 模型会从 DFS 读取并写入数据,因此这是必要的。如果此操作成功完成,所有在 DFS 上创建的临时文件将被删除。默认值为 /tmp/mlflow。对于在 pyspark.ml.connect 模块中定义的模型,此参数将被忽略。

  • registered_model_name – 如果提供,将在 registered_model_name 下创建一个模型版本,并在不存在同名的注册模型时创建该注册模型。

  • signature

    A Model Signature object that describes the input and output Schema of the model. The model signature can be inferred using infer_signature function of mlflow.models.signature. Note if your Spark model contains Spark ML vector type input or output column, you should create SparkMLVector vector type for the column, infer_signature function can also infer SparkMLVector vector type correctly from Spark Dataframe input / output. When loading a Spark ML model with SparkMLVector vector type input as MLflow pyfunc model, it accepts Array[double] type input. MLflow internally converts the array into Spark ML vector and then invoke Spark model for inference. Similarly, if the model has vector type output, MLflow internally converts Spark ML vector output data into Array[double] type inference result.

    from mlflow.models import infer_signature
    from pyspark.sql.functions import col
    from pyspark.ml.classification import LogisticRegression
    from pyspark.ml.functions import array_to_vector
    import pandas as pd
    import mlflow
    
    train_df = spark.createDataFrame(
        [([3.0, 4.0], 0), ([5.0, 6.0], 1)], schema="features array<double>, label long"
    ).select(array_to_vector("features").alias("features"), col("label"))
    lor = LogisticRegression(maxIter=2)
    lor.setPredictionCol("").setProbabilityCol("prediction")
    lor_model = lor.fit(train_df)
    
    test_df = train_df.select("features")
    prediction_df = lor_model.transform(train_df).select("prediction")
    
    signature = infer_signature(test_df, prediction_df)
    
    with mlflow.start_run() as run:
        model_info = mlflow.spark.log_model(
            lor_model,
            "model",
            signature=signature,
        )
    
    # The following signature is outputted:
    # inputs:
    #   ['features': SparkML vector (required)]
    # outputs:
    #   ['prediction': SparkML vector (required)]
    print(model_info.signature)
    
    loaded = mlflow.pyfunc.load_model(model_info.model_uri)
    
    test_dataset = pd.DataFrame({"features": [[1.0, 2.0]]})
    
    # `loaded.predict` accepts `Array[double]` type input column,
    # and generates `Array[double]` type output column.
    print(loaded.predict(test_dataset))
    

  • input_example – 一个或多个有效模型输入实例。输入示例用于提示应向模型提供何种数据。它将被转换为一个 Pandas DataFrame,然后使用 Pandas 的 split-oriented 格式序列化为 json,或者转换为一个 numpy array,其中示例将通过将其转换为列表的方式序列化为 json。字节使用 base64 编码。当 signature 参数为 None 时,输入示例用于推断模型签名。

  • await_registration_for – 等待模型版本完成创建并处于 READY 状态的秒数。默认情况下,函数等待五分钟。指定 0 或 None 可跳过等待。

  • pip_requirements – 要么是一个可迭代的 pip 依赖项字符串集合(例如 ["pyspark", "-r requirements.txt", "-c constraints.txt"])或本地文件系统上 pip requirements 文件的字符串路径(例如 "requirements.txt")。如果提供,则描述该模型应运行的环境。如果 None,则由当前软件环境通过 mlflow.models.infer_pip_requirements() 推断出默认的依赖列表。如果依赖项推断失败,则改为使用 get_default_pip_requirements。依赖项和约束会被自动解析并分别写入 requirements.txtconstraints.txt 文件,并作为模型的一部分进行存储。这些依赖还会写入模型的 conda 环境(conda.yaml)文件的 pip 部分。

  • extra_pip_requirements

    要么是一个 pip 需求字符串的可迭代对象(例如 ["pandas", "-r requirements.txt", "-c constraints.txt"]),要么是本地文件系统上 pip requirements 文件的字符串路径(例如 "requirements.txt")。如果提供,该参数描述了附加的 pip 依赖,这些依赖会被追加到基于用户当前软件环境自动生成的默认 pip 依赖集合中。requirements 和 constraints 会被自动解析并分别写入 requirements.txtconstraints.txt 文件,并作为模型的一部分存储。依赖项也会被写入模型的 conda 环境(conda.yaml)文件的 pip 部分。

    警告

    以下参数不能同时指定:

    • conda_env

    • pip_requirements

    • extra_pip_requirements

    This example 演示了如何使用 pip_requirementsextra_pip_requirements 指定 pip 依赖。

  • metadata – 传递给模型并存储在 MLmodel 文件中的自定义元数据字典。

Returns

一个 ModelInfo 实例,包含已记录模型的元数据。

Example
from pyspark.ml import Pipeline
from pyspark.ml.classification import LogisticRegression
from pyspark.ml.feature import HashingTF, Tokenizer

training = spark.createDataFrame(
    [
        (0, "a b c d e spark", 1.0),
        (1, "b d", 0.0),
        (2, "spark f g h", 1.0),
        (3, "hadoop mapreduce", 0.0),
    ],
    ["id", "text", "label"],
)
tokenizer = Tokenizer(inputCol="text", outputCol="words")
hashingTF = HashingTF(inputCol=tokenizer.getOutputCol(), outputCol="features")
lr = LogisticRegression(maxIter=10, regParam=0.001)
pipeline = Pipeline(stages=[tokenizer, hashingTF, lr])
model = pipeline.fit(training)
mlflow.spark.log_model(model, "spark-model")
mlflow.spark.save_model(spark_model, path, mlflow_model=None, conda_env=None, code_paths=None, dfs_tmpdir=None, signature: mlflow.models.signature.ModelSignature = None, input_example: Union[pandas.core.frame.DataFrame, numpy.ndarray, dict, list, csr_matrix, csc_matrix, str, bytes, tuple] = None, pip_requirements=None, extra_pip_requirements=None, metadata=None)[source]

将 Spark MLlib Model 保存到本地路径。

默认情况下,此函数使用 Spark MLlib 的持久化机制保存模型。

Parameters
  • spark_model – 要保存的 Spark 模型 - MLflow 只能保存实现了 MLReadable 和 MLWritable 的 pyspark.ml.Model 或 pyspark.ml.Transformer 的子类。

  • path – 模型要保存的本地路径。

  • mlflow_model – MLflow 模型配置,此 flavor 正在被添加到该配置。

  • conda_env

    可以是 Conda 环境的字典表示或 conda 环境 yaml 文件的路径。若提供,则描述了应在其中运行该模型的环境。至少,它应当指定包含在 get_default_conda_env() 中的依赖项。若 None,则会将一个其 pip 依赖由 mlflow.models.infer_pip_requirements() 推断的 conda 环境添加到模型中。如果依赖推断失败,则回退使用 get_default_pip_requirements。来自 conda_env 的 pip 依赖将被写入到一个 pip requirements.txt 文件中,完整的 conda 环境写入到 conda.yaml。下面是一个 示例 的 conda 环境字典表示:

    {
        "name": "mlflow-env",
        "channels": ["conda-forge"],
        "dependencies": [
            "python=3.8.15",
            {
                "pip": [
                    "pyspark==x.y.z"
                ],
            },
        ],
    }
    

  • code_paths

    本地文件系统中指向 Python 文件依赖(或包含文件依赖的目录)路径的列表。这些文件在模型加载时会被预先添加到系统路径中。如果为某个模型声明了依赖文件且多个文件之间存在导入依赖关系,那么这些文件应从一个共同的根路径声明相对导入,以避免在加载模型时发生导入错误。

    有关 code_paths 功能、推荐的使用模式和限制的详细说明,请参阅 code_paths usage guide

  • dfs_tmpdir – 在分布式(Hadoop)文件系统(DFS)上的临时目录路径,或在本地模式下运行时的本地文件系统路径。模型将写入此目的地,然后复制到请求的本地路径。这是必要的,因为在集群上运行时,Spark ML 模型会从并写入 DFS。如果此操作成功完成,所有在 DFS 上创建的临时文件都将被移除。默认值为 /tmp/mlflow

  • signature – 请参阅参数 signaturemlflow.spark.log_model() 中的文档。

  • input_example – 一个或多个有效模型输入实例。输入示例用于提示应向模型提供何种数据。它将被转换为一个 Pandas DataFrame,然后使用 Pandas 的 split-oriented 格式序列化为 json,或者转换为一个 numpy array,其中示例将通过将其转换为列表的方式序列化为 json。字节使用 base64 编码。当 signature 参数为 None 时,输入示例用于推断模型签名。

  • pip_requirements – 要么是 pip 需求字符串的可迭代对象(例如 ["pyspark", "-r requirements.txt", "-c constraints.txt"])要么是本地文件系统上 pip 需求文件的字符串路径(例如 "requirements.txt")。如果提供,该项描述了应在何种环境下运行此模型。如果 None,则默认的依赖列表由 mlflow.models.infer_pip_requirements() 从当前软件环境推断得出。如果依赖推断失败,则回退使用 get_default_pip_requirements。依赖和约束会被自动解析并分别写入 requirements.txtconstraints.txt 文件,并作为模型的一部分存储。依赖也会被写入模型的 conda 环境(conda.yaml)文件的 pip 部分。

  • extra_pip_requirements

    要么是一个 pip 需求字符串的可迭代对象(例如 ["pandas", "-r requirements.txt", "-c constraints.txt"]),要么是本地文件系统上 pip requirements 文件的字符串路径(例如 "requirements.txt")。如果提供,该参数描述了附加的 pip 依赖,这些依赖会被追加到基于用户当前软件环境自动生成的默认 pip 依赖集合中。requirements 和 constraints 会被自动解析并分别写入 requirements.txtconstraints.txt 文件,并作为模型的一部分存储。依赖项也会被写入模型的 conda 环境(conda.yaml)文件的 pip 部分。

    警告

    以下参数不能同时指定:

    • conda_env

    • pip_requirements

    • extra_pip_requirements

    This example 演示了如何使用 pip_requirementsextra_pip_requirements 指定 pip 依赖。

  • metadata – 传递给模型并存储在 MLmodel 文件中的自定义元数据字典。

Example
from mlflow import spark
from pyspark.ml.pipeline import PipelineModel

# your pyspark.ml.pipeline.PipelineModel type
model = ...
mlflow.spark.save_model(model, "spark-model")