# Лаба 08. Живой поток и свежесть данных > Формат: **практика** — запустишь живой поток, остановишь его и продолжишь с > сохранённого состояния. > Пререквизит: пройдена лаба 07; стенд возвращён к эталонному миру, а > `world_next_day` стоит на паузе. > Эталонные пути: > [`Makefile`](../../../Makefile) и > [`airflow/dags/etl_pipeline_dag.py`](../../../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. Руки: оживляем мониторинг и догоняем витрины Подготовь стенд по [канонической инструкции курса](../README.md#подготовка-и-канонический-сброс). После `world_init` живой генератор выключен, а эталонный мир заканчивается на одной и той же модельной границе у всех менти. ### Запоминаем начальную свежесть Открой SQL-консоль `http://localhost:9123/play` (пользователь `default`, пароль `123456`) и выполни: ```sql SELECT 'STG' AS layer, max(parseDateTime64BestEffortOrNull( JSONExtractString(raw, 'event_timestamp'), 6, 'UTC' )) 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; ``` На эталонном мире правые границы слоёв должны совпадать. ### Запускаем живое продолжение В терминале выполни: ```bash 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`](../../../Makefile) и найди две команды этой лабы: ```make 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`](../../../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. Останови генератор ```bash 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. Продолжи тот же мир ```bash make generator-continue ``` Проверь обратную картину: - target `generator` снова в `UP`; - график событий и offset-ы снова растут; - `Kafka No Messages Produced` возвращается в норму при ближайшей оценке; - правая граница STG стала больше значения, записанного до остановки. Генератор взял прежнюю популяцию и прежнюю точку продолжения. Остановка не превратила поток в новый независимый мир. ### Шаг 4. Догони аналитику и увидь расхождение мира Запусти `etl_pipeline` с пустой формой. После зелёного прогона выполни: ```sql WITH parseDateTime64BestEffort('', 6, 'UTC') 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; ``` Подставь вместо `` границу из манифеста. Успех — хотя бы одна строка. Иначе подожди несколько минут, перезапусти `etl_pipeline` и повтори запрос. Это симметрия лаб: в лабе 07 визит пересекал полночь, здесь пользователь с разными визитами пересекает границу замороженного мира. Мир расходится: итоги превышают эталонные и различаются у менти. Это правильно. ### Верни как было Останови живой поток: ```bash make generator-down ``` После `continue` мир стал недетерминированным: его числа зависят от длительности запуска. Верни стенд к эталону по [канонической инструкции курса](../README.md#подготовка-и-канонический-сброс). Не пересказывай шаги по памяти: эта ссылка остаётся единственным учебным описанием полного сброса. --- ## 5. Проверь себя | Действие | Где смотреть | Что ожидать | |----------|--------------|-------------| | выполнить `make generator-continue` | Prometheus Targets | `generator` переходит в `UP` | | измерить свежесть до и после ETL | запрос из секции 2 и часы | записаны модельное отставание и настенное ожидание | | выполнить `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 минут плюс время прогона. Для минутной свежести нужна другая обработка.