MLOps Path

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:

Cost awareness. 10 executors × 4h × $ = bill. Для учебного кластера MDP — лимиты namespace ResourceQuota; для prod — spot instances / autoscaling.

Что сделать после занятия

Официальные материалы

Открыть интерактивную версию