Machine learning on Spark usually starts as a notebook of DataFrame operations: fill some nulls, index a few string columns, assemble a feature vector, train a model. It works once. Then someone has to score new data every night with exactly the same transformations, retrain next month, and tune hyperparameters without accidentally letting the test data leak into the preprocessing. Spark ML Pipelines exist for that problem. A pipeline packages every step that learns from data together with the model, so that fitting, evaluation, persistence and scoring all go through one object.

This article explains the contract from first principles, walks through a complete churn model, shows how to tune a whole pipeline correctly and how to write your own stages, and then covers persistence, scoring in batch and streaming, performance and the failures you will actually hit. It uses the DataFrame-based pyspark.ml API; the older RDD-based MLlib API is in maintenance mode and is not covered. If DataFrames themselves are new to you, start with DataFrames and Datasets.

Advertisement

The contract: Transformers, Estimators and Params

Everything in pyspark.ml is one of two kinds of stage. A Transformer has a transform(df) method that takes a DataFrame and returns a new DataFrame, usually with extra columns. VectorAssembler, SQLTransformer and every trained model are transformers. An Estimator has a fit(df) method that looks at data, learns something, and returns a Transformer holding what it learned. StringIndexer is an estimator because it must see the data to learn which labels exist; its fit returns a StringIndexerModel that holds the label-to-index map. LogisticRegression is an estimator whose fit returns a model holding coefficients.

Both kinds carry Params: named, typed, documented settings such as inputCol, handleInvalid or regParam, each with an optional default. Params are what makes tuning generic, because a tuner can set any param of any stage by reference without knowing what the stage does. Every stage also has a uid, a unique identifier that persistence and param maps use to tell stages apart.

A Pipeline is itself an estimator made of a list of stages. Calling pipeline.fit(train) walks the stages in order. A transformer stage is applied to the data and the result passed on. An estimator stage is fitted on the data as transformed by all earlier stages, and the fitted model then transforms the data for the stages after it. The result is a PipelineModel: the same list of stages, with every estimator replaced by its fitted model. A PipelineModel is a transformer, so scoring is a single transform call.

Pipeline.fit(train): estimators are fitted in order on the output of the stages before themImputerEstimatorStringIndexerEstimatorOneHotEncoderEstimatorVectorAssemblerTransformerLogisticRegressionEstimatorfit()PipelineModel: every stage is now a Transformer holding learned stateImputerModelmeans / mediansStringIndexerModellabel to index mapOneHotEncoderModelcategory sizesVectorAssemblerunchangedLR modelcoefficientsmodel.save(path)one directory, all stagesBatch scoringPipelineModel.load + transformStreaming scoringsame transform on a streamEverything learned from data is learned inside fit(), from the training split only.
Simplified (the code below also has a StandardScaler): Pipeline.fit turns a list of estimators and transformers into a PipelineModel of transformers only; that one object is saved, loaded and used for batch and streaming scoring.

Why this matters: leakage and training-serving skew

Two classic bugs disappear when all learned preprocessing lives inside the pipeline. The first is leakage: if you compute the median for imputation, or the mean and variance for scaling, on the whole dataset before splitting it, the test set has influenced the features and your evaluation is optimistic. In a pipeline, those statistics are learned in fit, and fit only ever sees the training split, so test and future data are transformed with statistics they did not influence.

The second is training-serving skew. When the scoring job reimplements preprocessing by hand, it drifts: a different null rule, a new category order, a forgotten log transform. Saving the PipelineModel and loading it in the scoring job means there is only one implementation of the features, and it is the one the model was trained with.

Advertisement

A complete worked example

The model below predicts customer churn from three numeric and two categorical columns. It imputes numeric nulls with the median, indexes and one-hot encodes the categoricals, assembles and scales the features, and fits a logistic regression. Note that OneHotEncoder has been an estimator since Spark 3.0, because it learns the number of categories per column; older examples that call it as a transformer no longer apply.

from pyspark.sql import SparkSession, functions as F
from pyspark.ml import Pipeline
from pyspark.ml.feature import (Imputer, StringIndexer, OneHotEncoder,
                                VectorAssembler, StandardScaler)
from pyspark.ml.classification import LogisticRegression
from pyspark.ml.evaluation import BinaryClassificationEvaluator

spark = SparkSession.builder.appName("churn").getOrCreate()
df = spark.read.parquet("s3://lake/churn/snapshot=2026-09-01/")

train, test = df.randomSplit([0.8, 0.2], seed=42)
train = train.cache()

num_cols = ["tenure_days", "monthly_spend", "support_tickets"]
cat_cols = ["plan", "country"]

imputer = Imputer(inputCols=num_cols,
                  outputCols=[c + "_imp" for c in num_cols], strategy="median")
indexer = StringIndexer(inputCols=cat_cols,
                        outputCols=[c + "_idx" for c in cat_cols],
                        handleInvalid="keep")            # unseen label -> extra index
