Зачем это нужно
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-системе.
Ingestion — приём из источников (API, CDC, файлы в object storage).
Processing — трансформации, агрегаты, feature engineering.
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:
Откуда inference берёт 12 фич? (online store vs API vs in-request)
Как часто обновляется каждая фича? (real-time vs T-1 day)
Что произойдёт, если batch job опоздал на 6 часов?
Где хранится snapshot обучающей выборки для воспроизводимости?
Типичная ошибка новичка: считать, что «данные для обучения» и «данные для inference» — один и тот же SQL-запрос, просто в разное время. На деле join'ы, фильтры и окна агрегации часто расходятся.
Что сделать после занятия
Нарисуйте data flow для учебного ML-проекта: источники → обработка → train → serve.
Отметьте на схеме, где batch, а где streaming.
Выпишите 3 вопроса про lineage, на которые ваша команда не может ответить за 5 минут — это зоны риска.