Конвейеры подготовки признаков: DAGs и оркестрация
Конвейеры подготовки признаков — это управляемые, повторяемые и контролируемые последовательности задач, которые превращают сырой поток данных в качественные, пригодные для обучения и инференса признаки. В контексте Lakehouse они становятся важной связующей частью между данными в хранилище и моделями, которые эти данные учат и применяют на проде. Правильная оркестрация DAGs (Directed Acyclic Graphs) обеспечивает:
- повторяемость и воспроизводимость экспериментов;
- прозрачность lineage признаков (от источника данных до модели);
- разделение ролей: аналитики формируют признаки, data scientists — эксперименты и выбор моделей, инженеры — инфраструктуру и качество данных;
- баланс между batch-обработкой и стриминговыми обновлениями;
- безопасное и управляемое обновление онлайн-слоя признаков без нарушения производительности сервиса.
Эта глава даст теоретические основы DAGs и оркестрации, а затем перейдёт к практическим примерам с указанием открытых инструментов и российских решений, сценариев развёртывания, ограничений и рисков.
Что такое DAG в контексте подготовки признаков
DAG (Directed Acyclic Graph) — граф, где рёбра направлены и не образуют циклов. В контексте подготовки признаков DAG представляет набор задач (tasks), каждая из которых имеет входы и выходы, и зависимости между задачами задают порядок выполнения. Ключевые характеристики DAG для признаков:
- идемпотентность: повторный запуск даёт идентичный результат;
- детерминированность: одинаковые входы — одинаковые признаки;
- явная зависимость от источников данных, вычисляемых параметров и конфигурации;
- переменная длительность выполнения и возможность параллельного исполнения;
- поддержка возвратов и повторного выполнения без side effects.
Архитектура конвейеров признаков
Общая картина в Lakehouse-подходе:
- Источники данных (raw), обычно в Data Lake/Raw-хранилище (S3/ADLS/Яндекс.ObjectStorage);
- Offline Store — хранилище для обучения и репродукции признаков (чаще Parquet/ORC в дата-логах, может быть Hive/Delta/Apache Iceberg);
- Online Store — низкол latency хранилище признаков для онлайн-инференса (Redis, Cassandra, DynamoDB и т.п.);
- Feature Registry/Store — база описаний признаков: имена, типы, входы, версии, TTL;
- Конвейер обработки признаков — ETL/ELT-пайплайн, который вычисляет признаки и записывает их в Offline и Online Stores;
- Инструменты контроля качества данных, тестирования и мониторинга.
Ключевые концепты:
- Feature Definition (описание признака): имя, тип данных, источник входных данных, формула расчёта, версия, TTL;
- Feature View (обозначение набора признаков, сборка из нескольких признаков, доступ к ним через единый API);
- Online против Offline: offline-признаки для обучения и ретро-поиск, online-признаки — для продакшна;
- Feature Serving слои: механизм подстановки признаков в модели и сервисы по инференсу;
- Data Quality Gates: проверки корректности данных до публикации признаков;
- Data Lineage и Versioning: прослеживаемость происхождения признаков, влияние изменений на модели и выводы;
- Контракты данных: соглашения об ожидаемом формате данных и уровне качества.
Методологии и принципы
- ELT-подход к подготовке признаков: извлечение и загрузка, затем вычисление признаков в рамках гибкой обработки; баланс между вычислительной стоимостью и скоростью обновления.
- Управление версиями признаков: версионирование по имени признака и версии, поддержка отката к предыдущим версиям в случае проблем.
- Контракты данных и тестирование: использование Great Expectations или аналогов для проверки входных и выходных данных.
- Нормализация и стандартизация признаков: единая сигнатура типов, обработка пропусков, масштабирование и кодирование категориальных признаков.
- Метрики качества признаков: корреляции, устойчивость к дрейфу, влияние на качество моделей (производительность, лаги).
- Контроль доступа и безопасность: разграничение прав на источники данных, конфигурации конвейеров, журналирование и аудит.
Инструменты и подходы (обзор)
- Оркестраторы: Apache Airflow, Dagster, Kedro (для MLOps-подходов), Prefect.
- Фреймворки для признаков: Feast (open-source feature store) и его альтернативы/обертки; интеграции с Spark/BigQuery/Delta Lake.
- Хранилища и слои хранения: Parquet/Delta/Apache Iceberg в Data Lake; Redis/Cassandra/DynamoDB как online store.
- Контроль качества: Great Expectations, Deequ, OpenLineage для lineage и мониторинга.
- Эксперименты и версии моделей: MLflow, Kubeflow, Metaflow (и локальные решения).
- Российские платформы: Яндекс DataSphere и облачные решения СберОблако/MLOps позволяют интегрировать конвейеры, данные и эксперименты внутри локальной инфраструктуры и в рамках российского облака; возможность использования локальных хранилищ и соблюдения регуляторных требований.
| Техника | Что даёт | Примечания |
|---|---|---|
| DAG-оркестрация (Airflow, Dagster, Kedro) | управление задачами, мониторинг, повторяемость | выбор зависит от экосистемы и интеграций |
| Feature Store ( Feast и аналоги) | централизованный доступ к признакам, онлайн/оффлайн слои | версия признаков, совместная работа аналитиков и DS |
| Offline Store | хранение больших массивов признаков для обучения | Parquet/Delta/ICEBERG, lakehouse-архитектура |
| Online Store | быстрый доступ к признакам во время инференса | Redis/Cassandra/DS, latency 1-10 мс |
| Контроль качества | проверки данных перед публикацией | Great Expectations/OpenLineage |
| Эксперименты и версии моделей | воспроизводимость, сравнение подходов | MLflow, Kubeflow, Metaflow |
Практические примеры
Ниже рассмотрены два кейса: первый — классическая открытая экосистема на базе Airflow + Feast; второй — пример использования российских платформ, ориентированных на локальную инфраструктуру и облако.
Пример на Open-Source стеке: Airflow + Feast + Kedro + Spark
Цель: построить конвейер подготовки признаков для задачи персонализированных рекомендаций на Lakehouse, где признаки формируются пакетно (еженедельно) и обновляются для онлайн-использования.
Архитектура и данные
- Источник: логи взаимодействий пользователя, транзакционные таблицы.
- Offline Store: Parquet в Data Lake (S3/ADLS) или Delta Lake.
- Online Store: Redis для инференса с очень низкой задержкой.
- Feature Registry: Feast для описания признаков, версий и зависимостей.
- Оркестрация: Apache Airflow.
Пошаговый пайплайн
- Стадия 1: Ingest raw data из источников в Data Lake.
- Стадия 2: Применение трансформаций (Grouping, Window functions, агрегаты) с использованием Spark.
- Стадия 3: Определение признаков через Feast: создаём метаданные признаков, формулы и официальную версию.
- Стадия 4: Заполнение offline store (Parquet/Delta) и публикация в Feast.
- Стадия 5: Обновление online store (Redis) для продакшна.
- Стадия 6: Валидации качества и тесты на линии признаков.
- Стадия 7: Использование признаков в обучении (training) и эксперименты с моделями.
Пример кода: Airflow DAG (Python)
# airflow_dag_feature_pipeline.py
from datetime import datetime, timedelta
from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.operators.bash import BashOperator
def extract_raw():
# код источников, загрузка данных в_raw
pass
def transform_features():
# Spark-трансформации, создание признаков
pass
def push_to_offline_store():
# сохранение в Parquet/Delta, публикация в Feast offline store
pass
def update_online_store():
# обновление Redis онлайн склада
pass
def quality_checks():
# проверки данных
pass
default_args = {
'owner': 'ml-team',
'depends_on_past': False,
'start_date': datetime(2025, 1, 1),
'retries': 1,
'retry_delay': timedelta(minutes=15),
}
with DAG('feature_conveyor_open_source',
default_args=default_args,
schedule_interval='0 2 * * *',
catchup=False) as dag:
t1 = PythonOperator(task_id='extract_raw', python_callable=extract_raw)
t2 = PythonOperator(task_id='transform_features', python_callable=transform_features)
t3 = PythonOperator(task_id='push_to_offline', python_callable=push_to_offline_store)
t4 = PythonOperator(task_id='update_online', python_callable=update_online_store)
t5 = PythonOperator(task_id='quality_checks', python_callable=quality_checks)
t1 >> t2 >> t3 >> t4 >> t5
Пример определения признаков в Feast (feature.yaml)
# feature_definitions.yaml
entities:
- name: user_id
join_key: user_id
description: Unique user identifier
features:
- name: user_avg_session_duration
uid: 1
value_type: DOUBLE
input:
- table: db.analytics.user_session
timestamp: event_time
keys: [user_id]
transform:
expression: "avg(session_duration) over (partition by user_id)"
ttl: 30d
- name: user_total_purchases_last7d
uid: 2
value_type: INT64
input:
- table: db.transactions
timestamp: event_time
keys: [user_id]
transform:
expression: "sum(amount) filter where event_time > now()-7d"
ttl: 14d
Контроль качества и тестирование
- Great Expectations тесты на входные данные и границы значений;
- OpenLineage для lineage и мониторинга.
Роль Kedro
- Kedro даёт структуру проекта и повторяемый код для подготовки данных и признаков, обеспечивает модульность и тестируемость.
Преимущества данного подхода:
- ясная диспозиция задач и зависимостей;
- единый источник truth через Feast;
- возможность использования как batch, так и микро-батч режимов;
- прозрачность и аудит изменений признаков.
Недостатки:
- сложность развёртывания и поддержки, особенно в больших командах;
- потребность в грамотной стратегии таймингов и backfill;
- риск задержек из-за зависимостей между задачами.
Пример на российской платформе: Яндекс DataSphere / СберОблако
Цель: продемонстрировать развёртывание конвейера на российской платформе, сохраняя совместимость с открытыми инструментами и соблюдением локальных регуляторных требований.
Архитектура и данные
- Источники: логи сервиса и бизнес-данные, размещённые в облачном хранилище российского провайдера.
- Offline Store: локальные хранилища на платформе (напр., Yandex Object Storage) с возможностью экспорта в Parquet/Delta.
- Online Store: Redis или аналогичный высокопроизводительный слот на платформе.
- Оркестрация: встроенная в Яндекс DataSphere или СберОблако инфраструктура конвейеров (вариант — интеграция с внешними инструментами).
- Контроль качества и lineage: средства платформы и интегрированные внешние решения.
Пошаговый сценарий развертывания
- Создать проект и определить конфигурации конвейера под требования регуляторов;
- Подключиться к источникам данных через коннекторы DataSphere;
- Определить признаки через централизованный реестр признаков;
- Развернуть пайплайн: обработка данных, вычисление признаков, сохранение в offline/online stores;
- Включить проверки качества и аудит изменений;
- Интегрировать пайплайн с обучением и инференсом.
Технические детали и примеры
Пример YAML-конфигурации пайплайна в DataSphere (условно, синтаксис может отличаться в зависимости версии):
pipeline:
name: feature_pipeline_rs
schedule: "0 3 * * *"
stages:
- name: extract_raw
type: data_ingest
source: "rs_service_logs"
- name: feature_engineering
type: spark_transform
script: "transform_features.py"
- name: store_offline
type: write_offline
target: "yds_parquet"
- name: push_online
type: write_online
target: "redis_cluster"
- name: quality_checks
type: data_quality
-
Встроенные средства мониторинга и lineage позволяют отслеживать происхождение признаков и влияние изменений на модели.
-
Варианты интеграции: DataSphere может работать с внешними инструментами (Airflow/Dedicated pipelines) через коннекторы и API, что позволяет строить гибридные архитектуры.
Преимущества российского решения:
- соответствие требованиям локализации, безопасности и регуляторным нормам;
- интеграция с отечественными облачными сервисами и хранилищами;
- поддержка сценариев совместной работы аналитиков и инженеров в одном пространстве.
Ограничения и вызовы:
- возможно меньшая экосистема готовых интеграций по сравнению с англосаксонскими экосистемами;
- требования к лицензированию и стоимости использования сервисов;
- необходимость локальной экспертизы и адаптации к инфраструктуре организации.
Архитектура конвейера признаков
- Источники данных (data sources)
- Промежуточные вычисления (transforms)
- Offline Store (хранение признаков для обучения)
- Feature Registry (описания признаков)
- Online Store (подача признаков в сервис инференса)
- Контроль качества и мониторинг
- Эксперименты и репродуктивность
Примеры конфигураций и кода
1) YAML-конфигурация Feast (пример)
entities:
- name: user_id
description: "Уникальный идентификатор пользователя"
features:
- name: user_purchase_count_last_30d
value_type: INT64
input_source: transactions
transform: "sum(purchases) over last 30 days"
ttl: 30d
- name: user_average_session_time
value_type: DOUBLE
input_source: sessions
transform: "avg(session_time)"
ttl: 14d
2) Пример SQL для вычисления признаков (псевдо-цель)
-- Пример: признак среднего времени сессии по пользователю
SELECT
user_id,
AVG(session_time) AS user_avg_session_time
FROM raw_user_sessions
WHERE event_time >= now() - INTERVAL '30 days'
GROUP BY user_id;
3) Пример данных контракта и теста качества (Python + Great Expectations)
# tests/test_feature_quality.py
from great_expectations.dataset import PandasDataset
import pandas as pd
class FeatureDataset(PandasDataset):
@property
def expect_feature_within_range(self):
return self.expect_column_values_to_be_between(
column='user_avg_session_time', min_value=0.0, max_value=3600.0
)
def test_feature(df: pd.DataFrame):
ds = FeatureDataset(df)
assert ds.expect_feature_within_range().success
4) Пример DAG на Dagster (Python)
from dagster import pipeline, solid, ModeDefinition
@solid
def extract_raw(context):
# загрузка данных
return "raw_data"
@solid
def transform_features(context, raw):
# вычисления признаков
features = {"user_avg_session_time": 123.4}
return features
@solid
def publish_offline(context, features):
# сохранить в offline store
pass
@solid
def publish_online(context, features):
# обновить online store
pass
@pipeline(mode_defs=[ModeDefinition()])
def feature_pipeline():
raw = extract_raw()
feats = transform_features(raw)
publish_offline(feats)
publish_online(feats)
5) Контроль качества данных (Great Expectations)
- Схема данных: user_id, feature_name, value, timestamp;
- Требование: отсутствие пропусков в user_id и feature_name; значения в допустимом диапазоне.
Риски и ограничения
- Дрейф признаков (feature drift): характеристики признаков меняются со временем, что ухудшает качество моделей. Решение: периодические проверки, тесты на сценарии регрессии, мониторинг изменений.
- Несоответствие онлайн и оффлайн признаков: неверные версии признаков, задержки в обновлениях, рассинхрон между режимами. Решение: строгие версии, синхронизация времени, откаты к предыдущим версиям.
- Управление версиями: смешение версий признаков может привести к непредсказуемым результатам. Решение: политика версионирования, фиксация зависимостей.
- Производительность и стоимость: частые перерасчёты и обновления признаков требуют вычислительных ресурсов и баланса между частотой обновления и задержкой.
- Безопасность и доступ: ограничение доступа к чувствительным данным, аудит изменений; соответствие требованиям регуляторов (GDPR и пр.).
- Совместимость инструментов: различные версии инструментов могут иметь несовместимости, особенно в гибридных облачных средах.
- Миграции и backfill: обновление существующих признаков может потребовать backfill, что влияет на производительность и может привести к временным ошибкам.
- Репродуктивность экспериментов: если пайплайн не фиксирует версии источников и параметров, воспроизвести результаты может быть сложно.
- Зависимости от облака и vendor lock-in: выбор конкретного облака может ограничить гибкость в будущем.
Рекомендации по снижению рисков:
- Вводить строгие контракты данных и тесты на входные/выходные данные.
- Вести версионирование признаков и пайплайнов.
- Использовать OpenLineage/MLflow для аудита и воспроизводимости экспериментов.
- Настраивать артефакты пайплайна так, чтобы повторная обработка не приводила к конфликтам.
- Реализовать мониторинг производительности пайплайна и задержек.
- Разрабатывать на стеке с поддержкой локальных и облачных инстансов, чтобы легко мигрировать между средами.
- Применять режимы ограниченного доступа и шифрование данных.
Выводы
- Конвейеры подготовки признаков и DAGs играют критическую роль в Lakehouse-архитектуре ML и продвинутой аналитики. Они обеспечивают воспроизводимость, контроль качества, прозрачность lineage и устойчивость к изменениям данных.
- Выбор инструментов зависит от контекста: открытые экосистемы (Airflow, Dagster, Feast, Kedro) дают гибкость и масштабируемость; российские решения (Яндекс DataSphere, СберОблако) — локальную интеграцию, соответствие требованиям локализации и регуляторным ограничениям.
- Важнейшие практики: единая дефиниция признаков в feature registry, версия признаков, интеграция с онлайн и оффлайн складами, контроль качества, мониторинг и аудит, поддержка экспериментов.
- Реализация требовательна к инфраструктуре и организационной культуре: необходимы четкие процессы для тестирования, развёртывания и совместной работы аналитиков и data scientists.
FAQ (Вопрос–Ответ)
1) Что такое конвейер признаков и чем он отличается от обычного ETL?
- Конвейер признаков — это специализированный пайплайн, который не просто переносит данные, а производит признаки через трансформации и агрегации, сохраняет их в оффлайн и онлайн Stores, регистрирует их версии и зависимости. Он ориентирован на повторяемость, качество данных и быструю подачу признаков в модели. В отличие от обычного ETL, конвейер признаков фокусируется на управляемости признаков как артефакт, который подлежит версии и мониторингу.
2) Чем отличается offline store от online store?
- Offline store содержит признаки для обучения и ретро-анализа; он ориентирован на большие объёмы данных и долговременное хранение. Online store — это низколатентный доступ к признакам для инференса в режиме онлайн. Разные слои обеспечивают баланс между скоростью и объёмами данных.
3) Какие инструменты рекомендуется использовать для начинающего проекта?
- Открытая связка: Airflow (оркестрация) + Feast (feature store) + Spark/Delta Lake (обработка и хранение) + Great Expectations (качество). В зависимости от требований можно добавить Dagster/Kedro для лучшей тестируемости и модульности.
4) Как обеспечить качество признаков и защиту от дрейфа?
- Введите контракты данных и тесты на каждую версию признаков; используйте мониторинг изменений признаков и drift-домены; реализуйте регулярные backtests и re-training план. Great Expectations и OpenLineage помогают в этом.
5) Как выбрать между открытой экосистемой и российскими платформами?
- Открытая экосистема даёт широкую совместимость и гибкость, но требует самостоятельной инфраструктуры и SRE-подхода. Российские платформы лучше подходят для локальной инфраструктуры и соответствия требованиям регуляторов, часто имеют готовые интеграции с локальными сервисами и хранилищами.
6) Какие риски связаны с онлайн-признаками?
- Риск помехи в инференсе при задержках обновления признаков, несоответствия между онлайн и оффлайн версиями, проблемы консистентности и устойчивости к дрейфу. Рекомендации: ограничение отклонений, строгие версии, мониторинг latency.
7) Как организовать совместную работу аналитиков и data scientists в конвейерах признаков?
- Введите концепцию feature registry и governance: аналитики определяют признаки и формулы, DS — экспериментируют с моделями и проверяют влияние признаков; инженеры — инфраструктура и эксплуатацию пайплайна. Регулярно проводите ревью признаков и автоматизируйте тесты.
8) Можно ли применить такие конвейеры в реальном времени?
- Да, но это требует эффективного онлайн-store и подхода к стримингу. Часто онлайн-признаки обновляются по графику или в режиме near-real-time, тогда часть вычислений переносится в потоковую обработку (пример — Spark Structured Streaming).
9) Какие ограничения следует учитывать в российских платформах?
- Ограничения обычно касаются экосистемы интеграций, лицензий и стоимости, а также зависимости от конкретной инфраструктуры. В то же время они обеспечивают локальную безопасность, соответствие нормам и поддержку отечественных сервисов.
10) Какие шаги для старта проекта по конвейерам признаков?
- Определить бизнес-задачи и метрики; выбрать стек инструментов (open-source или российскую платформу); спроектировать архитектуру оффлайн/онлайн-store; создать registry признаков; построить начальный DAG; внедрить тесты и мониторинг; запустить пилотный проект и постепенно расширять пайплайн.
Если вы рассматриваете переход к архитектуре Lakehouse, мы поможем оценить текущую data-инфраструктуру, спроектировать целевую архитектуру и подготовить поэтапный план внедрения. Узнайте больше о Lakehouse.