encoder = OneHotEncoder(inputCols=[c + "_idx" for c in cat_cols],
                        outputCols=[c + "_ohe" for c in cat_cols])
assembler = VectorAssembler(
    inputCols=[c + "_imp" for c in num_cols] + [c + "_ohe" for c in cat_cols],
    outputCol="raw_features")
scaler = StandardScaler(inputCol="raw_features", outputCol="features")
lr = LogisticRegression(featuresCol="features", labelCol="churned", maxIter=50)

pipeline = Pipeline(stages=[imputer, indexer, encoder, assembler, scaler, lr])
model = pipeline.fit(train)                  # fits every estimator, in order

scored = model.transform(test)
auc = BinaryClassificationEvaluator(labelCol="churned",
                                    metricName="areaUnderROC").evaluate(scored)
print(f"test AUC = {auc:.3f}")
model.write().overwrite().save("s3://models/churn/v12")

Three details in that code are deliberate. The split happens before any fitting, so every statistic comes from train only. The training DataFrame is cached, because the pipeline reads it once per estimator stage and logistic regression iterates over it many times; without the cache Spark would re-read and re-transform the Parquet source on each pass. And the model is saved as one directory: every stage's params and learned state go together, so there is no way to load the coefficients with the wrong indexer.

Nulls, unseen categories and other data surprises

Production data will contain values the training data did not. Each stage has an explicit policy, and the defaults are designed to fail loudly, which is correct in development and painful at 3 a.m. Decide each one on purpose.

StageSurpriseOptions
StringIndexerA label not seen in training, or nullhandleInvalid: error (default), skip (drop the row), keep (map to an extra index)
VectorAssemblerNull or NaN in an input columnhandleInvalid: error (default), skip, keep (NaN in the vector)
ImputerNulls in numeric columnsstrategy mean, median or mode; the missingValue param marks a sentinel as missing
OneHotEncoderIndex beyond the categories seen in fithandleInvalid: error or keep; dropLast controls the reference category

For scoring, skip is dangerous because it silently drops customers from the output; downstream joins then see missing predictions and nobody knows why. Prefer keep, which gives unseen categories their own bucket, and count how many rows land in it. A rising share of unseen labels is a drift signal and a reason to retrain. Also watch cardinality: StringIndexer builds its label map on the driver, and one-hot encoding a column with millions of distinct values produces enormous sparse vectors. Bucket rare categories into an other value first, or use feature hashing for very high-cardinality inputs.

Tuning the whole pipeline without leaking

CrossValidator and TrainValidationSplit are estimators that take another estimator, a grid of param maps and an evaluator. Pass them the entire pipeline, not just the final model. Then, for each fold, the imputer, indexer and scaler are fitted on that fold's training portion only, which is the only way the validation score means anything. You can put preprocessing params in the grid too, for example the imputer strategy.

from pyspark.ml.tuning import CrossValidator, ParamGridBuilder

grid = (ParamGridBuilder()
        .addGrid(lr.regParam, [0.001, 0.01, 0.1])
        .addGrid(lr.elasticNetParam, [0.0, 0.5])
        .build())                                   # 6 candidate settings

cv = CrossValidator(estimator=pipeline,             # the WHOLE pipeline
                    estimatorParamMaps=grid,
                    evaluator=BinaryClassificationEvaluator(labelCol="churned"),
                    numFolds=3,
                    parallelism=4,                  # candidate fits run concurrently
                    seed=7)
cv_model = cv.fit(train)                            # 6 x 3 = 18 pipeline fits + 1 refit
best = cv_model.bestModel                           # a PipelineModel
print(list(zip(cv_model.avgMetrics, grid)))

Budget the cost before you press run. The work is grid size times folds full pipeline fits, plus one final fit of the best settings on all of train: here 6 times 3 plus 1, so 19 fits. A grid of four params with four values each is 256 candidates, and with five folds that is 1,280 fits. The parallelism param runs several candidate fits concurrently, which helps when each fit uses only part of the cluster, but each concurrent fit holds its own cached data and models, so raise it only while executors have memory to spare. When data is large, TrainValidationSplit, a single split, costs a fraction of k-fold and is usually good enough. Keep the untouched test split for one final evaluation of cv_model.bestModel.

For time-dependent problems such as churn, random folds leak the future into the past. Spark's CrossValidator has a foldCol param (Spark 3.1 and later) that lets you assign folds yourself, for example by month, and you should evaluate on a later time window than you train on.

Writing custom stages

Two built-ins cover many custom needs without code. SQLTransformer applies a SQL statement in which __THIS__ stands for the input DataFrame, useful for ratios, filters and derived columns. RFormula builds features and a label from an R-style formula string.

from pyspark.ml.feature import SQLTransformer

ratio = SQLTransformer(statement="""
    SELECT *, monthly_spend / GREATEST(tenure_days, 1) AS spend_per_day
    FROM __THIS__
""")

