6.1.1 · блок 6
Пути данных в ML-системе
Пути данных в ML-системе
Зачем это нужно
ML-модель «видит» только то, что до неё дошло. На практике данные проходят десятки систем: CRM, логи приложения, DWH, Kafka, Feature Store, inference API. Если не понимать путь данных (data lineage), невозможно ответить на простые вопросы: «почему score изменился?», «какая версия фичи в prod?», «можно ли воспроизвести обучение?».
Для начинающего MLOps-инженера путь данных — это карта местности. Без неё вы строите пайплайны «вслепую» и ловите расхождения train/serve уже после релиза.
Основные идеи
Batch vs streaming — два режима движения данных.
| Режим | Когда данные приходят | Типичный инструмент | Пример в ML |
|-------|----------------------|---------------------|-------------|
| **Batch** | Пакетами по расписанию | Spark, SQL, Airflow/Argo | Ночной пересчёт фич, retrain раз в неделю |
| **Streaming** | Непрерывно, событие за событием | Kafka, Flink | Online scoring, fraud detection в реальном времени |
| **Near real-time** | Микробатчи (секунды–минуты) | Spark Structured Streaming | Агрегаты «за последний час» для рекомендаций |
Data pipeline vs ML pipeline. Data pipeline готовит сырые данные: очистка, дедупликация, join. ML pipeline использует подготовленные данные: train, validate, export model. Это разные пайплайны с разными SLA, но общий контракт на схему.
Три зоны данных в типичной ML-системе.
1. Ingestion — приём из источников (API, CDC, файлы в object storage).
2. Processing — трансформации, агрегаты, feature engineering.
3. Serving — данные или фичи доступны модели в runtime (online store, кэш, batch table).
Lineage (происхождение). Для каждой фичи полезно знать: исходная таблица → трансформация → версия → потребитель (train job / inference). Без lineage debugging «train F1=0.9, prod F1=0.6» занимает недели.
Контракт на схему. Поля имеют тип, nullable, допустимый диапазон. Изменение схемы upstream (добавили колонку, переименовали поле) — breaking change для downstream ML, если не версионировать.
Cold path vs hot path.
- Cold path — тяжёлые batch-джобы: обучение, backfill фич за год.
- Hot path — latency-sensitive: запрос пользователя → фичи → inference < 200 ms.
Одна и та же логика агрегации часто реализуется дважды (batch + streaming) — классический источник train/serve skew. Feature Store (модуль 6.4) частично решает это.
Data lake / lakehouse. Сырые данные лежат в object storage (MinIO, S3) в формате Parquet/Delta. ML читает оттуда через Spark или SQL-движок. Важно разделять bronze (raw), silver (cleaned), gold (business-ready) — хотя бы концептуально.
Как это выглядит на практике
Сценарий: скоринг оттока клиентов.
CRM (PostgreSQL)
│ CDC / nightly export
▼
Object Storage (Parquet, bronze/)
│ Spark job (06-data pipeline)
▼
Feature table: customer_features_daily (gold/)
│ ├─► ClearML train job (batch, раз в неделю)
│ └─► Feast offline store
│
Kafka topic: customer.events
│ streaming aggregator
▼
Feast online store (Redis)
│ feature lookup
▼
Inference API (KServe + MLServer)
Вопросы, которые задаёт инженер на design review:
1. Откуда inference берёт 12 фич? (online store vs API vs in-request)
2. Как часто обновляется каждая фича? (real-time vs T-1 day)
3. Что произойдёт, если batch job опоздал на 6 часов?
4. Где хранится snapshot обучающей выборки для воспроизводимости?
Типичная ошибка новичка: считать, что «данные для обучения» и «данные для inference» — один и тот же SQL-запрос, просто в разное время. На деле join'ы, фильтры и окна агрегации часто расходятся.
Что сделать после занятия
- [ ] Нарисуйте data flow для учебного ML-проекта: источники → обработка → train → serve.
- [ ] Отметьте на схеме, где batch, а где streaming.
- [ ] Выпишите 3 вопроса про lineage, на которые ваша команда не может ответить за 5 минут — это зоны риска.