skip to content

How do you save and reload a fitted Spark MLlib PipelineModel, and what are the limits?

level: seniorimportance: should knowfreq 30%

answer

  1. not one file on disk
  2. a matching class must read it back
  3. guarantees stop at the major boundary
  4. scoring still needs a Spark job

basics

~10 s

Write a fitted pipeline with model.write().overwrite().save(path), which produces a directory on any Hadoop-compatible filesystem, and read it back with PipelineModel.load(path). Loading is guaranteed across minor and patch Spark versions, best-effort across majors.

solid answer

~40 s

Call `model.write().overwrite().save(path)` on the fitted `PipelineModel`; unfitted `Pipeline` objects save the same way. The output is a **directory**, not a file — pipeline metadata plus one subdirectory per stage holding that stage's parameters and learned data — written to whatever Hadoop-compatible filesystem the path names, so HDFS or object storage works. Reload with the matching class: `PipelineModel.load(path)`, not `Pipeline.load(path)`. On compatibility, MLlib is backwards compatible for minor and patch versions; across major versions there are no guarantees, only best effort, and breaking changes are reported in release notes. The persistence *format* itself is explicitly not stable — only loading is designed to be compatible. Two practical limits: a custom stage must be writable or the whole save fails, and models saved from R use a modified format loadable only in R.

code

python · 8 lines
python
from pyspark.ml import PipelineModel

model = pipeline.fit(training)
model.write().overwrite().save("s3a://models/churn/v7")

# later, in a completely separate job (Python or Scala)
loaded = PipelineModel.load("s3a://models/churn/v7")
scored = loaded.transform(today_batch)

go deeper

for a junior

Know that a fitted pipeline is written with write().save(path) to a directory and read back with PipelineModel.load(path), using the fitted class rather than the unfitted Pipeline.

for a middle

Explain what the saved directory contains — pipeline metadata plus each stage's parameters and learned data — and that it is written through Spark's filesystem layer, so HDFS or object storage is fine.

for a senior

Show you have a plan for Spark upgrades, since loading is only guaranteed across minor and patch versions, and that you know transform() is a batch path so online serving means batch scoring or an export.

for a principal

Own the artifact lifecycle: where models live, how versions are pinned to a Spark release, who is allowed to re-fit, how the training snapshot is retained, and what the rollback path is when a reload fails after an upgrade.

