What does calling fit() on a Spark MLlib Pipeline produce, and how does each stage run?
answer
- stages run in the order you listed them
- one type of stage changes on the way through
- estimators are replaced by what they produce
- fit returns a model of the whole chain
basics
~20 sPipeline.fit() returns a PipelineModel. The stages run in the declared order: Transformer stages transform the DataFrame, and each Estimator stage is fitted and replaced by the Model it produced, so the result contains only Transformers.
solid answer
~40 sA `Pipeline` is a sequence of `PipelineStage`s and is itself an `Estimator`, so `fit()` returns a `PipelineModel`. Fitting walks the stages in order against the DataFrame. For a Transformer stage, `transform()` is called and the frame passes on. For an Estimator stage, `fit()` is called on the current frame to produce a Transformer, which becomes the corresponding stage of the `PipelineModel`; if further stages follow, that new Transformer's `transform()` is applied before the frame moves down. The resulting `PipelineModel` has the same number of stages, all Transformers, and scoring is a single `transform()` call over the whole chain. Stages are wired by column names (`outputCol` of one feeding `inputCol` of the next), validated against the schema at runtime before any work starts.
code
python · 12 linesfrom pyspark.ml import Pipeline
from pyspark.ml.feature import Tokenizer, HashingTF
from pyspark.ml.classification import LogisticRegression
tokenizer = Tokenizer(inputCol="text", outputCol="words")
hashing_tf = HashingTF(inputCol="words", outputCol="features", numFeatures=1000)
lr = LogisticRegression(maxIter=10, regParam=0.001)
pipeline = Pipeline(stages=[tokenizer, hashing_tf, lr])
model = pipeline.fit(training) # -> PipelineModel
predictions = model.transform(test).select("id", "probability", "prediction")go deeper
Know that a pipeline is an ordered list of stages, that fit() trains it, and that the object you get back is what you call transform() on to make predictions.
Explain the stage-by-stage mechanics: transformers transform, estimators are fitted and replaced by their models, and the result is a PipelineModel of the same length containing only transformers.
Show why this matters in production — the fitted pipeline freezes every learned feature statistic, so scoring cannot drift from training, and schema checking fails fast before an expensive job starts.
Own the argument that the fitted pipeline is the deployable unit. Push back on designs that ship a model file plus separately maintained feature SQL, since those two artifacts inevitably diverge across releases.
## What a Pipeline is In Spark's DataFrame-based ML API, a `Pipeline` is an ordered sequence of `PipelineStage`s, where each stage is either a `Transformer` (implements `transform`) or an `Estimator` (implements `fit`). It exists because real ML work is never one call: you tokenize text, vectorise it, maybe scale it, then train. Bundling that sequence into one object is what makes it possible to tune, persist and reload the whole workflow as a unit, and it is what guarantees training data and scoring data pass through *identical* feature processing. Crucially, `Pipeline` is itself an `Estimator`. That is not a curiosity — it is the reason the whole thing can be handed to `CrossValidator`, and the reason `fit()` gives back a fitted object rather than a DataFrame. ## What fit() actually does, stage by stage `pipeline.fit(trainingDf)` runs the stages in the order they were listed, threading a DataFrame through them: - **Transformer stage**: `transform()` is called on the current DataFrame. The stage is carried into the fitted pipeline unchanged, because there is nothing to learn. - **Estimator stage**: `fit()` is called on the current DataFrame, producing a `Model` — a Transformer. *That produced Transformer* becomes the corresponding stage of the resulting `PipelineModel`. If any stages follow, the new Transformer's `transform()` is then applied to the DataFrame before it moves downstream, so later stages see the columns it appends. For the final stage this transform is unnecessary and is not performed. The classic example is `[Tokenizer, HashingTF, LogisticRegression]`. `Tokenizer.transform()` appends a words column. `HashingTF.transform()` turns words into a `features` vector column. `LogisticRegression` is an Estimator, so `fit()` runs on that frame and yields a `LogisticRegressionModel`. ## What comes out: PipelineModel `fit()` returns a `PipelineModel`, which is a Transformer. It has the same number of stages as the original pipeline, but every Estimator has become the Model it produced. At test or scoring time, `pipelineModel.transform(testDf)` sends the data through the fitted stages in order, each one updating the DataFrame and passing it along, and the predictions come out the far end. That symmetry is the whole point. The scaler applies the *training* mean, the indexer applies the *training* vocabulary, and the model applies the coefficients learned alongside them — because they all live in one fitted artifact rather than in code you might edit on one side only. ## How stages are wired Stages are not wired by position. They are wired by column names: a stage's `outputCol` (or `outputCols`) must match the `inputCol` of whatever consumes it. The array order determines execution order, but the data flow is implicit in the column parameters. The examples in the Spark guide are all *linear* pipelines, where each stage consumes the previous stage's output, but non-linear pipelines are legal as long as the data flow forms a directed acyclic graph — and in that case you must list the stages in topological order yourself, because Spark will not reorder them. ## Runtime checking, not compile-time Because pipelines operate on DataFrames with varied types, they cannot use compile-time type checking. Instead, `Pipeline` and `PipelineModel` do *runtime* checking against the DataFrame schema before actually running. This is more useful than it sounds: a typo in an `inputCol`, or an Estimator expecting a vector column that is currently a string, fails up front rather than after forty minutes of shuffling. ## Stages must be distinct instances Each Transformer and Estimator instance carries a unique id, and pipeline stages must be unique instances. Putting the same `myHashingTF` object into a pipeline twice is invalid; two separate instances, `myHashingTF1` and `myHashingTF2`, are fine even though they are the same class, because they get different ids. The ids are also what lets a `ParamMap` name a parameter on one specific stage during hyperparameter search. ## Parameters Stage parameters can be set on the instance (`lr.setMaxIter(10)`) or supplied as a `ParamMap` passed to `fit()`, where the map overrides earlier setters. Because parameters are instance-scoped, a grid can address `hashingTF.numFeatures` and `lr.regParam` in the same search without ambiguity — the mechanism `ParamGridBuilder` and `CrossValidator` are built on. ## What interviewers are checking They want to hear that `fit()` returns a *fitted pipeline*, not just a trained model, and that the feature stages inside it are frozen with their training-time state. Candidates who say "fit trains the model and then I apply my feature code again at scoring time" have described precisely the bug the abstraction exists to prevent.
- Why must a Spark MLlib Pipeline's stages be distinct instances?Every Transformer and Estimator carries a unique id, and pipeline stages are required to be unique instances, so the same object cannot appear twice. Two separate instances of the same class are fine, because they get different ids. The ids are also how a `ParamMap` targets a parameter on one specific stage during a grid search.
- How does Spark catch a misspelled inputCol before the job runs?Pipelines cannot type-check column names at compile time, so `Pipeline` and `PipelineModel` perform runtime schema checking before executing. Each stage's expected input columns and types are validated against the DataFrame schema up front, so a typo or a string column where a vector is required fails immediately rather than deep into a long job.
- Can a Spark MLlib Pipeline be non-linear?Yes. The stages are given as an ordered array, but the data flow is implicit in each stage's input and output column names, so any directed acyclic graph is expressible — for example two independent featurisation branches merged by a `VectorAssembler`. Spark will not reorder for you, so you must list the stages in topological order.
saying these in an interview costs you the question
- Thinks Pipeline.fit() returns only the trained final model
- Says pipeline stages execute in parallel
- Believes the feature stages are re-fit during scoring
- Wires stages by array position rather than column names
- Reuses one transformer instance twice in the same pipeline