When you need logic of your own, write a Transformer if the step is stateless, or an Estimator plus a Model if it learns from data. Mixing in the shared param classes gives you inputCol and outputCol with getters and setters, and DefaultParamsReadable and DefaultParamsWritable make the stage persist with the rest of the pipeline, as long as all of its state lives in params.

from pyspark import keyword_only
from pyspark.ml import Transformer
from pyspark.ml.param.shared import HasInputCol, HasOutputCol
from pyspark.ml.util import DefaultParamsReadable, DefaultParamsWritable
from pyspark.sql import functions as F

class LogTransformer(Transformer, HasInputCol, HasOutputCol,
                     DefaultParamsReadable, DefaultParamsWritable):
    """log1p of a non-negative numeric column. Stateless, so a Transformer."""

    @keyword_only
    def __init__(self, inputCol=None, outputCol=None):
        super().__init__()
        self._set(**self._input_kwargs)

    def _transform(self, dataset):
        c = self.getInputCol()
        return dataset.withColumn(self.getOutputCol(),
                                  F.log1p(F.greatest(F.col(c), F.lit(0.0))))

Prefer built-in column expressions, as above, over Python UDFs inside _transform: expressions run in the JVM and are optimised by Catalyst, while row-at-a-time Python UDFs ship every row to a Python worker and back. If you must use Python logic, use a pandas UDF; the PySpark performance guide explains the cost. One portability caveat: a pipeline containing a Python-defined stage can be loaded only from Python with that class importable, not from a Scala or Java job.

Persistence, versioning and scoring

model.write().overwrite().save(path) writes a directory with metadata JSON for the pipeline and a subdirectory per stage holding its params and learned data. PipelineModel.load(path) restores it. Treat that directory as an immutable artifact: write each run to a new versioned path with its data snapshot, code version and metrics, and promote by changing a pointer, never by overwriting. Loading models saved by older Spark versions generally works within a major line, but test it as part of every Spark upgrade.

Batch scoring is the natural fit: load, transform, write predictions with their keys. Most feature transformers and fitted models transform row by row, so the same PipelineModel can also score a Structured Streaming DataFrame with no changes; stages that need global operations are the exception, so test your specific pipeline against a stream before relying on it. For single-request, millisecond serving, Spark is the wrong runtime. Either precompute scores in batch and serve them from a key-value store, or re-implement the small final model in the serving language and verify parity on a sample. With Spark Connect, the client talks to a remote Spark server; the Spark 4.0.0 release notes list ML support over Connect, so check which stages your version supports there before moving a pipeline to it, and see the Spark Connect article for the architecture.

Performance at scale

  • Cache the right DataFrame. Cache the assembled training input for iterative learners, and unpersist it when done; Spark memory management explains what that storage competes with.
  • Right-size partitions. Iterative algorithms run one job per iteration; thousands of tiny partitions add scheduling overhead to each, while too few leave cores idle. A few partitions per core is a sensible start.
  • Keep vectors sparse. One-hot output is sparse; a StandardScaler with withMean=True has to subtract the mean from every element, which turns sparse vectors dense and can multiply memory use. Leave withMean off for sparse inputs.
  • Watch the driver. Label maps, model coefficients and tree ensembles live on the driver; very high cardinality or very large ensembles need driver memory.

Failure modes

FailureCauseFix
Scoring job fails on a new categoryStringIndexer handleInvalid left at errorUse keep and monitor the unseen-label share
Great validation score, poor production scorePreprocessing fitted outside the pipeline, or random folds on time dataPut all learned steps inside the pipeline; use time-based folds
Tuning run takes all nightGrid times folds far larger than expectedBudget fits up front, use TrainValidationSplit, shrink the grid
Executors out of memory during fitDense vectors from scaling sparse data, or high parallelismwithMean off, lower parallelism, fewer features
Model cannot be loaded in the JVM jobCustom Python stage in the pipelineRewrite the stage with built-ins or in Scala
Predictions silently missing rowshandleInvalid skip at scoring timeUse keep; reconcile input and output row counts

What to do next

  1. Take one existing notebook model and move every step that computes a statistic from data (imputation, indexing, scaling) into a Pipeline, splitting train and test before the fit.
  2. Set handleInvalid explicitly on every indexer, encoder and assembler, and add a metric for rows mapped to the unseen bucket.
  3. Wrap the pipeline in TrainValidationSplit or CrossValidator, count the fits before you start, and use time-based folds if the data is temporal.
  4. Save each trained PipelineModel to a new versioned path with its data snapshot, code version and metrics, and make the scoring job load by version.
  5. Reconcile row counts between scoring input and output on every run.
  6. Profile one fit in the Spark UI: check caching, partition count and whether any stage produced dense vectors.
Key takeaway: A Spark ML Pipeline is an estimator that fits its stages in order and returns a PipelineModel in which every learned step is a transformer. Keeping all learned preprocessing inside that object is what prevents leakage, lets you tune the whole workflow correctly, and removes training-serving skew because the saved model is the only implementation of the features. Set explicit policies for invalid data, budget tuning runs before starting them, version every saved model, and score in batch or streaming with the same artifact.