Как посчитать Data Freshness в SQL

Проверь себя · 1/3разбор после ответа
Вы строите выдачу «последние заказы» и используете LIMIT 50. Поле created_at не уникально (много заказов в одну секунду). Какой ORDER BY лучше, чтобы порядок был детерминированным?

Зачем Data Freshness

Data freshness (свежесть данных) — это лаг между тем, как событие произошло в реальности, и тем, как оно стало доступным в хранилище (DWH). Устаревшие (stale) данные приводят к неверным решениям: аналитик смотрит на дашборд, думает, что видит вчерашнюю картину, а на самом деле pipeline встал три дня назад и данные не обновлялись. Именно поэтому свежесть — одна из ключевых характеристик качества данных, а SLA на неё прописывают отдельно: для критичных метрик обычно требуют лаг менее часа, для ежедневных отчётов — менее суток.

Для data-инженера это метрика самоконтроля: она позволяет заметить сломанный или отстающий pipeline раньше, чем на это пожалуется бизнес.

Формула

Freshness Lag = now() - max(event_time)

То есть свежесть — это возраст самой свежей записи в таблице: сколько времени прошло с момента последнего события до текущего момента. Для пакетных (batch) задач часто удобнее считать иначе — от последнего успешного запуска pipeline:

Freshness = now() - last_successful_run_timestamp

Эти два определения дают разные цифры, и на собесе важно не путать их: первое меряет возраст данных, второе — возраст самого прогона.

Базовый расчёт

Считаем лаг как разницу между текущим временем и самым свежим событием. EXTRACT(EPOCH FROM ...) даёт разницу в секундах, делим на 60 для минут:

SELECT
    MAX(event_time) AS latest_event,
    NOW() AS now_ts,
    EXTRACT(EPOCH FROM (NOW() - MAX(event_time))) / 60 AS lag_minutes
FROM events;

Чтобы одним запросом проверить свежесть нескольких таблиц, собираем их через UNION ALL:

SELECT
    'events' AS table_name,
    MAX(event_time) AS latest,
    EXTRACT(EPOCH FROM (NOW() - MAX(event_time))) / 60 AS lag_minutes
FROM events
UNION ALL
SELECT
    'transactions',
    MAX(created_at),
    EXTRACT(EPOCH FROM (NOW() - MAX(created_at))) / 60
FROM transactions;

По pipelines

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

SELECT
    pipeline_name,
    last_run_timestamp,
    last_successful_timestamp,
    EXTRACT(EPOCH FROM (NOW() - last_successful_timestamp)) / 60 AS minutes_since_success,
    rows_processed_last_run,
    status_last_run
FROM pipeline_metadata
WHERE pipeline_name LIKE 'critical_%'
ORDER BY minutes_since_success DESC;

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

Мониторинг SLA

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

WITH freshness AS (
    SELECT
        pipeline_name,
        EXTRACT(EPOCH FROM (NOW() - last_successful_timestamp)) / 60 AS lag_min,
        sla_minutes
    FROM pipeline_metadata
)
SELECT
    pipeline_name,
    lag_min,
    sla_minutes,
    CASE
        WHEN lag_min > sla_minutes * 1.5 THEN 'CRITICAL'
        WHEN lag_min > sla_minutes THEN 'WARN'
        ELSE 'OK'
    END AS status,
    lag_min / sla_minutes AS sla_ratio
FROM freshness
ORDER BY sla_ratio DESC;

sla_ratio показывает, во сколько раз лаг превысил порог: значение больше 1 — нарушение SLA, а сортировка по нему выводит самые проблемные pipeline наверх.

Закрепи формулу data freshness в Карьернике
Запомнить надолго — 5 коротких сессий с задачами на эту тему. Бесплатно
Тренировать data freshness в Telegram

История свежести

Мониторить только текущий лаг недостаточно — полезно видеть тренд, чтобы ловить деградацию заранее. Если писать замеры свежести в лог, можно считать среднюю, максимальную и P95-свежесть по часам за неделю:

SELECT
    DATE_TRUNC('hour', TIMESTAMP) AS hour,
    pipeline_name,
    AVG(lag_minutes) AS avg_lag,
    MAX(lag_minutes) AS max_lag,
    PERCENTILE_CONT(0.95) WITHIN GROUP (ORDER BY lag_minutes) AS p95_lag
