Why should a Spark MLlib CrossValidator wrap the whole Pipeline instead of just the final estimator?
answer
- what the folds are supposed to be hiding
- some feature stages learn from data too
- statistics computed before the split
- fold count multiplies the grid
basics
~20 sA Spark Pipeline is itself an Estimator, so CrossValidator can re-fit every feature stage inside each fold's training split. Fitting scalers or indexers once over the full dataset leaks held-out statistics into training and inflates the score.
solid answer
~40 s`CrossValidator` takes an Estimator, a set of `ParamMap`s and an `Evaluator`, and because a `Pipeline` *is* an Estimator, the whole workflow can go inside. That matters for two reasons. First, correctness: stages like `StandardScaler`, `StringIndexer`, `IDF` and `Imputer` learn statistics from whatever data they are fitted on, so fitting them once over the full dataset and cross-validating only the classifier lets each fold's held-out rows influence its own training — the metric comes out optimistic. Wrapping the pipeline re-fits every stage inside each fold's training split. Second, reach: the grid can then tune feature parameters such as `hashingTF.numFeatures` jointly with `lr.regParam`. The cost is real — folds times grid size fits — so manage it with `parallelism`, caching, a smaller grid, or `TrainValidationSplit`.
code
python · 10 linesfrom pyspark.ml.feature import StandardScaler
from pyspark.ml.classification import LogisticRegression
from pyspark.ml.tuning import CrossValidator
# LEAKY: the scaler sees every row, including each fold's held-out third
scaler_model = StandardScaler(inputCol="features", outputCol="scaled").fit(data)
scaled = scaler_model.transform(data)
cv = CrossValidator(estimator=LogisticRegression(featuresCol="scaled"),
estimatorParamMaps=grid, evaluator=ev, numFolds=3)
cv.fit(scaled)go deeper
Know that cross-validation splits data into k folds, trains on the rest and evaluates on the held-out fold, and that Spark's CrossValidator needs an estimator, a parameter grid and an evaluator.
Explain that a Spark Pipeline is itself an Estimator, so the whole thing can be handed to CrossValidator, and be able to compute how many fits a given grid and fold count implies.
Demonstrate that you would spot leakage from a scaler or indexer fitted outside the folds, and that you would manage the cost deliberately with parallelism, caching, a coarser grid or TrainValidationSplit.
Own the search budget. Decide what a wider grid is actually worth in cluster hours, when a coarse search on a sample beats an exhaustive one on the full data, and what metric the organisation is selecting on.
## What CrossValidator needs MLlib's model-selection tools, `CrossValidator` and `TrainValidationSplit`, take three things: an **Estimator** (an algorithm *or a Pipeline*) to tune, a set of **ParamMaps** — the parameter grid, usually built with `ParamGridBuilder` — and an **Evaluator** that scores a fitted model on held-out data (`RegressionEvaluator`, `BinaryClassificationEvaluator`, `MulticlassClassificationEvaluator`, `MultilabelClassificationEvaluator` or `RankingEvaluator`, each with a `setMetricName` to choose the metric). `CrossValidator` splits the data into k folds. With k=3 it generates three (training, test) pairs, each training on two thirds and testing on the remaining third. For every `ParamMap` in the grid, it fits the estimator on each fold's training split, scores with the evaluator, and averages the k metrics. The winning `ParamMap` is then used to **re-fit the estimator on the entire dataset**, and that final fitted object is `bestModel`. ## Why the pipeline belongs inside Because `Pipeline` is an Estimator, you can tune an entire pipeline at once rather than tuning each element separately. The correctness argument is data leakage. Many feature stages are Estimators — they learn from data. `StandardScaler` learns each column's mean and standard deviation. `StringIndexer` learns a label vocabulary and its frequency ordering. `IDF` learns document frequencies. `Imputer` learns the mean, median or mode it substitutes. `CountVectorizer` learns a vocabulary. If you fit those stages once over the whole dataset, transform everything, and only then hand the classifier to `CrossValidator`, every fold's "held-out" third has already contributed to the scaler's mean, the indexer's vocabulary and the imputer's median. The evaluation is no longer measuring performance on unseen data, and the reported metric is optimistic — sometimes mildly, sometimes catastrophically when the leaking statistic is target-related or when the dataset is small. Put the pipeline inside the `CrossValidator` and each fold re-fits every stage on that fold's training rows only, which is the honest estimate. The second argument is reach. Once the pipeline is the estimator, the grid can address parameters on *any* stage, because parameters are scoped to instances: `ParamGridBuilder().addGrid(hashingTF.numFeatures, [10, 100, 1000]).addGrid(lr.regParam, [0.1, 0.01])` tunes featurisation and regularisation jointly. Feature-space size and model regularisation interact strongly, so tuning them separately finds a worse optimum than tuning them together. ## The cost, and how to control it Spark's own guide is blunt: cross-validation over a grid is expensive. The grid above has 3 × 2 = 6 combinations; at 2 folds that is 12 model fits. With realistic grids and the common k=3 or k=10 the number climbs fast, and each fit re-runs the *whole* pipeline, including the shuffles in the feature stages. Levers: - **`parallelism`**. By default parameter sets are evaluated in *serial* (a value of 1). Setting `parallelism` to 2 or more evaluates several `ParamMap`s concurrently. Choose it to maximise parallelism without exceeding cluster resources; larger values do not always help, since each individual fit already parallelises across the cluster, and generally a value up to 10 suffices for most clusters. - **Caching**. The input DataFrame is read repeatedly across folds and grid points; persisting it (and any expensive deterministic preprocessing that is genuinely *not* learned, such as parsing or joins) avoids recomputing it every fit. - **A smaller grid**. A coarse search followed by a narrow refinement usually beats one exhaustive sweep for the same cluster hours. - **`TrainValidationSplit`**. It evaluates each parameter combination once against a single (training, test) pair split by `trainRatio` — with `trainRatio=0.75`, 75% trains and 25% validates. It is therefore roughly k times cheaper than `CrossValidator`, but it will not produce results as reliable when the training dataset is not sufficiently large. Like `CrossValidator`, it finishes by fitting the estimator on the entire dataset with the best `ParamMap`. ## Things candidates get wrong The first is describing `bestModel` as "the best fold's model". It is not; the folds only *score* parameter combinations, and the returned model is a fresh fit on all the data using the winning parameters. The second is forgetting that when the estimator was a `Pipeline`, `bestModel` is a `PipelineModel` — the deployable artifact, feature stages included. The third is leaving the evaluator's metric implicit: set `metricName` deliberately, because selecting on the wrong metric quietly optimises the wrong thing, and an imbalanced binary problem selected on the wrong criterion can look excellent and be useless. ## How to answer Open with "a Pipeline is an Estimator, so it can be the thing being cross-validated". Then give the leakage mechanism with a concrete stage — the scaler's mean — say what the correct arrangement does instead, add the joint-tuning benefit, and close by acknowledging the multiplicative cost and naming two levers you would actually pull.
- How many models does a Spark CrossValidator fit for a 3-by-2 grid over 3 folds?Eighteen during the search — six parameter combinations times three folds — and then one more refit of the estimator on the entire dataset using the winning `ParamMap`, so nineteen fits in total. Each of those is a full run of the pipeline, feature shuffles included, which is why grid size and fold count are budget decisions rather than defaults to accept.
- What does the parallelism parameter on CrossValidator change, and what limits it?It controls how many `ParamMap`s are evaluated concurrently; the default is serial evaluation. Raising it trades cluster resources for wall-clock time, and it should be chosen to maximise parallelism without exceeding what the cluster can hold, since each individual fit already parallelises across executors. Spark's guidance is that a value up to about ten suffices for most clusters.
- When is TrainValidationSplit the better choice than CrossValidator?When the search cost dominates and the dataset is large. `TrainValidationSplit` makes a single training/test pair using `trainRatio` and evaluates each parameter combination once instead of k times, so it is roughly k times cheaper. The tradeoff is reliability: with a training set that is not sufficiently large, one split gives a noisier estimate than averaging over folds.
- What exactly is CrossValidator.bestModel when the estimator was a Pipeline?A `PipelineModel` — not the model from the winning fold. After the folds identify the best `ParamMap`, the estimator is re-fitted on the entire dataset with those parameters, so `bestModel` contains every feature stage fitted on all the data plus the trained final model, which is exactly the artifact you would save and deploy.
Fitting the scaler on all the data before cross-validating is like letting a student revise from the exam paper: the score stops measuring what they would do on questions they have not seen.
saying these in an interview costs you the question
- Fits the scaler on all the data, then cross-validates only the classifier
- Says cross-validation tunes the model but not the feature stages
- Ignores that fold count multiplies the grid size
- Raises parallelism blindly to make the search finish sooner
- Thinks bestModel is the highest-scoring fold's model