log signature and input data for Spark LinearRegression

MohsenJ
Databricks Partner

I am looking for a way to log my `pyspark.ml.regression.LinearRegression` model with input and signature ata. The usual example that I found around are using sklearn and they can simply do 

 

# Log the model with signature and input example
signature = infer_signature(X_train, pd.DataFrame(y_train))
input_example = X_train.head(3)
mlflow.sklearn.log_model(rf_model, "rf_model", signature=signature, input_example=input_example)

 

but it doesn't work for `LinearRegression` because the "feature" column should be Spark `VectorUDT ` and it is not support by mlflow. This is how I generate my feature column

 

cat_input_cols = ["gender","occupation","zip_code","age_category"]
cat_index_output_cols = [x + '_index' for x in cat_input_cols]
ohe_output_cols = [x + '_ohe' for x in cat_input_cols]
stringIndexer = StringIndexer(inputCols=cat_input_cols, outputCols=cat_index_output_cols, handleInvalid="error", stringOrderType="alphabetDesc")
ohe_encoder = OneHotEncoder(inputCols=cat_index_output_cols, outputCols=ohe_output_cols)
 assembler = VectorAssembler(inputCols=features, outputCol="features")

pipeline = Pipeline(stages=[stringIndexer, ohe_encoder,assembler])
df_training_transformed = pipeline.fit(df_trainig_set).transform(df_trainig_set)

 

now when I build my model I like to use `infer_signature` but I'm not able to find the right way to do it.

 

with mlflow.start_run(run_name="content_based_linReg") as run:
   # ---- set hyperparameters
    lr = LinearRegression(featuresCol=COL_FEATURES, labelCol=COL_LABEL)
    lr.setMaxIter(MAX_ITER)
    lr.setRegParam(REG_PARAM)
    lr.setElasticNetParam(ELASTIC_NET_PARAM)
    lr.setFitIntercept(FIT_INTERCEPT)
   
    # ----  Split data into training and validation  sets
    (df_training_sampled, df_validation_sampled) = 
    df_content.randomSplit([0.6,0.4],SEEDS)
    #Create the model.
    lr_model = lr.fit(df_training_sampled)
    
    # Log the model parameters used for this run.
    mlflow.log_param("MAX_ITER", MAX_ITER)
    mlflow.log_param("REG_PARAM", REG_PARAM)
    mlflow.log_param("ELASTIC_NET_PARAM", ELASTIC_NET_PARAM)
    mlflow.log_param("FIT_INTERCEPT", FIT_INTERCEPT)

    # Run the model to create a prediction. Predict against the validation_df.
    df_validation_predictions = lr_model.transform(df_validation_sampled)

    # -->> this won't work
    # Log the model with signature and input example
    signature = infer_signature(df_training_sampled["features"], df_validation_predictions["prediction"])
    input_example = df_training_sampled.select(col("features"),col("rating")).head(3)

    mlflow.spark.log_model(lr_model, "LinearRegression",sample_input=input_example,signature=signature)

 

here is the error I get when I call the `infer_signature`

 

Exception: Unsupported Spark Type '<class 'pyspark.ml.linalg.VectorUDT'>', MLflow schema is only supported for scalar Spark types.

 

any idea how should I go about it?