Зачем: тикет #23 — каждое цитируемое число курса сверено с живым стендом; до этого в уроках стояли заглушки. Что: - уроки 00–04: 18 маркеров заполнены числами свежего импорта (280 437 событий, нули в таблицах ошибок, счётчики ODS/DDS); - лаба 07: таблица manifest после дня 4 (374 092 / 34 801 / 5 388), переходящие визиты по стыкам (35/22/26), числа после дня 5; - лаба 08: каноническая граница трёх дней, пример замера свежести (лаг 3:45 модельного времени до догона, ETL ~29 с) и вернувшегося пользователя; две живые поправки разбора времени: убран принудительный UTC в разборе STG и суффикс +00:00 в сравнении границы (ловились только на живом стенде). Проверка: детерминизм подтверждён двумя независимыми циклами сброс→импорт→инкремент (числа manifest и checksum_sha256 дней 4 и 5 совпали бит в бит); chain-check зелёный; сценарий лабы 08 прогнан вживую, включая стоп/продолжение и красный full_refresh=false из урока 4; grep «сверить-на-стенде» пуст; ссылки и якоря целы; make test (219+31) и make lint зелёные; /ai-text-lint по лабам — без существенных находок. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
16 KiB
Лаба 08. Живой поток и свежесть данных
Формат: практика — запустишь живой поток, остановишь его и продолжишь с сохранённого состояния. Пререквизит: пройдена лаба 07; стенд возвращён к эталонному миру, а
world_next_dayстоит на паузе. Эталонные пути:Makefileиairflow/dags/etl_pipeline_dag.py.Поток данных одной строкой:
generator-continue → Kafka → STG (сразу) → etl_pipeline (по запуску) → ODS → DDS → DMО чём лаба простыми словами: события будут приезжать непрерывно. Ты увидишь, что Kafka и STG уже обновились, а аналитические слои ещё ждут пакетного запуска.
1. Зачем и где в проде
У данных есть свежесть — задержка между событием в источнике и моментом, когда его увидит потребитель. Для одного отчёта допустимы сутки, для другого нужны минуты. Обещание «данные готовы не позже чем через N минут» влияет на всю цепочку.
В этом стенде приём потоковый: ClickHouse постоянно читает Kafka, а
Materialized View сразу складывает сообщения в STG. Дальнейшая обработка
пакетная: ODS, DDS и DM обновляет только запуск etl_pipeline.
Поэтому потоковый приём не равен потоковой обработке. Новые строки уже могут лежать в STG, пока витрина для аналитика показывает старую правую границу.
В проде иначе. Требование к свежести задают вместе с потребителем данных. Под него выбирают расписание, размер порции, проверки и допустимое отставание. Слово «поток» само по себе ничего не обещает о времени появления данных в отчёте.
2. Руки: оживляем мониторинг и догоняем витрины
Подготовь стенд по
канонической инструкции курса.
После world_init живой генератор выключен, а эталонный мир заканчивается на
одной и той же модельной границе у всех менти.
Запоминаем начальную свежесть
Открой SQL-консоль http://localhost:9123/play (пользователь default, пароль
123456) и выполни:
SELECT
'STG' AS layer,
max(parseDateTime64BestEffortOrNull(
JSONExtractString(raw, 'event_timestamp'), 6
)) AS max_event_ts
FROM stg.browser_raw
UNION ALL
SELECT 'ODS', max(event_ts)
FROM ods.browser_event
UNION ALL
SELECT 'DDS', max(event_ts)
FROM dds.event
UNION ALL
SELECT 'DM', max(event_ts)
FROM dm.v_events_enriched;
На эталонном мире правые границы слоёв должны совпадать.
После импорта трёх эталонных дней все четыре слоя заканчиваются на
2026-01-03 23:59:59.066766.
Запускаем живое продолжение
В терминале выполни:
make generator-continue
Команда запускает генератор в режиме continue: он берёт сохранённое состояние
эталонного мира и продолжает его, не создавая новый мир с нуля.
Вернись к экранам из урока 5:
- в Prometheus target
generatorпереходит вUP; - в
Generator OverviewоживаетTotal Events/min (all 4 topics); - в
Kafka Overviewрастут сообщения и конечные offset-ы партиций; - в Kafka UI из урока 0 у новых сообщений растут offset-ы;
Kafka No Messages Producedвыходит изAlerting, когда Grafana увидит устойчивый поток.
У алерта есть окно for: 10m: в Alerting он переходит только после десяти
минут тишины. Когда поток вернётся, условие перестанет выполняться при
ближайшей оценке. График событий и offset-ы должны начать двигаться раньше.
Подожди несколько минут и снова выполни запрос по слоям. Правая граница STG
уйдёт вперёд. ODS, DDS и DM останутся на границе последнего
etl_pipeline. Запиши время по обычным часам и границы STG и DM. Границу STG
считай контрольной.
Догоняем пакетные слои
Открой Airflow и запусти etl_pipeline через Trigger DAG с пустой формой.
По умолчанию это полный пересчёт. При запуске включи секундомер.
После зелёного прогона останови секундомер и повтори запрос. Запиши новые границы STG и DM: DM должна достичь контрольной границы STG. Генератор всё ещё работает, поэтому STG вскоре снова может оказаться чуть свежее. Это ожидаемое расслоение, а не потеря строк.
3. Загляни внутрь
Открой Makefile и найди две команды этой лабы:
generator-continue:
COMPOSE_BIN="$(COMPOSE)" bash ./scripts/run_generator.sh continue
generator-down:
$(COMPOSE) stop generator
generator-down останавливает только контейнер генератора. Kafka, ClickHouse и
читатели STG продолжают работать. Поэтому они успевают дочитать уже записанные
сообщения.
generator-continue запускает генератор с сохранением прежнего состояния. В нём
остаются популяция пользователей, активные визиты, случайное состояние и
модельные часы. После рестарта продолжается тот же мир.
Модельное время. Профиль
daily-waveускоряет внутренние часы стенда в 60 раз: модельные сутки проходят примерно за полчаса. Поэтому «сейчас» на дашборде может обгонять настенные часы. Ручка тренажёра называетсяGEN_MODEL_TIME_SPEED; в этой лабе её не меняем.
Где проходит граница между потоком и пакетом
В уроке 1 ты видел Kafka engine и Materialized View. Пока сервисы работают,
этот путь сам переносит сообщения из Kafka в stg.*_raw.
У etl_pipeline в
airflow/dags/etl_pipeline_dag.py
стоит schedule=None. ODS, DDS и DM не догоняют поток сами: нужен ручной
запуск DAG. Получается две разные гарантии свежести:
- Kafka → STG работает постоянно;
- STG → ODS → DDS → DM обновляется по запуску
etl_pipeline.
Если потребителю нужна витрина не старше пяти минут, ручной DAG такую гарантию не даёт. Понадобится расписание или другая обработка, а затем измерение фактического отставания.
4. Управляемая правка: останавливаем и продолжаем поток
Проверим, переживает ли мир остановку генератора.
Шаг 1. Зафиксируй состояние до остановки
Пока генератор работает, запиши:
- правую границу STG из запроса секции 2;
- значение
Total Events/min (all 4 topics); - текущие конечные offset-ы в
Kafka Overview.
Шаг 2. Останови генератор
make generator-down
После остановки:
- target
generatorпереходит вDOWN; Total Events/min (all 4 topics)падает к нулю или перестаёт обновляться;- конечные offset-ы перестают расти;
- читатели ClickHouse дочитывают уже записанные сообщения, и отставание стекает к нулю;
- после окна ожидания
Kafka No Messages Producedпереходит вAlerting.
Как и в уроке 5, Kafka UI или панель lag могут не показать прогресс групп
ch_stg_*. Тогда смотри system.kafka_consumers со стороны ClickHouse.
Главный признак здесь такой: новые offset-ы больше не появляются, а строки,
уже записанные до остановки, доходят в STG.
Шаг 3. Продолжи тот же мир
make generator-continue
Проверь обратную картину:
- target
generatorснова вUP; - график событий и offset-ы снова растут;
Kafka No Messages Producedвозвращается в норму при ближайшей оценке;- правая граница STG стала больше значения, записанного до остановки.
Генератор взял прежнюю популяцию и прежнюю точку продолжения. Остановка не превратила поток в новый независимый мир.
Шаг 4. Догони аналитику и увидь расхождение мира
Запусти etl_pipeline с пустой формой. После зелёного прогона выполни:
WITH parseDateTime64BestEffort(
substring('<model_t_end>', 1, 19), 6
) AS boundary
SELECT
user_domain_id,
countIf(visit_start < boundary) AS visits_before,
countIf(visit_start >= boundary) AS visits_after,
min(visit_start) AS first_visit_ts, max(visit_start) AS last_visit_ts
FROM (
SELECT user_domain_id, click_id, min(event_ts) AS visit_start
FROM dm.v_events_enriched
WHERE user_domain_id IS NOT NULL
GROUP BY user_domain_id, click_id
)
GROUP BY user_domain_id
HAVING min(visit_start) < boundary AND max(visit_start) >= boundary
LIMIT 10;
Подставь вместо <model_t_end> границу из манифеста. Успех — хотя бы одна
строка. Иначе подожди несколько минут, перезапусти etl_pipeline и повтори
запрос. Это симметрия лаб: в лабе 07 визит пересекал полночь, здесь
пользователь с разными визитами пересекает границу замороженного мира.
Мир расходится: итоги превышают эталонные и различаются у менти. Это правильно.
Например, для границы 2026-01-06T00:00:00+00:00 найден пользователь
5938d296-14cf-49cb-801e-82370518cf59: 13 визитов до границы и 2 после
неё. Твои числа и идентификатор будут другими.
Верни как было
Останови живой поток:
make generator-down
После continue мир стал недетерминированным: его числа зависят от длительности
запуска. Верни стенд к эталону по
канонической инструкции курса.
Не пересказывай шаги по памяти: эта ссылка остаётся единственным учебным
описанием полного сброса.
5. Проверь себя
| Действие | Где смотреть | Что ожидать |
|---|---|---|
выполнить make generator-continue |
Prometheus Targets | generator переходит в UP |
| измерить свежесть до и после ETL | запрос из секции 2 и часы | например: до ETL — 3:45:01.519535, после — 0:45:55.977582 модельного времени; настенное ожидание — около 29 секунд |
выполнить make generator-down |
Grafana и Kafka | события перестают поступать, отставание читателей стекает к нулю |
снова выполнить make generator-continue |
Grafana, Kafka и STG | поток продолжается, offset-ы и правая граница снова растут |
| найти вернувшегося пользователя | запрос по границе model_t_end |
есть пользователь с визитами до и после границы |
Ответь своими словами:
- почему STG опережает DM, а
etl_pipelineдогоняет лишь очередной снимок; - что сохраняет
continueи почему итоги менти расходятся; - какую свежесть можно честно обещать при ручном запуске и при расписании каждые 30 минут;
- почему минутная свежесть лежит за пределами этого стенда.
6. Что должно получиться
После лабы сохрани:
- скрин
generatorвDOWNпослеmake generator-downи скринgeneratorвUPпосле повторногоmake generator-continue; - границы STG до остановки и после продолжения;
- строку пользователя с визитами до и после
model_t_end; - замеры до и после
etl_pipeline: разницуSTG − DMв модельном времени и настенное ожидание до появления зафиксированной границы STG в DM; - вывод своими словами: что можно обещать при ручном запуске, что — при расписании каждые 30 минут и почему этот стенд не обещает минутную свежесть.
В конце generator должен быть остановлен, а стенд — возвращён к эталонному
миру.
Мост после курса
Замер отвечает на вопрос из секции 1: ручной запуск не ограничивает ожидание, а расписание раз в 30 минут дало бы около 30 минут плюс время прогона. Для минутной свежести нужна другая обработка.