Data Engineering и ETL Аналитика, метрики и стык с машинным обучением
0%

Аналитика, метрики и стык с машинным обучением

Аналитика, метрики и стык с машинным обучением

Мы прошли весь путь: извлекли и загрузили данные, спроектировали модель, посчитали её батчем и потоком, оркестрировали, положили в нужный формат и обвесили проверками качества. Осталась последняя миля — та, ради которой всё и строилось: данные должны превратиться в решения.

Решения бывают двух видов, и они предъявляют к данным разные, местами противоположные требования. Человек принимает решение, глядя на метрику в дашборде: ему нужны понятность, однозначность и стабильность определений во времени. Модель принимает решение автоматически, миллион раз в сутки: ей нужны признаки, воспроизводимые с точностью до момента времени, и одинаковые в обучении и в проде.

Эта статья — про обе последние мили. Первая половина про то, как метрика становится инженерным объектом, а не строкой SQL в чьём-то блокноте. Вторая — про самый недооценённый интерфейс в индустрии: стык между дата-инженером и ML-инженером, где ломается больше моделей, чем от всех неудачных архитектур нейросетей вместе взятых.

Последняя миля: почему витрина — это ещё не продукт

Типичная зрелая платформа выглядит благополучно: слой marts собран, тесты зелёные, SLA соблюдается. И при этом на совещании три человека называют три разных числа выручки за квартал. Так происходит не потому, что кто-то ошибся в JOIN. Так происходит потому, что определение метрики нигде не живёт как код.

Финансист считает выручку по факту оплаты и без НДС. Продакт — по факту заказа и с НДС. Маркетинг — по атрибуции последнего клика и в валюте кампании. Все три числа корректны внутри своей логики; ни одно не документировано; все три посчитаны копипастой SQL, разошедшейся по десяткам дашбордов. Это классическая проблема расхождения метрик (metric drift в организационном смысле), и она не решается ни хранилищем, ни BI-инструментом. Она решается тем, что определение метрики становится версионируемым артефактом с владельцем и тестами.

Сформулируем правило, которое стоит всей остальной статьи:

Витрина отвечает на вопрос «какие данные есть». Метрика отвечает на вопрос «что мы считаем правдой». Первое — задача моделирования, второе — задача семантики, и путать их дорого.

Анатомия метрики

Метрика — не число и не запрос. Это кортеж из шести элементов, и пропуск любого из них порождает спор.

  1. Мера — что агрегируем: sum(order_amount_net).
  2. Агрегация — как: sum, count_distinct, avg, перцентиль, отношение двух мер.
  3. Зерно события — на каком уровне лежит факт: строка заказа? заказ? платёж? (см. моделирование — зерно факта определяет всё).
  4. Временная привязка — по какой дате метрика ложится в период: дата заказа, дата оплаты, дата отгрузки. Это самый частый источник расхождений.
  5. Фильтры-по-умолчанию — исключаем ли тестовые заказы, внутренних сотрудников, отменённые транзакции, возвраты.
  6. Допустимые разрезы — по каким измерениям метрику вообще законно резать, и какие разрезы бессмысленны.

Шестой пункт неочевиден, поэтому пример. Метрика «средний чек» разрезается по стране, каналу, категории. А вот метрика «доля активных пользователей» не разрезается по заказу — нет такого зерна. Семантический слой обязан уметь запретить бессмысленный разрез, иначе кто-нибудь построит график и будет им управлять.

Аддитивность: три класса мер

Самое важное техническое свойство меры — можно ли складывать её значения вдоль измерения.

Практические следствия, за которые платят инцидентами:

  • Аддитивные меры можно предагрегировать в rollup-таблицу по дням и потом суммировать по неделям — результат точный. Это основа быстрых дашбордов.
  • Полуаддитивные ломаются, если кто-то просуммировал остатки за 30 дней и получил «месячный остаток». Правильная агрегация по времени — last_value (остаток на конец периода) или среднее. В OLAP-движках вроде ClickHouse это делают через argMax(balance, snapshot_date).
  • Неаддитивные нельзя предагрегировать вообще. Конверсия за неделю — это sum(конверсий)/sum(визитов), а не среднее семи дневных конверсий. Разница между этими двумя числами известна как парадокс Симпсона в бытовой форме, и она бывает не косметической: если в один день был всплеск дешёвого трафика с нулевой конверсией, «среднее средних» покажет ровную линию там, где реальная конверсия просела вдвое.

