06 · Данные для ML6.1–6.4 · Потоки, батч и признаки6.3.2сложный

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.

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