Аналитика, метрики и стык с машинным обучением
Мы прошли весь путь: извлекли и загрузили данные, спроектировали модель, посчитали её батчем и потоком, оркестрировали, положили в нужный формат и обвесили проверками качества. Осталась последняя миля — та, ради которой всё и строилось: данные должны превратиться в решения.
Решения бывают двух видов, и они предъявляют к данным разные, местами противоположные требования. Человек принимает решение, глядя на метрику в дашборде: ему нужны понятность, однозначность и стабильность определений во времени. Модель принимает решение автоматически, миллион раз в сутки: ей нужны признаки, воспроизводимые с точностью до момента времени, и одинаковые в обучении и в проде.
Эта статья — про обе последние мили. Первая половина про то, как метрика становится инженерным объектом, а не строкой SQL в чьём-то блокноте. Вторая — про самый недооценённый интерфейс в индустрии: стык между дата-инженером и ML-инженером, где ломается больше моделей, чем от всех неудачных архитектур нейросетей вместе взятых.
Последняя миля: почему витрина — это ещё не продукт
Типичная зрелая платформа выглядит благополучно: слой marts собран, тесты зелёные, SLA соблюдается. И при этом на совещании три человека называют три разных числа выручки за квартал. Так происходит не потому, что кто-то ошибся в JOIN. Так происходит потому, что определение метрики нигде не живёт как код.
Финансист считает выручку по факту оплаты и без НДС. Продакт — по факту заказа и с НДС. Маркетинг — по атрибуции последнего клика и в валюте кампании. Все три числа корректны внутри своей логики; ни одно не документировано; все три посчитаны копипастой SQL, разошедшейся по десяткам дашбордов. Это классическая проблема расхождения метрик (metric drift в организационном смысле), и она не решается ни хранилищем, ни BI-инструментом. Она решается тем, что определение метрики становится версионируемым артефактом с владельцем и тестами.
Сформулируем правило, которое стоит всей остальной статьи:
Витрина отвечает на вопрос «какие данные есть». Метрика отвечает на вопрос «что мы считаем правдой». Первое — задача моделирования, второе — задача семантики, и путать их дорого.
Анатомия метрики
Метрика — не число и не запрос. Это кортеж из шести элементов, и пропуск любого из них порождает спор.
- Мера — что агрегируем:
sum(order_amount_net). - Агрегация — как:
sum,count_distinct,avg, перцентиль, отношение двух мер. - Зерно события — на каком уровне лежит факт: строка заказа? заказ? платёж? (см. моделирование — зерно факта определяет всё).
- Временная привязка — по какой дате метрика ложится в период: дата заказа, дата оплаты, дата отгрузки. Это самый частый источник расхождений.
- Фильтры-по-умолчанию — исключаем ли тестовые заказы, внутренних сотрудников, отменённые транзакции, возвраты.
- Допустимые разрезы — по каким измерениям метрику вообще законно резать, и какие разрезы бессмысленны.
Шестой пункт неочевиден, поэтому пример. Метрика «средний чек» разрезается по стране, каналу, категории. А вот метрика «доля активных пользователей» не разрезается по заказу — нет такого зерна. Семантический слой обязан уметь запретить бессмысленный разрез, иначе кто-нибудь построит график и будет им управлять.
Аддитивность: три класса мер
Самое важное техническое свойство меры — можно ли складывать её значения вдоль измерения.
Практические следствия, за которые платят инцидентами:
- Аддитивные меры можно предагрегировать в 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), а не среднее дневных чеков. Ошибку «среднее средних» невозможно совершить, потому что её негде совершить.
сырьё 1:1"] --> core["core / dwh
факты и измерения"] core --> marts["marts
витрины"] end marts --> sem["Семантический слой
метрики, измерения,
аддитивность, фильтры"] sem --> bi["BI-дашборды"] sem --> nb["Ноутбуки / ad-hoc"] sem --> api["Metrics API
для продуктовых сервисов"] sem --> exp["Платформа экспериментов
A/B, guardrail-метрики"] sem --> feat["Feature engineering
признаки для ML"] marts -.->|"антипаттерн:
SQL мимо слоя"| bi classDef good fill:#10b98122,stroke:#10b981,stroke-width:2px classDef bad fill:#ef444422,stroke:#ef4444,stroke-width:2px,stroke-dasharray:5 4 class sem good class marts bad
Пунктирная стрелка — то, что происходит в реальности в первый же спринт после внедрения: кто-то торопится и пишет прямой 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) — ситуация, когда в обучающую выборку попадает информация, недоступная в момент, когда модель реально должна принять решение. Это ошибка данных, а не модели, и потому её обычно допускает дата-инженер, а расплачивается за неё дата-сайентист.
Канонический сюжет: скоринг заявок. Признак 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 решает ровно одну задачу: одно определение признака — два способа материализации.
(Iceberg / warehouse) participant FS as Feature store
(определения) participant ON as Online store
(Redis / DynamoDB) participant TR as Обучение participant SVC as Сервис инференса Note over DE,ON: Материализация — одно определение, две цели DE->>OFF: запись истории признаков (feature_ts, value) FS->>OFF: materialize(history) — вся история FS->>ON: materialize(latest) — только последний срез + TTL Note over TR,OFF: Обучение: точка во времени TR->>FS: get_historical_features(entity_df с event_ts) FS->>OFF: ASOF join по feature_ts <= event_ts OFF-->>TR: обучающая выборка без утечек Note over SVC,ON: Инференс: точка "сейчас" SVC->>FS: get_online_features(customer_id=42) FS->>ON: GET по ключу ON-->>SVC: вектор признаков за ~5 мс SVC-->>SVC: predict() Note over OFF,ON: Инвариант: обе ветки читают ОДНО определение
расхождение → training/serving skew
Минимальный пример на 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-кодам этого не заметит.
Жизненный цикл признака
Признак — не строка кода, а объект со своим жизненным циклом и владельцем. Полезно явно описать состояния, иначе в сторе через год накапливаются сотни признаков, из которых половину никто не использует, но все материализуются и стоят денег.
есть определение и тесты Валидируется --> Отклонён: нет прироста качества
или невозможен онлайн Валидируется --> Активен: прошёл офлайн-валидацию
и point-in-time проверку Активен --> Активен: материализация по расписанию,
мониторинг дрифта и freshness Активен --> Деградирует: PSI выше порога /
источник сменил семантику Деградирует --> Активен: пересчёт, фикс источника,
переобучение модели Деградирует --> Устарел: источник умер или
признак заменён новым Активен --> Устарел: ни одна активная модель
не использует признак Устарел --> [*]: материализация выключена,
определение оставлено для аудита Отклонён --> [*] note right of Устарел Признак нельзя удалить сразу: по нему обучались модели, которые ещё могут откатить в прод. Deprecation window — минимум один цикл переобучения. end note
Практическая деталь: переход «Активен → Устарел» невозможно определить без линиджа потребления — надо знать, какие модели читают какой признак. Это ровно тот же механизм, что и линидж таблиц из статьи про качество и governance, просто продлённый на один шаг вправо, за границу хранилища.
Мониторинг: данные ломаются тише, чем код
Отказавший сервис виден сразу. Модель, которой подсунули сместившиеся признаки, продолжает отвечать 200 OK и выдавать правдоподобные числа — просто хуже. Поэтому мониторинг на этом стыке трёхслойный:
- Качество данных на входе — те же проверки, что и везде: свежесть, объём, доля null, диапазоны, уникальность ключей.
- Дрифт распределений — сравнение распределения признака в проде с эталонным (обучающим). Стандартная метрика — PSI (Population Stability Index).
- Качество модели — как только приезжают метки: 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-половине.
- Среднее средних. Недельная конверсия как
avgдневных. Лечится объявлением метрики как ratio в семантическом слое. - Суммирование полуаддитивных мер по времени. «Остаток за месяц» = сумма дневных остатков. Число получается абсурдным, но график выглядит гладко.
- Одно имя — два определения. «Активный пользователь» у продакта и у маркетинга. Метрика без владельца и без кода обречена размножиться.
- Незакрытые периоды в когортах и трендах. Последняя точка графика всегда ниже — и каждый раз кто-то бьёт тревогу.
- Дашборд поверх сырья мимо семантического слоя. Работает, пока определение не изменится; после — тихо расходится.
- Утечка из будущего. Признак из snapshot-таблицы «как сейчас» вместо ASOF-join. Симптом — неправдоподобно высокое офлайн-качество.
- Забытое окно созревания метки. Свежие наблюдения без успевших проявиться событий смещают baseline.
- Два кода расчёта признака — офлайн и онлайн. Расхождение неизбежно и обнаруживается через месяцы.
- Mutable источник признаков. Backfill переписал историю — воспроизвести обучение и объяснить старое решение невозможно (а для регулируемых доменов это ещё и юридическая проблема).
- Тихое заполнение пропусков. Признак не пришёл → подставился ноль → модель решила, что у клиента нулевой баланс. Пропуск обязан быть явным (
is_missing-флаг) и отслеживаемым метрикой покрытия. - Предсказания без версии модели и порога. Любой ретроспективный анализ качества становится гаданием.
- Отсутствие контрольной группы. Эффект модели меряют «до/после» и приписывают ей сезонный рост.
Мини-итог
- Метрика — инженерный объект: мера, агрегация, зерно, временная привязка, фильтры, допустимые разрезы. Всё это должно жить в коде, а не в устной традиции.
- Аддитивность решает, что можно предагрегировать. Неаддитивные меры (отношения, перцентили, уникумы) считаются из фактов или через сливаемые скетчи.
- Семантический слой делает расхождение метрик технически невозможным; платите за это косвенностью и ограниченной выразительностью.
- Когорты показывают то, что кросс-секционные метрики прячут; главное — обрезать незавершённые периоды.
- ML меняет контракт к данным: не «как сейчас», а «как было известно в момент T». Отсюда point-in-time join, иммутабельная история и окно созревания меток.
- Feature store оправдан при онлайн-инференсе, переиспользовании признаков и требовании аудита; в остальных случаях достаточно дисциплины и одного общего модуля расчёта.
- Мониторьте вход (качество), середину (дрифт, freshness) и выход (качество модели и бизнес-эффект). Предсказания возвращайте в хранилище с версией модели, порогом и веткой эксперимента.
Источники
- Ralph Kimball, Margy Ross. The Data Warehouse Toolkit, 3rd ed. — аддитивность мер, зерно факта, когортные конструкции: https://www.kimballgroup.com/data-warehouse-business-intelligence-resources/books/data-warehouse-dw-toolkit/
- dbt MetricFlow — устройство семантического слоя и типы метрик: https://docs.getdbt.com/docs/build/about-metricflow
- Cube — headless BI и метрик-API: https://cube.dev/docs/product/introduction
- Feast — открытый feature store, point-in-time join и материализация: https://docs.feast.dev/
- D. Sculley et al. Hidden Technical Debt in Machine Learning Systems, NeurIPS 2015 — почему код модели это малая часть системы: https://papers.nips.cc/paper_files/paper/2015/hash/86df7dcfd896fcaf2674f757a2463eba-Abstract.html
- Google. Rules of Machine Learning: Best Practices for ML Engineering — Правило 29 про training/serving skew: https://developers.google.com/machine-learning/guides/rules-of-ml
- Chip Huyen. Designing Machine Learning Systems, O’Reilly, 2022 — главы про feature engineering и мониторинг: https://huyenchip.com/books/
- Evidently — практики мониторинга дрифта данных и качества моделей: https://docs.evidentlyai.com/
- Kaufman, Rosset, Perlich. Leakage in Data Mining: Formulation, Detection, and Avoidance — строгая формализация утечки: https://dl.acm.org/doi/10.1145/2020408.2020496
- Ron Kohavi, Diane Tang, Ya Xu. Trustworthy Online Controlled Experiments — guardrail-метрики и корректное измерение эффекта: https://experimentguide.com/
Что дальше
Это последняя статья трека Data Engineering. Данные доехали до дашборда и до модели — дальше начинаются соседние дисциплины, и каждая из них продолжает ровно ту линию, которую мы здесь оборвали.
- Machine Learning — что происходит с признаками после того, как вы их корректно отдали: постановка задачи, валидация, метрики качества моделей.
- Нейронные сети — отдельный класс потребителей данных с собственными требованиями к объёму и препроцессингу.
- Архитектурные паттерны — как дата-платформа встраивается в общую архитектуру системы: события, интеграции, границы сервисов.
- DDD — язык, на котором договариваются о смысле сущностей; половина споров про определение метрики — это на самом деле спор про ubiquitous language.
- Product Management — сторона, которая эти метрики заказывает: продуктовые метрики, эксперименты, приоритизация.
- Алгоритмы и Структуры данных — фундамент под всем, что мы обсуждали: сортировки и хэши в основе join’ов, вероятностные структуры в основе быстрых уникумов.
- Языковые треки, если хочется писать пайплайны руками: Go, Elixir, TypeScript, C#.
Общая карта портала и рекомендованный порядок изучения — в роадмапе. Начать трек заново или свериться с картой можно с обзорной статьи.