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

Пути данных в 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 минут — это зоны риска.

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