Уникальные пользователи — отдельная боль: count(distinct user_id) за месяц не равен сумме дневных уникальных. Если нужна быстрая аппроксимация, предагрегируют не число, а скетч — HyperLogLog-состояние, которое умеет объединяться (hllMerge в ClickHouse, HLL_COUNT.MERGE в BigQuery, APPROX_COUNT_DISTINCT в Snowflake/Spark). Скетч аддитивен там, где само число неаддитивно. Ошибка при типичной точности — около 1–2%, и это, как правило, дешевле, чем честный distinct по миллиардам строк.

-- BigQuery: аддитивный скетч уникальных пользователей на дневном зерне
create or replace table marts.daily_users_sketch
partition by activity_date as
select
    activity_date,
    country,
    hll_count.init(user_id, 15) as users_hll,   -- 15 → ~0.4% ошибки
    count(*)                    as events
from marts.fct_user_activity
group by 1, 2;

-- корректный месячный уникум: сливаем скетчи, а не суммируем числа
select
    date_trunc(activity_date, month) as month,
    hll_count.merge(users_hll)       as monthly_users
from marts.daily_users_sketch
group by 1;

Семантический слой: метрика как код

Идея семантического (метрик-) слоя проста: определение метрики хранится один раз в декларативном виде, а SQL под конкретный разрез генерируется. BI, ноутбук, API и ML-пайплайн обращаются к одному и тому же определению — расхождение становится технически невозможным.

