log signature and input data for Spark LinearRegression
Options
- Mark as New
- Bookmark
- Subscribe
- Mute
- Subscribe to RSS Feed
- Permalink
- Report Inappropriate Content
02-15-2024 08:48 AM
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?