6.3.2 · блок 6
Spark для ML-пайплайна
Spark для ML-пайплайна
Зачем это нужно
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.
1. Spark пишет Parquet (сжатый, partitioned).
2. Train Pod (GPU) читает через s3/pandas/polars или spark.read только нужные columns.
3. Альтернатива: 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:
1. Воскресенье 02:00 — Argo CronWorkflow стартует.
2. Spark job churn-features-weekly: 90 days history, 200M rows, 45 min on 10 executors.
3. Quality gate (6.1.2): schema hash, null rates, label balance.
4. Feast materialize-incremental для offline store sync.
5. ClearML task: train LightGBM на dataset_v5, log metrics, register churn/v45.
6. 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.