Исторически это умели «толстые» BI-инструменты (LookML в Looker, Cube.js), сейчас индустрия сдвинулась к headless-подходу: слой отделён от визуализации и доступен всем потребителям. Каноническая реализация в dbt — MetricFlow (документация: https://docs.getdbt.com/docs/build/about-metricflow).

# models/semantic/orders.yml — семантическая модель и метрики поверх факта
semantic_models:
  - name: orders
    model: ref('fct_orders')
    defaults:
      agg_time_dimension: ordered_at      # временная привязка по умолчанию
    entities:
      - name: order_id
        type: primary
      - name: customer_id
        type: foreign
    dimensions:
      - name: ordered_at
        type: time
        type_params: { time_granularity: day }
      - name: paid_at
        type: time
        type_params: { time_granularity: day }
      - name: channel
        type: categorical
      - name: country
        type: categorical
    measures:
      - name: revenue_net
        expr: order_amount_net
        agg: sum
        # тестовые заказы отсекаем ЗДЕСЬ, а не в каждом дашборде
      - name: orders_count
        expr: order_id
        agg: count_distinct

metrics:
  - name: revenue
    label: "Выручка (нетто, по дате заказа)"
    type: simple
    type_params: { measure: revenue_net }
    filter: "{{ Dimension('orders__is_test') }} = false"

  - name: revenue_recognized
    label: "Выручка признанная (по дате оплаты)"
    type: simple
    type_params: { measure: revenue_net }
    # тот же измеритель, другая временная привязка — и это ДРУГАЯ метрика
    config: { meta: { agg_time_dimension: paid_at } }

  - name: aov
    label: "Средний чек"
    type: ratio                     # отношение, а не avg — неаддитивность учтена
    type_params:
      numerator: revenue
      denominator: orders_count

Ключевой выигрыш — в третьей метрике. aov объявлен как ratio, поэтому при разрезе по неделям движок посчитает sum(revenue)/sum(orders), а не среднее дневных чеков. Ошибку «среднее средних» невозможно совершить, потому что её негде совершить.

Пунктирная стрелка — то, что происходит в реальности в первый же спринт после внедрения: кто-то торопится и пишет прямой SQL в дашборде. Это не катастрофа, но за долей таких обходов надо следить: если 40% дашбордов ходят мимо слоя, слоя фактически нет.

Trade-off честно. Семантический слой добавляет уровень косвенности: сгенерированный SQL сложнее отлаживать, редкие сложные запросы (оконные функции, самосоединения, хитрые полуаддитивные логики) в декларативный YAML ложатся плохо, и появляется соблазн описывать метрики «почти-SQL»-строками, теряя половину выгоды. Разумная стратегия — покрывать слоем 20–30 метрик, которые реально используются в управлении, и не пытаться загнать туда весь ad-hoc.

Когорты и удержание: где аналитика уже почти инженерия

Кросс-секционная метрика («сколько активных пользователей в марте») почти всегда обманывает, потому что смешивает разнородные группы. Когортный анализ фиксирует пользователей по моменту входа и наблюдает за группой во времени.

Когортный треугольник удержания

Треугольная форма — не украшение, а суть. Правый край каждой строки обрезан, потому что для молодых когорт данных о зрелом возрасте физически ещё нет. Главная ошибка при построении такой таблицы — не обрезать незавершённые периоды: если сегодня 15 марта, а вы считаете «удержание M0 для мартовской когорты», вы сравниваете половину месяца с полными месяцами и получаете фальшивое падение. Это тот же класс ошибки, что и незакрытая партиция в батче.

-- Когортное удержание по месяцам. Работает в любом ANSI-SQL движке.
with first_activity as (
    select
        user_id,
        date_trunc('month', min(activity_ts))            as cohort_month
    from marts.fct_user_activity
    group by user_id
),
monthly as (
    select distinct
        user_id,
        date_trunc('month', activity_ts)                 as active_month
    from marts.fct_user_activity
),
joined as (
    select
        f.cohort_month,
        -- возраст когорты в месяцах: 0 = месяц регистрации
        (extract(year from m.active_month) - extract(year from f.cohort_month)) * 12
        + (extract(month from m.active_month) - extract(month from f.cohort_month)) as age_months,
        m.user_id
    from first_activity f
    join monthly m using (user_id)
),
sizes as (
    select cohort_month, count(*) as cohort_size
    from first_activity group by cohort_month
)
select
    j.cohort_month,
    j.age_months,
    s.cohort_size,
    count(distinct j.user_id)                                       as retained,
    round(100.0 * count(distinct j.user_id) / s.cohort_size, 1)     as retention_pct
from joined j
join sizes s using (cohort_month)
-- ОБРЕЗКА: показываем только те ячейки, чей период полностью завершён
where j.cohort_month + (j.age_months + 1) * interval '1 month' <= date_trunc('month', current_date)
group by 1, 2, 3
order by 1, 2;

Сложность. Тяжёлая часть — count(distinct) по большой таблице активности: в MPP-движке это shuffle по user_id, то есть O(N) сетевого трафика на каждый запуск. При десятках миллиардов строк дневной активности такую таблицу материализуют инкрементально: считают агрегат на уровне (user_id, month) один раз, а треугольник строят уже поверх него — это превращает распределённый distinct над сырьём в дешёвый join над сжатой таблицей размера «пользователи × активные месяцы».

Как читать результат: по строке — форма кривой удержания (интересует не M1, а выход на плато; продукт без плато не удерживает никого); по столбцу — улучшается ли онбординг от когорты к когорте; по диагонали — общие внешние события (диагональ = один календарный месяц, и если она просела целиком, виноват не онбординг, а инцидент, сезон или релиз).

Смена контракта: аналитика против ML

Дальше начинается вторая половина последней мили. Данные те же, потребитель другой — и требования сдвигаются радикально.

Аспект Аналитическое потребление ML-потребление
Читатель человек, глазами код, автоматически
Частота десятки-сотни запросов в день обучение — редко; инференс — тысячи запросов в секунду
Латентность секунды приемлемы инференс: единицы-десятки миллисекунд
Отношение ко времени «как сейчас» или «за период» строго «как было известно в момент T»
Реакция на пересчёт истории обычно благо: цифры стали точнее часто катастрофа: обучающая выборка перестала соответствовать проду
Реакция на пропуски видно глазом, аналитик спросит молча заполнится дефолтом и сместит предсказание
Цена ошибки в определении неверный слайд неверные решения в проде месяцами, пока не заметят

Строка про пересчёт истории — самая недооценённая. Дата-инженер привык, что backfill улучшает данные. Для ML backfill означает, что признак, на котором обучалась модель в проде, задним числом изменил значения, и воспроизвести обучение больше нельзя. Отсюда правило: таблицы, из которых берутся признаки, должны быть либо иммутабельными, либо версионированными — snapshot-таблицы, SCD2 или time travel в Iceberg/Delta.

Point-in-time correctness: главная ошибка на этом стыке

Утечка из будущего (target leakage) — ситуация, когда в обучающую выборку попадает информация, недоступная в момент, когда модель реально должна принять решение. Это ошибка данных, а не модели, и потому её обычно допускает дата-инженер, а расплачивается за неё дата-сайентист.

Point-in-time join и утечка из будущего

Канонический сюжет: скоринг заявок. Признак avg_balance_30d берут из витрины dim_customer_current, где лежит текущее значение. Для клиента, который взял кредит в марте и ушёл в дефолт в мае, текущий баланс — почти ноль. Модель мгновенно «понимает», что низкий баланс = дефолт, показывает ROC-AUC 0.97 на валидации и разваливается в проде, потому что на момент заявки баланс был нормальным.

Признаки утечки, которые стоит воспринимать как красный флаг:

  • метрика качества подозрительно высока (AUC > 0.95 на реальной бизнес-задаче — почти всегда утечка);
  • один признак даёт почти всю важность в модели;
  • офлайн-качество катастрофически расходится с онлайн;
  • признак построен из таблицы, которая обновляется merge/overwrite без истории.

Лечится это единственным способом — ASOF-join (point-in-time join): для каждой строки обучающей выборки берём то значение признака, у которого feature_ts <= event_ts, и максимальное среди таких.

-- Вариант 1: движки с нативным ASOF JOIN (ClickHouse, DuckDB, Snowflake)
select
    l.event_id,
    l.customer_id,
    l.event_ts,
    l.target_default_90d,
    f.avg_balance_30d,
    f.tx_count_7d,
    f.feature_ts                    -- ОБЯЗАТЕЛЬНО тащим наружу для аудита
from ml.labels          as l
asof join ml.customer_features as f
    on l.customer_id = f.customer_id
   and l.event_ts    >= f.feature_ts;   -- строго "не позже момента события"
-- Вариант 2: портируемо, через оконную функцию. Работает везде, дороже.
with candidates as (
    select
        l.event_id,
        l.event_ts,
        l.target_default_90d,
        f.avg_balance_30d,
        f.tx_count_7d,
        f.feature_ts,
        row_number() over (
            partition by l.event_id
            order by f.feature_ts desc          -- ближайший снимок ИЗ ПРОШЛОГО
        ) as rn
    from ml.labels as l
    join ml.customer_features as f
      on f.customer_id = l.customer_id
     and f.feature_ts <= l.event_ts
     -- ограничение снизу спасает от разрастания join'а и от протухших признаков
     and f.feature_ts >= l.event_ts - interval '90 day'
)
select event_id, event_ts, target_default_90d,
       avg_balance_30d, tx_count_7d, feature_ts
from candidates
where rn = 1;

Сложность и цена. Наивный вариант через join + row_number материализует декартово произведение «метка × все её прошлые снимки признаков» до фильтрации: при L метках и в среднем K снимков на сущность в окне это O(L·K) строк и сортировка O(L·K·log K). Нативный ASOF-join делает merge-join по отсортированным потокам — O((L + F)·log(L + F)) на сортировку и линейный проход, память — O(1) на сущность. Разница на практике: часы против минут. Если ASOF недоступен, обязательно ограничивайте окно снизу (>= event_ts - interval) и партиционируйте обе стороны по сущности.

На малых данных то же самое в pandas — и здесь ловушка в том, что merge_asof требует глобальной сортировки по ключу времени, а не по группам:

import pandas as pd

# ВАЖНО: merge_asof требует сортировки по столбцу времени по всему датафрейму
labels = labels.sort_values("event_ts")
features = features.sort_values("feature_ts")

train = pd.merge_asof(
    labels,                      # левая сторона: что предсказываем и когда
    features,                    # правая: история значений признаков
    left_on="event_ts",
    right_on="feature_ts",
    by="customer_id",            # склейка внутри сущности
    direction="backward",        # только прошлое; forward = гарантированная утечка
    tolerance=pd.Timedelta("90d"),  # снимок старше 90 дней считаем протухшим
)

# Инвариант, который должен стоять в тестах пайплайна, а не в голове:
assert (train["feature_ts"] <= train["event_ts"]).all(), "утечка из будущего"
# Доля строк без признаков — тоже метрика качества выборки, а не мелочь
coverage = train["avg_balance_30d"].notna().mean()
assert coverage > 0.8, f"покрытие признаками всего {coverage:.1%}"

Отдельно про метки. Целевая переменная тоже имеет временную геометрию: «дефолт в течение 90 дней» означает, что метку для события от 1 марта нельзя считать раньше 30 мая. Наблюдения из последних 90 дней должны быть исключены из обучения — иначе часть положительных примеров ещё не успела проявиться, и модель обучится на систематически заниженном уровне позитивов. Этот зазор называют label maturity window, и его забывают примерно так же часто, как point-in-time.

Feature store: зачем нужен отдельный компонент

Когда признаки нужны и в обучении (большими батчами, за историю), и в инференсе (по одному ключу, за 10 мс), возникает соблазн написать два независимых куска кода. Именно так рождается training/serving skew — расхождение между тем, как признак считался при обучении, и тем, как он считается в проде. В офлайне avg_balance_30d посчитан по календарным дням в UTC из витрины, в онлайне — по скользящему окну от текущего момента в локальном времени из Redis. Формально «то же самое». Фактически — другое распределение, и модель деградирует тихо.

Feature store решает ровно одну задачу: одно определение признака — два способа материализации.

Минимальный пример на Feast (открытый feature store, документация: https://docs.feast.dev/):

from datetime import timedelta
from feast import Entity, FeatureView, Field, FileSource
from feast.types import Float32, Int64

customer = Entity(name="customer", join_keys=["customer_id"])

# Источник обязан нести две временные колонки:
# event_timestamp — когда значение стало ВЕРНЫМ (используется в point-in-time join)
# created_timestamp — когда строка была ЗАПИСАНА (для дедупликации поздних данных)
source = FileSource(
    path="s3://lake/features/customer_stats/",
    timestamp_field="event_timestamp",
    created_timestamp_column="created_timestamp",
)

customer_stats = FeatureView(
    name="customer_stats",
    entities=[customer],
    ttl=timedelta(days=90),        # снимки старше 90 дней не подставляются
    schema=[
        Field(name="avg_balance_30d", dtype=Float32),
        Field(name="tx_count_7d",     dtype=Int64),
        Field(name="days_since_signup", dtype=Int64),
    ],
    source=source,
    online=True,                    # материализуется и в онлайн-хранилище
)
# Обучение: point-in-time join делает сам стор — руками ASOF писать не нужно
training_df = store.get_historical_features(
    entity_df=labels_df,            # обязан содержать customer_id и event_timestamp
    features=[
        "customer_stats:avg_balance_30d",
        "customer_stats:tx_count_7d",
        "customer_stats:days_since_signup",
    ],
).to_df()

# Инференс: тот же список признаков, другой метод — и это весь смысл стора
vector = store.get_online_features(
    features=[
        "customer_stats:avg_balance_30d",
        "customer_stats:tx_count_7d",
        "customer_stats:days_since_signup",
    ],
    entity_rows=[{"customer_id": 42}],
).to_dict()

Когда feature store не нужен. Честный ответ: чаще, чем его внедряют. Если у вас одна-две батчевые модели, скорящие раз в сутки, и признаки берутся из тех же витрин, что и дашборды, — достаточно дисциплины: снапшоты с историей плюс один общий модуль расчёта признаков. Feature store окупается, когда есть (а) реальный онлайн-инференс с жёсткой латентностью, (б) переиспользование признаков между несколькими командами и моделями, (в) требование аудита «на каких данных обучалась модель, отказавшая клиенту в кредите». Вне этих условий вы получаете ещё одну распределённую систему с собственным дежурством.

Trade-off онлайн-хранилища. Redis даёт единицы миллисекунд, но это ещё одна копия данных, которую надо консистентно обновлять; расхождение между офлайном и онлайном (stale features) — самый частый инцидент feature store. Обязательно мониторьте freshness lag: возраст самого свежего значения в онлайн-сторе. Если материализация упала ночью, модель продолжит отвечать — просто вчерашними числами, и никакой алерт по HTTP-кодам этого не заметит.

Жизненный цикл признака

Признак — не строка кода, а объект со своим жизненным циклом и владельцем. Полезно явно описать состояния, иначе в сторе через год накапливаются сотни признаков, из которых половину никто не использует, но все материализуются и стоят денег.

Практическая деталь: переход «Активен → Устарел» невозможно определить без линиджа потребления — надо знать, какие модели читают какой признак. Это ровно тот же механизм, что и линидж таблиц из статьи про качество и governance, просто продлённый на один шаг вправо, за границу хранилища.

Мониторинг: данные ломаются тише, чем код

Отказавший сервис виден сразу. Модель, которой подсунули сместившиеся признаки, продолжает отвечать 200 OK и выдавать правдоподобные числа — просто хуже. Поэтому мониторинг на этом стыке трёхслойный:

  1. Качество данных на входе — те же проверки, что и везде: свежесть, объём, доля null, диапазоны, уникальность ключей.
  2. Дрифт распределений — сравнение распределения признака в проде с эталонным (обучающим). Стандартная метрика — PSI (Population Stability Index).
  3. Качество модели — как только приезжают метки: precision/recall, калибровка, бизнес-метрика. Отложено на длину label maturity window, поэтому первые два слоя и нужны: они дают сигнал раньше.
import numpy as np

def psi(expected: np.ndarray, actual: np.ndarray, bins: int = 10) -> float:
    """Population Stability Index — насколько распределение уехало от эталона.

    Интерпретация (индустриальная эвристика из кредитного скоринга):
        < 0.10  — сдвига нет
        0.10..0.25 — умеренный сдвиг, стоит разобраться
        > 0.25  — существенный сдвиг, кандидат на переобучение

    Сложность: O(n log n) на квантили + O(n) на раскладку по корзинам.
    """
    # Границы корзин берём ПО ЭТАЛОНУ: иначе сравниваем разные шкалы
    edges = np.quantile(expected, np.linspace(0, 1, bins + 1))
    edges[0], edges[-1] = -np.inf, np.inf
    edges = np.unique(edges)                      # схлопываем дубли на константах

    exp_pct = np.histogram(expected, bins=edges)[0] / len(expected)
    act_pct = np.histogram(actual,   bins=edges)[0] / len(actual)

    # Сглаживание: пустая корзина даёт log(0) и бесконечный PSI
    eps = 1e-6
    exp_pct = np.clip(exp_pct, eps, None)
    act_pct = np.clip(act_pct, eps, None)

    return float(np.sum((act_pct - exp_pct) * np.log(act_pct / exp_pct)))


# Ежедневная проверка всех признаков модели
for name in feature_names:
    score = psi(train_ref[name].to_numpy(), prod_today[name].to_numpy())
    if score > 0.25:
        alert(f"PSI по признаку {name}: {score:.3f} — вероятен сдвиг источника")

Важная оговорка про интерпретацию: PSI детектирует сдвиг, а не поломку. Скачок PSI по признаку «регион» после выхода на новый рынок — ожидаемое поведение бизнеса, а не инцидент. А вот скачок PSI одновременно по десятку не связанных признаков почти всегда означает не смену поведения пользователей, а техническую причину: сменился парсер, поехал часовой пояс, апстрим переименовал enum, партиция долилась не полностью. Именно поэтому дрифт-мониторинг должен смотреться рядом с линиджем — иначе разбор каждого алерта превращается в расследование с нуля.

Для отслеживания дрифта категориальных признаков PSI применим напрямую (корзины = категории), для непрерывных иногда предпочитают критерий Колмогорова–Смирнова или расстояние Вассерштейна — они менее чувствительны к выбору корзин. Готовые реализации: библиотека Evidently (https://docs.evidentlyai.com/) и NannyML.

Обратный поток: предсказания возвращаются в хранилище

Стык двусторонний. Модель не только читает данные — она их производит, и эти предсказания должны вернуться в аналитический контур, иначе оценить эффект модели невозможно.

Минимальный контракт таблицы предсказаний:

create table ml.predictions_credit_scoring (
    prediction_id   uuid          not null,
    entity_id       bigint        not null,   -- кому/чему предсказали
    predicted_at    timestamptz   not null,   -- когда (время инференса)
    model_name      text          not null,
    model_version   text          not null,   -- ОБЯЗАТЕЛЬНО: без версии анализ невозможен
    features_hash   text          not null,   -- хэш вектора признаков для воспроизведения
    score           double precision not null,
    decision        text          not null,   -- что решил бизнес-слой поверх score
    threshold       double precision not null,-- порог на момент решения (он меняется!)
    experiment_arm  text,                     -- ветка A/B, если модель катится экспериментом
    primary key (prediction_id)
) partition by range (predicted_at);

Три поля здесь неочевидны и все три критичны. model_version — потому что без неё нельзя разделить «модель стала хуже» и «выкатили другую модель». threshold — потому что порог отсечения меняют вручную чаще, чем переобучают модель, и без записи порога распределение решений выглядит необъяснимым. experiment_arm — потому что честно измерить эффект модели можно только сравнением с контрольной группой: сравнение «до/после» смешивает эффект модели с сезонностью и всеми параллельными релизами.

Дальше эта таблица джойнится с фактическими исходами (когда те созреют) и даёт две вещи: онлайн-метрики качества модели и бизнес-эффект — то единственное число, ради которого модель строилась. Отдельный поток — reverse ETL: доставка предсказаний обратно в операционные системы (CRM, рассылки, интерфейс оператора). Технически это обычный ELT наоборот, но с двумя специфическими требованиями: идемпотентность по entity_id (перезапуск не должен слать письмо дважды) и защита от «зацикливания» — если предсказание модели попадёт в признаки той же модели без разделения, вы получите петлю обратной связи, в которой модель обучается на собственных решениях и уверенно уезжает от реальности.

Организационный стык: где проходит граница

Технические механизмы работают только поверх ясной границы ответственности. Рабочее разделение, которое чаще всего выживает:

  • Дата-инженер отвечает за то, что таблицы-источники признаков существуют, свежие, иммутабельные или версионированные, имеют контракт и SLA. За event_timestamp. За то, что backfill не ломает воспроизводимость.
  • ML-инженер / DS отвечает за содержание признака (какая логика), за валидацию модели, за то, что офлайн-расчёт и онлайн-расчёт идут из одного определения.
  • Совместно — за реестр признаков, мониторинг дрифта и процедуру депрекации.

Формализуется это тем же контрактом данных, что и в governance, только с ML-специфичными полями:

# contracts/customer_stats.yml — контракт на источник признаков
dataset: lake.features.customer_stats
owner: team-risk-data
consumers: [model.credit_scoring.v3, model.churn.v1, dashboard.risk_overview]

schema:
  - { name: customer_id,      type: bigint,    nullable: false }
  - { name: event_timestamp,  type: timestamp, nullable: false }
  - { name: avg_balance_30d,  type: double,    nullable: true  }
  - { name: tx_count_7d,      type: bigint,    nullable: false }

guarantees:
  freshness_sla: 6h            # максимальный возраст самого свежего снимка
  immutable_history: true      # прошлые снимки НЕ переписываются — иначе нет воспроизводимости
  late_arrival_window: 48h     # позже — приезжает как новый снимок, а не правка старого
  backfill_policy: append_only # backfill добавляет строки, не апдейтит

ml_specific:
  point_in_time_column: event_timestamp
  max_staleness_for_training: 90d   # снимок старше — не подставляется в обучение
  online_materialization: true
  online_freshness_alert: 2h

Типичные ошибки

Собрал те, которые встречаются чаще всего, — от аналитической половины к ML-половине.

  1. Среднее средних. Недельная конверсия как avg дневных. Лечится объявлением метрики как ratio в семантическом слое.
  2. Суммирование полуаддитивных мер по времени. «Остаток за месяц» = сумма дневных остатков. Число получается абсурдным, но график выглядит гладко.
  3. Одно имя — два определения. «Активный пользователь» у продакта и у маркетинга. Метрика без владельца и без кода обречена размножиться.
  4. Незакрытые периоды в когортах и трендах. Последняя точка графика всегда ниже — и каждый раз кто-то бьёт тревогу.
  5. Дашборд поверх сырья мимо семантического слоя. Работает, пока определение не изменится; после — тихо расходится.
  6. Утечка из будущего. Признак из snapshot-таблицы «как сейчас» вместо ASOF-join. Симптом — неправдоподобно высокое офлайн-качество.
  7. Забытое окно созревания метки. Свежие наблюдения без успевших проявиться событий смещают baseline.
  8. Два кода расчёта признака — офлайн и онлайн. Расхождение неизбежно и обнаруживается через месяцы.
  9. Mutable источник признаков. Backfill переписал историю — воспроизвести обучение и объяснить старое решение невозможно (а для регулируемых доменов это ещё и юридическая проблема).
  10. Тихое заполнение пропусков. Признак не пришёл → подставился ноль → модель решила, что у клиента нулевой баланс. Пропуск обязан быть явным (is_missing-флаг) и отслеживаемым метрикой покрытия.
  11. Предсказания без версии модели и порога. Любой ретроспективный анализ качества становится гаданием.
  12. Отсутствие контрольной группы. Эффект модели меряют «до/после» и приписывают ей сезонный рост.

Мини-итог

  • Метрика — инженерный объект: мера, агрегация, зерно, временная привязка, фильтры, допустимые разрезы. Всё это должно жить в коде, а не в устной традиции.
  • Аддитивность решает, что можно предагрегировать. Неаддитивные меры (отношения, перцентили, уникумы) считаются из фактов или через сливаемые скетчи.
  • Семантический слой делает расхождение метрик технически невозможным; платите за это косвенностью и ограниченной выразительностью.
  • Когорты показывают то, что кросс-секционные метрики прячут; главное — обрезать незавершённые периоды.
  • ML меняет контракт к данным: не «как сейчас», а «как было известно в момент T». Отсюда point-in-time join, иммутабельная история и окно созревания меток.
  • Feature store оправдан при онлайн-инференсе, переиспользовании признаков и требовании аудита; в остальных случаях достаточно дисциплины и одного общего модуля расчёта.
  • Мониторьте вход (качество), середину (дрифт, freshness) и выход (качество модели и бизнес-эффект). Предсказания возвращайте в хранилище с версией модели, порогом и веткой эксперимента.

Источники

Что дальше

Это последняя статья трека Data Engineering. Данные доехали до дашборда и до модели — дальше начинаются соседние дисциплины, и каждая из них продолжает ровно ту линию, которую мы здесь оборвали.

  • Machine Learning — что происходит с признаками после того, как вы их корректно отдали: постановка задачи, валидация, метрики качества моделей.
  • Нейронные сети — отдельный класс потребителей данных с собственными требованиями к объёму и препроцессингу.
  • Архитектурные паттерны — как дата-платформа встраивается в общую архитектуру системы: события, интеграции, границы сервисов.
  • DDD — язык, на котором договариваются о смысле сущностей; половина споров про определение метрики — это на самом деле спор про ubiquitous language.
  • Product Management — сторона, которая эти метрики заказывает: продуктовые метрики, эксперименты, приоритизация.
  • Алгоритмы и Структуры данных — фундамент под всем, что мы обсуждали: сортировки и хэши в основе join’ов, вероятностные структуры в основе быстрых уникумов.
  • Языковые треки, если хочется писать пайплайны руками: Go, Elixir, TypeScript, C#.

Общая карта портала и рекомендованный порядок изучения — в роадмапе. Начать трек заново или свериться с картой можно с обзорной статьи.

Нашли неточность? Выделите фрагмент текста — рядом появится жучок.

Нужен разбор именно вашей ситуации?

Статья описывает общий случай. Если у вас частный — можно разобрать его отдельно, платно. А если не хватает целого материала, предложите тему: её оплачивают вскладчину, и она выходит открытой для всех.

Доска запросов