FROM freshness_log
WHERE TIMESTAMP >= CURRENT_DATE - INTERVAL '7 days'
GROUP BY 1, 2
ORDER BY 1, 2;

Растущий из часа в час P95 — сигнал, что pipeline постепенно не справляется с объёмом, ещё до того как он формально нарушит SLA.

Как это спрашивают на собесе

«Как понять, что данные в дашборде свежие?» Сильный ответ — не «глазами посмотреть на дату», а метрика: считать лаг между MAX(event_time) и NOW() и сравнивать его с SLA. И обязательно уточнить, от чего меряем возраст — от события или от прогона pipeline. Прогнать такие вопросы по data quality с разбором можно в Карьернике.

«Таблица пустая, ваш запрос вернул NULL — это свежие данные или сломанный pipeline?» Ловушка: MAX(...) по пустой таблице даёт NULL, и его нельзя трактовать как «лаг ноль». Пустота — это, наоборот, повод для тревоги, и её нужно обрабатывать отдельно.

«Событие пришло с задержкой в сутки — как это влияет на freshness?» Здесь всплывает разница между временем события и временем загрузки: мобильные события часто приходят с задержкой из-за офлайн-синхронизации, поэтому event_time может быть старым при свежей загрузке. Хороший кандидат проговорит late-arriving events и watermarking.

Частые ошибки

Wall-clock против event-time. Нужно чётко определить, что именно вы меряете: лаг между временем загрузки и временем события, или между текущим прогоном и предыдущим. Это разные метрики, и смешивать их нельзя — сначала фиксируем определение, потом считаем.

Late-arriving events. События с мобильных клиентов часто приходят с задержкой из-за офлайн-синхронизации: event_time оказывается сильно старее момента загрузки. Из-за этого свежесть по event-time может выглядеть плохой, хотя pipeline работает штатно.

Пустая таблица не равна нулевому лагу. MAX(...) по таблице без строк вернёт NULL, а не ноль. Не путайте это со «свежими данными»: пустота обычно означает, что загрузка не отработала, и это отдельный класс проблем.

Разные часовые пояса. Если источник пишет время в UTC, а хранилище — в локальном поясе, лаг получится искусственно огромным или отрицательным. Приводите обе метки времени к одному поясу перед вычитанием.

Сравнивать streaming и batch по одной планке. У стримингового pipeline свежесть измеряется минутами, у пакетного — часами. Это нормально: нельзя требовать от ночного batch-отчёта той же свежести, что от near-real-time потока, и нельзя сравнивать их напрямую.

Связанные темы

FAQ

Какой SLA по свежести считать нормальным?

Зависит от типа pipeline. Для real-time потоков ориентируются на лаг менее 5 минут, для near-real-time — менее часа, для ежедневного batch — менее 24 часов. Конкретную планку задаёт бизнес-требование: критичность метрики определяет, сколько устаревания допустимо.

Что мерять — wall-clock или event-time?

Лучше и то, и другое, потому что они отвечают на разные вопросы. Wall-clock (время загрузки) — операционная метрика: работает ли pipeline прямо сейчас. Event-time (возраст события) — аналитическая: насколько актуальны данные, которые видит пользователь. При late-arriving events эти две цифры расходятся, и важно понимать, какая из них релевантна вашей задаче.

Как обрабатывать поздно приходящие события?

Стандартный подход — watermarking (концепция из потоковых движков Flink и Beam): система задаёт «водяной знак» — границу, до которой ждёт опоздавшие события, и всё, что пришло позже порога, отбрасывает или отправляет в отдельную обработку. Так удаётся не держать окно открытым бесконечно и при этом не терять разумно опоздавшие данные.

Чем freshness отличается от latency?

Latency — это задержка одного конкретного события на пути от источника до хранилища. Freshness — агрегатная характеристика таблицы или pipeline: обычно это возраст самой свежей записи (максимальный лаг). Latency отвечает на вопрос «как долго ехало вот это событие», freshness — «насколько устарели данные в целом».

Какую схему завести для отслеживания свежести?

Обычно заводят таблицу вида pipeline_metadata(pipeline, last_run, last_success, sla, status) для текущего состояния и отдельный append-only лог замеров для истории. По логу считают тренды (avg, max, P95 лага), а текущее состояние удобно держать в материализованном представлении, чтобы мониторинг читал его быстро.