Most of a tabular model's accuracy comes from its features, and most of the bugs in tabular models come from them too. Feature engineering in Spark ML means two kinds of work: computing raw signals from large tables with DataFrame and SQL operations, then turning them into the single numeric vector column that every Spark ML algorithm consumes, using pipeline stages that learn statistics such as medians, vocabularies and category means.
This article walks that whole path on one example, a subscriber churn model. It covers point-in-time aggregation so features never see the future, numeric and categorical treatment, the TargetEncoder added in Spark 4.0 and why its defaults leak, interactions and feature selection, and the performance and failure modes that show up at scale. The Spark ML pipelines guide covers the Transformer and Estimator contract in detail, and the MLlib overview covers hashing and the full transformer catalog; this page focuses on doing feature work correctly.
From columns to one feature vector
Spark ML algorithms read one column of type vector, conventionally features, plus a numeric label. Feature engineering therefore ends with a VectorAssembler that concatenates numeric columns and vectors into that column. Everything before it falls into two groups.
Stateless transforms depend only on the row: a log, a ratio, a date difference, a flag. Write them in SQL or as a SQLTransformer stage. Learned transforms need statistics from data: the median for imputation, the category order for StringIndexer, quartiles for scaling, the per-category label mean for target encoding. These are Estimators. Their statistics must be computed on training rows only and then frozen, which is exactly what fitting a Pipeline and reusing the resulting PipelineModel does.
The most damaging leaks happen before the pipeline, though, in how raw features are aggregated.
Point-in-time features and time-based splits
Churn is defined per customer at a cutoff date: did the customer cancel within 30 days after the cutoff? Every feature must use only events strictly before that cutoff. Aggregating all history and then joining labels is the classic leak, because a customer who churned stops transacting, and a 'days since last purchase' computed over all data quietly encodes the answer.
Make the cutoff part of the join condition, and use a left join so customers with no recent events keep a row:
from pyspark.sql import functions as F
l = spark.table("labels").alias("l") # customer_id, cutoff_ts, churned, plan, country
e = spark.table("events").alias("e") # customer_id, ts, amount, channel
j = l.join(e, (F.col("l.customer_id") == F.col("e.customer_id"))
& (F.col("e.ts") < F.col("l.cutoff_ts")), "left")
in30 = F.col("e.ts") >= F.col("l.cutoff_ts") - F.expr("INTERVAL 30 DAYS")
feats = (j.groupBy("l.customer_id", "l.cutoff_ts", "l.churned", "l.plan", "l.country")
.agg(F.count(F.when(in30, 1)).alias("txn_30d"),
F.sum(F.when(in30, F.col("e.amount"))).alias("spend_30d"),
F.countDistinct("e.channel").alias("channels"),
F.max("e.ts").alias("last_ts"))
.withColumn("days_since_last", F.datediff("cutoff_ts", "last_ts")))
train = feats.where(F.col("cutoff_ts") < "2026-05-01") # 30-day label window
test = feats.where(F.col("cutoff_ts") >= "2026-06-01") # one-month gapSplit by time, not randomly. A random split puts the same customer's March and April rows on both sides, so the model is evaluated on people it has effectively seen. With several cutoffs per customer, a time split plus a gap of at least the label window (30 days here) between train and test cutoffs is the honest setup. Event tables are often skewed toward a few heavy customers; if one task runs far longer than the rest, apply the fixes in handling data skew in Spark.
Numeric features: missing values, skew and scale
Decide why each value is missing before choosing how to fill it. spend_30d is null when a customer had no events, which means zero, not 'unknown'; fill it with 0 in SQL. days_since_last is null when the customer never transacted, which is informative; keep a flag for it and impute the number.
Imputer supports the strategies mean (the default), median and mode, treats nulls as missing in addition to its missingValue (NaN by default), and computes medians with approximate quantiles at a relative error of 0.001. It works on numeric columns only.
Heavy-tailed money values hurt linear models and distance-based methods. A log1p in SQL tames them. RobustScaler then centres on the median and scales by the range between the 25th and 75th percentiles (its defaults are lower=0.25, upper=0.75, withCentering=False, withScaling=True), so a few huge accounts do not dominate the scale. Scale only continuous columns: dividing one-hot indicators by an interquartile range is meaningless, and for a rare category both quartiles are zero. Tree ensembles need no scaling at all; skip it for them.
When the relationship is non-monotonic, QuantileDiscretizer turns a number into quantile buckets that a linear model can weight separately, and Bucketizer does the same with split points you choose, which keeps buckets stable across retrains.
Categorical features and TargetEncoder
StringIndexer maps strings to indices, most frequent first by default (frequencyDesc). Its handleInvalid defaults to error, so a category that first appears after training crashes scoring; set keep to map unseen values to an extra index. OneHotEncoder then produces sparse indicator vectors; its defaults are dropLast=True and handleInvalid=error. Drop the last category for unregularized linear models to avoid perfectly collinear columns; keep all of them for regularized or tree models.
One-hot fails for high-cardinality columns such as country, merchant or postcode: thousands of columns, most nearly empty. Target encoding replaces each category with the label's mean for that category, blended towards the overall mean for rare categories. Spark 4.0 added TargetEncoder. Its input must already be category indices, so it follows a StringIndexer; it supports binary and continuous targets, and it blends with weight n / (n + m), where n is the category count and m is the smoothing parameter.
Two traps. First, smoothing defaults to 0.0, so every category gets its raw mean, and a category seen once encodes its own label exactly. Second, the encoder is fitted and applied on the same training rows, so each row's own label is inside its feature. Trees and boosted models find that shortcut immediately: training metrics soar and test metrics fall. Always set smoothing explicitly, and for strong learners use out-of-fold encoding on the training set.
Worked numbers. Say the overall churn rate is 0.08 and Iceland has 12 rows with 6 churners. With smoothing 0 the encoding is 0.5. With m = 20, the weight is 12 / 32 = 0.375, and the encoding is 0.375 x 0.5 + 0.625 x 0.08 = 0.2375, a much more believable estimate from 12 rows.
Out-of-fold target encoding
Out-of-fold encoding gives each training row a value computed without its own fold. Hash a stable ID into k folds, compute per-fold sums, and subtract:
k, m = 5, 20.0
d = train.withColumn("fold", F.pmod(F.xxhash64("customer_id"), F.lit(k)))
prior = d.agg(F.avg("churned")).first()[0]
per_fold = d.groupBy("country", "fold").agg(F.sum("churned").alias("s"),
F.count("*").alias("n"))
totals = per_fold.groupBy("country").agg(F.sum("s").alias("S"), F.sum("n").alias("N"))
oof = (per_fold.join(totals, "country")
.withColumn("country_te",
(F.col("S") - F.col("s") + m * prior) / (F.col("N") - F.col("n") + m))
.select("country", "fold", "country_te"))
train_enc = d.join(oof, ["country", "fold"], "left").fillna({"country_te": prior})
full_map = totals.withColumn("country_te", (F.col("S") + m * prior) / (F.col("N") + m))
test_enc = test.join(full_map.select("country", "country_te"), "country", "left") \
.fillna({"country_te": prior})The formula (s + m x prior) / (n + m) is the same blend TargetEncoder uses. The cost is that this logic lives outside the pipeline, so you must save full_map with the model and join it at serving time. Use TargetEncoder with explicit smoothing for regularized linear models, where the leak is mild, and out-of-fold encoding for boosted trees.
Worked example: the full churn pipeline
With raw features in place, the pipeline assembles, scales and models. Continuous and categorical parts are assembled separately so only continuous values are scaled:
from pyspark.ml import Pipeline
from pyspark.ml.feature import (SQLTransformer, Imputer, StringIndexer, OneHotEncoder,
TargetEncoder, VectorAssembler, RobustScaler)
from pyspark.ml.classification import LogisticRegression
prep = SQLTransformer(statement='''
SELECT *, log1p(coalesce(spend_30d, 0)) AS log_spend,
CAST(days_since_last IS NULL AS DOUBLE) AS never_active
FROM __THIS__''')
impute = Imputer(strategy="median", inputCols=["days_since_last"], outputCols=["days_i"])
plan_ix = StringIndexer(inputCol="plan", outputCol="plan_ix", handleInvalid="keep")
plan_oh = OneHotEncoder(inputCols=["plan_ix"], outputCols=["plan_vec"], handleInvalid="keep")
ctry_ix = StringIndexer(inputCol="country", outputCol="ctry_ix", handleInvalid="keep")
ctry_te = TargetEncoder(inputCols=["ctry_ix"], outputCols=["ctry_te"], labelCol="churned",
targetType="binary", smoothing=20.0, handleInvalid="keep")
num_vec = VectorAssembler(inputCols=["txn_30d", "log_spend", "channels", "days_i"],
outputCol="num_raw")
scale = RobustScaler(inputCol="num_raw", outputCol="num_scaled")
assemble = VectorAssembler(inputCols=["num_scaled", "never_active", "plan_vec", "ctry_te"],
outputCol="features")
lr = LogisticRegression(labelCol="churned", regParam=0.01)
model = Pipeline(stages=[prep, impute, plan_ix, plan_oh, ctry_ix, ctry_te,
num_vec, scale, assemble, lr]).fit(train)
model.write().overwrite().save("gs://ml-models/churn/v7")Tune smoothing, regularization and bucket counts with time-aware validation as described in Spark ML cross-validation, always fitting the whole pipeline inside each fold so the encoders are refitted too.
Interactions and feature selection
Interaction multiplies columns, for example plan indicators by log spend, so a linear model can learn a different spend slope per plan. Use it sparingly, because the output width is the product of the input widths.
To prune features, VarianceThresholdSelector drops near-constant columns, and UnivariateFeatureSelector scores each feature against the label: chi-squared for categorical features and label, ANOVA F for continuous features with a categorical label, and the F-value for continuous both sides. Its default mode keeps the top 50 features. Univariate scores ignore interactions, so treat them as a first cut and confirm with validation metrics or model-based importance.
Performance at scale
- Cache the feature table before fitting. Each Estimator in
Pipeline.fitruns at least one job over its input; without caching, the point-in-time join is recomputed for each of them. - Quantile cost. Imputer medians, RobustScaler and QuantileDiscretizer each compute approximate quantiles; relative error trades accuracy for memory.
- Keep vectors sparse. One-hot output is sparse; centering with
withCentering=TrueorStandardScalerwithwithMean=Truemakes it dense and can multiply memory many times. - Set
VectorAssemblerhandleInvaliddeliberately. The defaulterrorfails on nulls;skipsilently drops rows, which changes your evaluation population;keepwrites NaN that many models reject. Prefer fixing nulls upstream.
Failure modes
- Great offline, poor in production. Usually a point-in-time leak or random split. Recompute features with the cutoff in the join and re-evaluate on a later time window.
- Scoring crashes on a new category.
StringIndexerorOneHotEncoderleft athandleInvalid=error. - Boosted model overfits a categorical column. Target encoding with default smoothing, fitted in-sample. Smooth and use out-of-fold encoding.
- Training-serving skew. Serving recomputes features with different SQL or freshness. Generate both from one feature definition and compare distributions daily.
- Silent row loss. An inner join to events or
handleInvalid=skipremoves the inactive customers, who are the likeliest churners.
Trade-offs
SQL or pipeline stage. Aggregations over event history belong in SQL or a feature table, because they need joins and windows that pipeline stages cannot express. Anything that learns a statistic belongs in the pipeline, so the statistic is fitted on training rows and saved with the model. Row-level arithmetic can go either way; putting it in a SQLTransformer keeps it versioned with the model, which helps serving.
Vocabulary, hashing or target encoding. Indexing plus one-hot is exact and interpretable but grows with cardinality and needs a fitted vocabulary. Feature hashing has a fixed width and no fit, at the price of collisions. Target encoding compresses any cardinality into one column per target but adds leakage risk and a dependency on the label distribution, which drifts. For very large vocabularies with strong learners, out-of-fold target encoding usually wins; for linear models with a few hundred categories, one-hot with regularization is simpler and safer.
Spark or a single machine. Spark earns its overhead when the event history is too large for one machine. If the final training table fits in memory, a common split is Spark for point-in-time aggregation and a single-node library for the last modelling steps, provided the same feature definitions are reused for serving.
What to do next
- Define the prediction time for every label row and add it to every feature join as a strict cutoff.
- Split by time with a gap of at least the label window.
- For each feature, write down why it can be missing and fill it accordingly, with flags for informative nulls.
- Set
handleInvalidexplicitly on every indexer, encoder and assembler. - Use one-hot for low cardinality and smoothed or out-of-fold target encoding for high cardinality.
- Scale only continuous columns, and only for models that need it.
- Fit one Pipeline on train, save the PipelineModel with any external maps, and use them unchanged for test and serving.