## The API ML persistence was added to the Pipeline API in Spark 1.6, and as of Spark 2.3 the DataFrame-based API in `spark.ml` and `pyspark.ml` has complete coverage — every standard estimator, model and pipeline can be written and read. The write side is `model.write().overwrite().save(path)`; without `overwrite()` an existing path is an error. Both fitted and unfitted objects are savable: a `Pipeline` (the recipe) and a `PipelineModel` (the recipe plus everything learned). The read side must use the class that matches what was written — `PipelineModel.load(path)` for a fitted pipeline, `Pipeline.load(path)` for an unfitted one. Loading a fitted pipeline with `Pipeline.load` is a common and confusing error. ML persistence works across Scala, Java and Python: a pipeline trained in PySpark loads in a Scala job, which is exactly what teams want when data scientists train in Python and the production scoring job is JVM code. R is the exception — it currently uses a modified format, so models saved in R can only be loaded back in R (tracked as SPARK-15572). ## What is on disk `save()` writes a **directory**, not a single serialized blob, and it writes through Spark's filesystem layer, so the path can be local, HDFS, or object storage. Inside are pipeline-level metadata (the uid, the stage ordering, the Spark version that wrote it) and a `stages` subdirectory containing one entry per stage. Each stage entry holds its own metadata — the stage class and its parameter values — and, for fitted stages, a data portion carrying the learned content: coefficients, a label vocabulary, a scaler's means and standard deviations, tree structures. That data is written in Spark's own formats, which is why the whole thing must be read back through Spark rather than parsed by hand. Two consequences follow. First, you cannot treat a saved model as a file to attach somewhere; copying it means copying a directory tree. Second, everything the pipeline learned is in there, which is precisely the point — the artifact is self-contained and cannot drift away from the feature code that produced it. ## Compatibility, stated precisely Spark's documented guarantees are worth quoting accurately, because candidates routinely overstate them. *Can a model saved by Spark X be loaded by Spark Y?* Across **major** versions: no guarantees, but best-effort. Across **minor and patch** versions: yes, these are backwards compatible. And there are no guarantees of a stable persistence *format* — the guarantee is about model loading being backwards compatible, not about the bytes staying the same. *Will it behave identically?* Across major versions, again no guarantees but best-effort. Across minor and patch versions, identical behaviour except for bug fixes. Any breaking change across a minor or patch version is reported in the release notes; if a breakage is not reported there, it should be treated as a bug. Operationally this means a Spark major upgrade needs a plan for models in flight. Load every production model against the new version in a staging job before the upgrade, keep the training code and the training data snapshot so a model can be re-fitted rather than rescued, and record the Spark version alongside every stored artifact. ## Custom stages If your pipeline contains a stage you wrote, that stage must be persistable or the whole pipeline fails to save — one unwritable stage poisons the entire artifact. In practice a custom transformer needs to implement Spark's writable/readable traits (in PySpark, `DefaultParamsWritable` and `DefaultParamsReadable`, with all its state expressed as `Param`s rather than ordinary attributes). A Python-only custom stage will also not load in a Scala job, which quietly breaks the cross-language promise for that pipeline. This is a strong argument for expressing custom logic with `SQLTransformer` or built-in stages wherever the logic allows it. ## The serving limit A loaded `PipelineModel` scores by `transform()` on a DataFrame. That is a distributed batch operation: it needs an active `SparkSession`, it plans a job, and it launches tasks. For scoring millions of rows nightly that is ideal. For answering one HTTP request in a few milliseconds it is the wrong tool by an order of magnitude — the job-launch overhead alone dwarfs the request budget. The realistic options are to batch-score into a serving store (a key-value store or an online feature store) and serve lookups, or to export the model into a non-Spark runtime. Spark's own export support is thin: `toPMML` exists only on the RDD-based `spark.mllib` side and covers a short list of models — `KMeansModel`, `LinearRegressionModel`, `RidgeRegressionModel`, `LassoModel`, `SVMModel` and binary `LogisticRegressionModel`. For unsupported models there is either no `toPMML` method or an `IllegalArgumentException`. Anything richer means a third-party exporter, and that is an architectural decision to take before training, not after. ## How to answer Give the two calls and the fact that the output is a directory. State the compatibility rule as minor/patch guaranteed, major best-effort. Then show operational judgment: name the custom-stage trap, and say plainly that `transform()` is a batch path so online serving needs batch scoring or an export.

  • Can a PipelineModel saved by Spark 3.5 be loaded by Spark 4.0?
    Possibly, but nothing guarantees it. MLlib promises backwards-compatible loading across minor and patch versions only; across major versions the commitment is best-effort. Treat a major upgrade as a migration: load every production model against the new version in staging first, keep the training code and data snapshot so anything that fails can simply be re-fitted, and record the writing Spark version with each artifact.
  • Why does saving a pipeline containing a custom transformer sometimes fail?
    Because persistence is per-stage, and every stage must be writable. A hand-written stage has to implement Spark's writable and readable traits — in PySpark, `DefaultParamsWritable` and `DefaultParamsReadable` — with its state expressed as `Param`s rather than plain attributes. One unwritable stage fails the whole save, and a Python-only custom stage will not load in a Scala job even if it does save.
  • How would you serve a Spark MLlib model for low-latency single-request predictions?
    Generally not with Spark. `transform()` plans and launches a distributed job over a DataFrame, so its overhead alone exceeds a typical request budget. Either batch-score in Spark and serve the results from a key-value store, or export the model into a non-Spark runtime. Spark's built-in `toPMML` is RDD-side only and covers a handful of models, so richer exports need a third-party tool chosen before training.

saying these in an interview costs you the question

  • Expects save() to write one serialized model file
  • Loads a fitted pipeline with Pipeline.load instead of PipelineModel.load
  • Promises a saved model loads on any future Spark major version
  • Plans per-request online serving by calling transform()
  • Saves only the estimator and reimplements the features at scoring time

context