Зачем это нужно
Spark — не только ETL. В ML-контуре он готовит training datasets, считает batch features, запускает распределённое обучение (Spark MLlib) и batch inference на миллионах строк. MLOps-инженер связывает Spark jobs с orchestrator (Argo Workflows), registry (ClearML/MMS) и Feature Store (Feast).
Урок показывает, где Spark стоит в end-to-end ML lifecycle — от сырого lake до артефакта для train job.
Основные идеи
Spark в ML lifecycle (типичные роли).
| Этап | Spark job | Output |
|---|---|---|
| Feature engineering | Aggregations, joins, encoding | Parquet в gold / Feast offline |
| Train/test split | Stratified sample, time-based split | train/, test/ paths |
| Batch scoring | Load model, predict on partition | Scores table |
| Hyperparameter (limited) | Spark MLlib CV | Model in MLlib format |
Большинство команд используют Spark для data prep, а обучение — XGBoost/LightGBM/PyTorch на выгрузке или через sparkdl.
Point-in-time correctness (критично для фич). При join исторических фич с label нельзя «подглядывать в будущее». Нужны as-of joins или pre-computed feature tables с feature_timestamp <= label_timestamp. Feast materialization решает часть проблемы (6.4).
Train snapshot contract:
s3a://ml-gold/churn/train/dataset_v5/ _SUCCESS part-*.parquet manifest.json # schema, row count, label distribution, git sha pipeline
Train job (ClearML) читает immutable snapshot; повторный train на том же path = воспроизводимость.
Spark → single-node train handoff.
Spark пишет Parquet (сжатый, partitioned).
Train Pod (GPU) читает через
s3/pandas/polarsилиspark.readтолько нужные columns.Альтернатива: Spark MLlib train entirely in Spark — реже для deep learning.
Batch inference pattern.
` from pyspark.ml import PipelineModel
model = PipelineModel.load("s3a://ml-models/churn/v44/spark_pipeline/") scored = model.transform(input_df) scored.select("customer_id", "prediction", "probability").write.parquet( "s3a://ml-scores/dt=2026-01-15/" ) `
Для tree models часто export ONNX → Triton (модуль 7); Spark batch — когда scoring SQL-friendly и без GPU.
Orchestration с Argo Workflows.
DAG: validate-raw → spark-feature-job → data-quality-gate → export-to-feast → trigger-clearml-train
Каждый step — отдельный Pod; Spark job может быть SparkApplication CR (Spark Operator).
Integration с ClearML. Spark driver логирует: input paths, row counts, output path как ClearML artifacts. MMS получает ссылку на dataset version после quality gate.
Spark Structured Streaming + ML (edge case). Online learning редок; чаще streaming пишет features, offline train периодически.
Как это выглядит на практике
Weekly retrain pipeline:
Воскресенье 02:00 — Argo CronWorkflow стартует.
Spark job
churn-features-weekly: 90 days history, 200M rows, 45 min on 10 executors.Quality gate (6.1.2): schema hash, null rates, label balance.
Feast
materialize-incrementalдля offline store sync.ClearML task: train LightGBM на
dataset_v5, log metrics, registerchurn/v45.Jenkins/MMS approval → deploy inference (модуль 7).
Failure handling:
Spark job failed stage 3 → retry с
--conf spark.speculation=trueесли straggler.Quality gate failed → Slack data team, no train.
Partial write без
_SUCCESS→ downstream не читает (idempotent paths сdt=partition).
Cost awareness. 10 executors × 4h × $ = bill. Для учебного кластера MDP — лимиты namespace ResourceQuota; для prod — spot instances / autoscaling.
Что сделать после занятия
Нарисуйте DAG: Spark feature job → quality → train → register model.
Объясните, зачем train читает immutable snapshot, а не «последний Parquet в bucket».
Выпишите 3 метрики Spark job, которые вы бы логировать в ClearML.