Как посчитать Data Freshness в SQL
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 наверх.
История свежести
Мониторить только текущий лаг недостаточно — полезно видеть тренд, чтобы ловить деградацию заранее. Если писать замеры свежести в лог, можно считать среднюю, максимальную и 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 потока, и нельзя сравнивать их напрямую.
Связанные темы
- Как посчитать data quality score в SQL
- Как посчитать duplicate rate в SQL
- SLA / SLO / SLI на собесе SA
- Data quality dimensions на собесе DE
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 лага), а текущее состояние удобно держать в материализованном представлении, чтобы мониторинг читал его быстро.