docs(course): переведены уроки на стартовую историю
- Зачем: - учебный путь должен идти через генерацию и штатный пайплайн, а не через архивный сид. - Что: - обновлены уроки 00-06 и стандарт урока под startup-history/backfill. - объяснено, что data/*.jsonl остаются кладовкой значений генератора. - тест-план переведён на новый штатный запуск и HITL-приёмку. - Проверка: - rg -n \"make clean/up/ddl/data/transform|make data|kafka_load|LIMIT=|2022-11-28|26 из 50\" docs/course docs/TEST_PLAN.md. - git diff --cached --check.
This commit is contained in:
+22
-40
@@ -8,15 +8,17 @@
|
|||||||
|
|
||||||
Проверить, что стек `Kafka + ClickHouse + Airflow + Superset + Prometheus/Grafana`:
|
Проверить, что стек `Kafka + ClickHouse + Airflow + Superset + Prometheus/Grafana`:
|
||||||
- стабильно поднимается;
|
- стабильно поднимается;
|
||||||
- загружает данные по пути `kafka_load -> STG -> etl_pipeline -> ODS/DDS/DM`;
|
- загружает данные по пути `startup-history/backfill -> Kafka -> STG -> ODS/DDS/DM`;
|
||||||
- корректно обрабатывает «грязные» записи (ошибки фиксируются в ODS, пайплайн не падает);
|
- корректно обрабатывает «грязные» записи (ошибки фиксируются в ODS, пайплайн не падает);
|
||||||
- отдает метрики и дашборды мониторинга.
|
- отдает метрики и дашборды мониторинга.
|
||||||
|
|
||||||
## Общие принципы
|
## Общие принципы
|
||||||
|
|
||||||
- По умолчанию используем малый срез (`limit=50`) для быстрых и повторяемых проверок.
|
- По умолчанию используем быстрый профиль стартовой истории (`PROFILE=ci`).
|
||||||
- Полный прогон (`limit=0`) выполняем отдельно как long-run сценарий.
|
- Полный прогон выполняем отдельно через `PROFILE=daily-wave`.
|
||||||
- Основной путь запуска — через Airflow DAG.
|
- Основной ручной путь запуска — через Airflow DAG `generator_control`.
|
||||||
|
- Консольный чистый прогон `make generated-history-analytics` остаётся коротким
|
||||||
|
повторяемым сценарием для smoke и CI.
|
||||||
- Критерий успеха: не только `Success` DAG, но и проверки данных/ошибок/мониторинга.
|
- Критерий успеха: не только `Success` DAG, но и проверки данных/ошибок/мониторинга.
|
||||||
|
|
||||||
---
|
---
|
||||||
@@ -28,10 +30,10 @@
|
|||||||
### A.1 Подготовка окружения
|
### A.1 Подготовка окружения
|
||||||
|
|
||||||
```bash
|
```bash
|
||||||
# Полная очистка стенда
|
# Чистый быстрый прогон: backfill -> Kafka -> STG -> ODS -> DDS -> DM -> Superset
|
||||||
make clean
|
make generated-history-analytics
|
||||||
|
|
||||||
# Запуск инфраструктуры
|
# Поднять остальные UI-сервисы после чистого прогона
|
||||||
make up
|
make up
|
||||||
|
|
||||||
# Проверка контейнеров
|
# Проверка контейнеров
|
||||||
@@ -42,16 +44,7 @@ docker compose ps
|
|||||||
- `airflow-init` в `Exited (0)`;
|
- `airflow-init` в `Exited (0)`;
|
||||||
- остальные сервисы в `Up` (включая `superset`, `prometheus`, `grafana`, `kafka-exporter`, `statsd-exporter`).
|
- остальные сервисы в `Up` (включая `superset`, `prometheus`, `grafana`, `kafka-exporter`, `statsd-exporter`).
|
||||||
|
|
||||||
### A.2 DDL и минимальная загрузка данных
|
### A.2 Проверка стартовой истории в STG
|
||||||
|
|
||||||
```bash
|
|
||||||
# Инициализация схемы
|
|
||||||
docker compose exec -T airflow-webserver airflow dags trigger ddl_init
|
|
||||||
|
|
||||||
# Быстрый ingest: по 50 строк на поток
|
|
||||||
docker compose exec -T airflow-webserver airflow dags trigger kafka_load \
|
|
||||||
--conf '{"limit": 50, "reset_topics": true}'
|
|
||||||
```
|
|
||||||
|
|
||||||
Проверки:
|
Проверки:
|
||||||
|
|
||||||
@@ -72,12 +65,7 @@ SELECT 'geo_raw', count() FROM stg.geo_raw
|
|||||||
|
|
||||||
Ожидаем: во всех 4 таблицах `cnt > 0`.
|
Ожидаем: во всех 4 таблицах `cnt > 0`.
|
||||||
|
|
||||||
### A.3 ETL и проверки слоев
|
### A.3 Проверки слоев
|
||||||
|
|
||||||
```bash
|
|
||||||
docker compose exec -T airflow-webserver airflow dags trigger etl_pipeline \
|
|
||||||
--conf '{"full_refresh": true}'
|
|
||||||
```
|
|
||||||
|
|
||||||
Проверки:
|
Проверки:
|
||||||
|
|
||||||
@@ -123,7 +111,7 @@ curl -s -u admin:admin "http://localhost:3000/api/dashboards/uid/airflow-overvie
|
|||||||
|
|
||||||
Критерий успеха smoke:
|
Критерий успеха smoke:
|
||||||
- сервисы подняты;
|
- сервисы подняты;
|
||||||
- 3 DAG (`ddl_init`, `kafka_load`, `etl_pipeline`) успешны;
|
- стартовая история доведена до DM;
|
||||||
- данные проходят до DM;
|
- данные проходят до DM;
|
||||||
- мониторинг и дашборды доступны.
|
- мониторинг и дашборды доступны.
|
||||||
|
|
||||||
@@ -133,29 +121,22 @@ curl -s -u admin:admin "http://localhost:3000/api/dashboards/uid/airflow-overvie
|
|||||||
|
|
||||||
Ожидаемая длительность: ~30-60 минут.
|
Ожидаемая длительность: ~30-60 минут.
|
||||||
|
|
||||||
### B.1 Полная загрузка и полный ETL
|
### B.1 Полная стартовая история и ETL
|
||||||
|
|
||||||
```bash
|
```bash
|
||||||
# Чистый старт
|
# История на 2 суток с видимой суточной волной
|
||||||
make clean && make up
|
PROFILE=daily-wave make generated-history-analytics
|
||||||
|
make up
|
||||||
# DDL
|
|
||||||
docker compose exec -T airflow-webserver airflow dags trigger ddl_init
|
|
||||||
|
|
||||||
# Полный ingest (все строки)
|
|
||||||
docker compose exec -T airflow-webserver airflow dags trigger kafka_load \
|
|
||||||
--conf '{"limit": 0, "reset_topics": true}'
|
|
||||||
|
|
||||||
# Полный ETL
|
|
||||||
docker compose exec -T airflow-webserver airflow dags trigger etl_pipeline \
|
|
||||||
--conf '{"full_refresh": true}'
|
|
||||||
```
|
```
|
||||||
|
|
||||||
### B.2 Проверка объемов и DQ
|
### B.2 Проверка объемов и DQ
|
||||||
|
|
||||||
```bash
|
```bash
|
||||||
# Фактические размеры входа
|
# Диапазон модельного времени в витрине
|
||||||
wc -l data/*.jsonl
|
docker compose exec -T clickhouse clickhouse-client --user=default --password=123456 --query "
|
||||||
|
SELECT min(event_ts) AS min_event_ts, max(event_ts) AS max_event_ts, count() AS events
|
||||||
|
FROM dm.v_events_enriched
|
||||||
|
"
|
||||||
|
|
||||||
# Сводка по слоям
|
# Сводка по слоям
|
||||||
docker compose exec -T clickhouse clickhouse-client --user=default --password=123456 --query "
|
docker compose exec -T clickhouse clickhouse-client --user=default --password=123456 --query "
|
||||||
@@ -185,6 +166,7 @@ docker compose exec -T clickhouse clickhouse-client --user=default --password=12
|
|||||||
Критерии успеха full:
|
Критерии успеха full:
|
||||||
- STG заполнен по всем 4 потокам;
|
- STG заполнен по всем 4 потокам;
|
||||||
- ODS/DDS/DM заполнены;
|
- ODS/DDS/DM заполнены;
|
||||||
|
- диапазон `event_ts` соответствует созданной стартовой истории;
|
||||||
- `dm.dq_summary` не пуста и содержит метрики всех слоев.
|
- `dm.dq_summary` не пуста и содержит метрики всех слоев.
|
||||||
|
|
||||||
### B.3 Тест восстановления stop/start
|
### B.3 Тест восстановления stop/start
|
||||||
|
|||||||
@@ -21,7 +21,7 @@
|
|||||||
запись и увидеть `+1` в таблице ошибок `parse_errors`; сменить `kafka_group_name`
|
запись и увидеть `+1` в таблице ошибок `parse_errors`; сменить `kafka_group_name`
|
||||||
и увидеть, как топик читается заново. Менти меняет — видит эффект — объясняет.
|
и увидеть, как топик читается заново. Менти меняет — видит эффект — объясняет.
|
||||||
**Верни как было** — каждая правка завершается явным шагом отката к чистому
|
**Верни как было** — каждая правка завершается явным шагом отката к чистому
|
||||||
состоянию (откатить изменение либо `make clean/up/ddl/data/transform`), чтобы
|
состоянию (откатить изменение либо `make generated-history-analytics && make up`), чтобы
|
||||||
самостоятельный менти не застрял со сломанным стендом без ментора.
|
самостоятельный менти не застрял со сломанным стендом без ментора.
|
||||||
(Урок 0 — без этого шага, только наблюдение.)
|
(Урок 0 — без этого шага, только наблюдение.)
|
||||||
5. **Проверь себя** — самопроверка (раздел 4).
|
5. **Проверь себя** — самопроверка (раздел 4).
|
||||||
@@ -112,7 +112,7 @@
|
|||||||
- штатные быстрые проверки (smoke) из `docs/TEST_PLAN.md`;
|
- штатные быстрые проверки (smoke) из `docs/TEST_PLAN.md`;
|
||||||
- встроенные проверки в `etl_pipeline_dag.py` (DAG падает на пустой витрине или
|
- встроенные проверки в `etl_pipeline_dag.py` (DAG падает на пустой витрине или
|
||||||
нарушении целостности);
|
нарушении целостности);
|
||||||
- `make clean/up/ddl/data/transform` для сброса и повтора.
|
- `make generated-history-analytics && make up` для чистого сброса и повтора штатного пути.
|
||||||
|
|
||||||
В каждом уроке — маленькая табличка самопроверки в формате
|
В каждом уроке — маленькая табличка самопроверки в формате
|
||||||
**действие → где смотреть → что ожидать**, а для управляемой правки — какой
|
**действие → где смотреть → что ожидать**, а для управляемой правки — какой
|
||||||
@@ -125,8 +125,8 @@
|
|||||||
- **Проверяй на стенде, а не «на глаз».** Любой категоричный claim про числа или целостность
|
- **Проверяй на стенде, а не «на глаз».** Любой категоричный claim про числа или целостность
|
||||||
подтверждай командой на стенде. Так нашёлся баг с `kafka_ts` в уроке 1 (значение типа
|
подтверждай командой на стенде. Так нашёлся баг с `kafka_ts` в уроке 1 (значение типа
|
||||||
`DateTime64` молча резалось через `toInt64`, и время по всему стенду уехало в `1970`) и
|
`DateTime64` молча резалось через `toInt64`, и время по всему стенду уехало в `1970`) и
|
||||||
подтвердилось, что «26 из 50» в уроке 2 — это схлопывание повторов по `click_id` (проверено
|
подтвердилось, что меньший размер таблиц контекста в уроке 2 — это схлопывание повторов
|
||||||
`uniqExact`), а не потеря данных.
|
по `click_id` (проверяется через `uniqExact`), а не потеря данных.
|
||||||
- **Имена таблиц и колонок в тексте сверяй с реальным DDL.** То, что написано в уроке, должно
|
- **Имена таблиц и колонок в тексте сверяй с реальным DDL.** То, что написано в уроке, должно
|
||||||
совпадать с `sql/ddl/...`, иначе менти запутается, когда выполнит запрос на стенде.
|
совпадать с `sql/ddl/...`, иначе менти запутается, когда выполнит запрос на стенде.
|
||||||
- **Один паттерн на урок; будущие темы только анонсируй.** Не раскрывай то, что относится к
|
- **Один паттерн на урок; будущие темы только анонсируй.** Не раскрывай то, что относится к
|
||||||
|
|||||||
+11
-6
@@ -25,16 +25,21 @@
|
|||||||
слова с живым стендом.
|
слова с живым стендом.
|
||||||
- **Железо:** стек тяжёлый — Kafka, ClickHouse, Airflow, Superset, Prometheus и Grafana
|
- **Железо:** стек тяжёлый — Kafka, ClickHouse, Airflow, Superset, Prometheus и Grafana
|
||||||
поднимаются одновременно. Нужна машина, которая это потянет.
|
поднимаются одновременно. Нужна машина, которая это потянет.
|
||||||
- **Подними стенд и залей малый срез** (из корня репозитория) — этого хватит, чтобы
|
- **Подними стенд и создай стартовую историю** (из корня репозитория) — этого хватит,
|
||||||
начать, и прогон быстрый:
|
чтобы начать, и прогон быстрый:
|
||||||
|
|
||||||
```bash
|
```bash
|
||||||
make up # поднять контейнеры
|
make generated-history-analytics
|
||||||
make ddl # создать схему в ClickHouse (базы, таблицы, VIEW)
|
make up
|
||||||
LIMIT=50 make data # залить по 50 строк на топик в Kafka
|
|
||||||
```
|
```
|
||||||
|
|
||||||
Дальше каждый урок в секции «Руки» сам напоминает, что перезапустить и с каким срезом.
|
Эта команда проводит штатный путь стенда: готовый источник данных создаёт стартовую
|
||||||
|
историю, события попадают в Kafka, затем в STG, ODS, DDS, DM и Superset. Файлы
|
||||||
|
`data/*.jsonl` пока остаются только кладовкой готовых значений для генератора
|
||||||
|
(браузеры, страны, устройства, UTM), а не источником аналитического контура. `make up`
|
||||||
|
после неё поднимает остальные UI-сервисы курса: Kafka UI, Airflow, Prometheus и Grafana.
|
||||||
|
|
||||||
|
Дальше каждый урок в секции «Руки» сам напоминает, что перезапустить.
|
||||||
Точные шаги, параметры и troubleshooting — в [`docs/OPERATIONS.md`](../OPERATIONS.md).
|
Точные шаги, параметры и troubleshooting — в [`docs/OPERATIONS.md`](../OPERATIONS.md).
|
||||||
|
|
||||||
## Уроки
|
## Уроки
|
||||||
|
|||||||
@@ -7,7 +7,7 @@
|
|||||||
> Эталонный путь: сам Kafka UI на `http://localhost:8082`.
|
> Эталонный путь: сам Kafka UI на `http://localhost:8082`.
|
||||||
>
|
>
|
||||||
> Поток данных одной строкой:
|
> Поток данных одной строкой:
|
||||||
> `make data → топики Kafka (партиции, offset'ы) → consumer-группа ClickHouse вычитывает`
|
> `готовый источник данных стенда → топики Kafka (партиции, offset'ы) → consumer-группа ClickHouse вычитывает`
|
||||||
>
|
>
|
||||||
> О чём урок простыми словами: ходим по Kafka UI и разглядываем поток — где лежат события,
|
> О чём урок простыми словами: ходим по Kafka UI и разглядываем поток — где лежат события,
|
||||||
> кто их читает и как Kafka помнит, до какого места уже дочитано.
|
> кто их читает и как Kafka помнит, до какого места уже дочитано.
|
||||||
@@ -22,11 +22,15 @@
|
|||||||
consumer читает в своём темпе, не трогая producer'а. Если читатель отстал или прилёг —
|
consumer читает в своём темпе, не трогая producer'а. Если читатель отстал или прилёг —
|
||||||
события не теряются, они лежат в топике и ждут.
|
события не теряются, они лежат в топике и ждут.
|
||||||
|
|
||||||
На нашем стенде роли уже расставлены: `make data` играет роль трекера и заливает события
|
На нашем стенде роли уже расставлены: готовый источник данных стенда играет роль трекера
|
||||||
в топики, а ClickHouse — это consumer, который их вычитывает. В этом уроке мы не запускаем
|
и пишет события в топики, а ClickHouse — это consumer, который их вычитывает. Для курса нам
|
||||||
пайплайн и ничего не меняем — мы открываем Kafka UI и **узнаём в живом кластере** те самые
|
важно не устройство источника, а сам путь данных: Kafka → STG → ODS → DDS → DM → Superset.
|
||||||
понятия из видео. Это разминка: в уроке 1 ты уже руками увидишь, как эти же сообщения
|
Файлы `data/*.jsonl` здесь не источник аналитики; пока это только кладовка готовых значений
|
||||||
доезжают до таблиц ClickHouse.
|
для источника данных стенда.
|
||||||
|
|
||||||
|
В этом уроке мы не запускаем пайплайн и ничего не меняем — мы открываем Kafka UI и
|
||||||
|
**узнаём в живом кластере** те самые понятия из видео. Это разминка: в уроке 1 ты уже руками
|
||||||
|
увидишь, как эти же сообщения доезжают до таблиц ClickHouse.
|
||||||
|
|
||||||
> **В проде так же, только крупнее.** Здесь один брокер и крошечный срез данных. В бою
|
> **В проде так же, только крупнее.** Здесь один брокер и крошечный срез данных. В бою
|
||||||
> брокеров несколько, топик разбит на много партиций, читателей в группе — тоже несколько,
|
> брокеров несколько, топик разбит на много партиций, читателей в группе — тоже несколько,
|
||||||
@@ -37,8 +41,14 @@ consumer читает в своём темпе, не трогая producer'а.
|
|||||||
|
|
||||||
## 2. Наблюдай: открой Kafka UI
|
## 2. Наблюдай: открой Kafka UI
|
||||||
|
|
||||||
Стенд уже должен быть поднят (`make up`) и в топиках должны лежать события
|
Стенд уже должен быть поднят, а в топиках должны лежать события после стартовой истории:
|
||||||
(`LIMIT=50 make data` из урока 1 — или любой прошлый прогон). Открой Kafka UI:
|
|
||||||
|
```bash
|
||||||
|
make generated-history-analytics
|
||||||
|
make up
|
||||||
|
```
|
||||||
|
|
||||||
|
Открой Kafka UI:
|
||||||
`http://localhost:8082`. Ходи по нему свободно — это режим чтения, сломать тут ничего нельзя.
|
`http://localhost:8082`. Ходи по нему свободно — это режим чтения, сломать тут ничего нельзя.
|
||||||
|
|
||||||
Пройди по трём экранам и просто посмотри.
|
Пройди по трём экранам и просто посмотри.
|
||||||
@@ -54,20 +64,20 @@ consumer читает в своём темпе, не трогая producer'а.
|
|||||||
топиков) — его Kafka использует сама, мы его не трогаем. У каждого нашего топика в колонке
|
топиков) — его Kafka использует сама, мы его не трогаем. У каждого нашего топика в колонке
|
||||||
с партициями стоит **1**: топик маленький, делить не на что.
|
с партициями стоит **1**: топик маленький, делить не на что.
|
||||||
|
|
||||||
**Сообщения в топике.** Открой `browser_events` → вкладку *Messages*. Это и есть события,
|
**Сообщения в топике.** Открой `browser_events` → вкладку *Messages*. Это и есть события
|
||||||
которые залил `make data`. У каждого сообщения видно:
|
стенда. У каждого сообщения видно:
|
||||||
|
|
||||||
- **Offset** — порядковый номер сообщения в партиции (0, 1, 2, …);
|
- **Offset** — порядковый номер сообщения в партиции (0, 1, 2, …);
|
||||||
- **Timestamp** — когда сообщение легло в Kafka (время *доставки*, не время самого события);
|
- **Timestamp** — когда сообщение легло в Kafka (время *доставки*, не время самого события);
|
||||||
- **Value** — тело: JSON события целиком, например:
|
- **Value** — тело: JSON события целиком, например:
|
||||||
|
|
||||||
```json
|
```json
|
||||||
{"event_id": "8cca1c7d-...", "event_timestamp": "2022-11-28 20:51:05.627882",
|
{"event_id": "8cca1c7d-...", "event_timestamp": "2026-01-01 00:01:00.000000",
|
||||||
"event_type": "pageview", "browser_name": "Chrome", "browser_language": "sat_IN"}
|
"event_type": "pageview", "browser_name": "Chrome", "browser_language": "sat_IN"}
|
||||||
```
|
```
|
||||||
|
|
||||||
Загляни внутрь Value: у события есть своё `event_timestamp` (когда оно случилось,
|
Загляни внутрь Value: у события есть своё `event_timestamp` (когда оно случилось,
|
||||||
здесь — 2022 год), и оно отличается от Kafka-Timestamp (когда оно попало в топик — при заливке стенда).
|
модельное время стенда), и оно отличается от Kafka-Timestamp (когда оно попало в топик).
|
||||||
Два разных времени у одной записи — запомни этот момент, в уроке 1 он всплывёт уже на
|
Два разных времени у одной записи — запомни этот момент, в уроке 1 он всплывёт уже на
|
||||||
стороне ClickHouse.
|
стороне ClickHouse.
|
||||||
|
|
||||||
|
|||||||
@@ -49,13 +49,14 @@
|
|||||||
|
|
||||||
## 2. Руки: убедись, что данные текут
|
## 2. Руки: убедись, что данные текут
|
||||||
|
|
||||||
Поднимаем стенд, создаём схему и заливаем **малый срез** (50 строк на топик — этого
|
Поднимаем стенд и создаём стартовую историю. Это штатный путь курса: готовый источник
|
||||||
хватает, чтобы всё увидеть, и прогон быстрый):
|
данных стенда пишет события в Kafka, ClickHouse читает их в STG, затем batch строит
|
||||||
|
ODS, DDS и DM. Файлы `data/*.jsonl` в этом пути не источник аналитики; пока это только
|
||||||
|
кладовка значений для источника данных стенда.
|
||||||
|
|
||||||
```bash
|
```bash
|
||||||
make up # поднять инфраструктуру (Kafka, ClickHouse, ...)
|
make generated-history-analytics
|
||||||
make ddl # создать базы и таблицы в ClickHouse (в т.ч. слой STG)
|
make up
|
||||||
LIMIT=50 make data # залить по 50 строк каждого файла в Kafka-топики
|
|
||||||
```
|
```
|
||||||
|
|
||||||
Теперь смотрим, что доехало до ClickHouse. Открой SQL-консоль:
|
Теперь смотрим, что доехало до ClickHouse. Открой SQL-консоль:
|
||||||
@@ -196,11 +197,12 @@ FROM stg.kafka_browser_raw;
|
|||||||
TRUNCATE TABLE stg.browser_raw;
|
TRUNCATE TABLE stg.browser_raw;
|
||||||
```
|
```
|
||||||
|
|
||||||
Затем перезаливаем — с пересозданием топиков, чтобы Kafka-движок перечитал сообщения
|
Затем возвращаем стенд в чистое состояние и заново создаём стартовую историю, чтобы
|
||||||
с начала (без этого он считает их уже прочитанными и ничего нового не подхватит):
|
Kafka-движок прочитал сообщения уже с новой схемой:
|
||||||
|
|
||||||
```bash
|
```bash
|
||||||
LIMIT=50 RESET_TOPICS=1 make data
|
make generated-history-analytics
|
||||||
|
make up
|
||||||
```
|
```
|
||||||
|
|
||||||
**Смотрим результат:**
|
**Смотрим результат:**
|
||||||
@@ -233,7 +235,8 @@ ALTER TABLE stg.browser_raw DROP COLUMN kafka_msg_ts;
|
|||||||
make ddl # пересоздаёт эталонный MV из 10_stg.sql — схема снова как в репозитории
|
make ddl # пересоздаёт эталонный MV из 10_stg.sql — схема снова как в репозитории
|
||||||
```
|
```
|
||||||
|
|
||||||
Если запутался в состоянии — всегда есть полный сброс: `make clean && make up && make ddl && LIMIT=50 make data`.
|
Если запутался в состоянии — всегда есть полный чистый прогон:
|
||||||
|
`make generated-history-analytics && make up`.
|
||||||
|
|
||||||
---
|
---
|
||||||
|
|
||||||
@@ -241,7 +244,7 @@ make ddl # пересоздаёт эталонный MV из 10_stg.sql —
|
|||||||
|
|
||||||
| Действие | Где смотреть | Что ожидать |
|
| Действие | Где смотреть | Что ожидать |
|
||||||
|----------|--------------|-------------|
|
|----------|--------------|-------------|
|
||||||
| `LIMIT=50 make data` | `SELECT count() FROM stg.browser_raw` | счётчик > 0 и растёт после загрузки |
|
| `make generated-history-analytics && make up` | `SELECT count() FROM stg.browser_raw` | счётчик > 0 |
|
||||||
| глянуть строку | `SELECT raw FROM stg.browser_raw LIMIT 1` | валидный JSON целиком, неразобранный |
|
| глянуть строку | `SELECT raw FROM stg.browser_raw LIMIT 1` | валидный JSON целиком, неразобранный |
|
||||||
| глянуть offset'ы | `SELECT kafka_offset FROM stg.browser_raw ORDER BY kafka_offset` | идут по возрастанию, без дублей |
|
| глянуть offset'ы | `SELECT kafka_offset FROM stg.browser_raw ORDER BY kafka_offset` | идут по возрастанию, без дублей |
|
||||||
| правка из секции 4 | `SELECT kafka_msg_ts FROM stg.browser_raw LIMIT 5` | колонка заполнена временем сообщения |
|
| правка из секции 4 | `SELECT kafka_msg_ts FROM stg.browser_raw LIMIT 5` | колонка заполнена временем сообщения |
|
||||||
|
|||||||
@@ -84,27 +84,27 @@ ODS мы пересобираем целиком, одной задачей Airf
|
|||||||
|
|
||||||
## 2. Руки: смотрим базовый прогон
|
## 2. Руки: смотрим базовый прогон
|
||||||
|
|
||||||
Поднимаем стенд, создаём схему, заливаем **малый срез** (50 строк на топик) и запускаем
|
Поднимаем стенд и создаём стартовую историю. Это штатный путь курса: готовый источник
|
||||||
трансформацию:
|
данных стенда пишет события в Kafka, ClickHouse читает их в STG, затем batch строит
|
||||||
|
ODS, DDS и DM. Файлы `data/*.jsonl` пока остаются только кладовкой значений для этого
|
||||||
|
источника, а не источником аналитического контура.
|
||||||
|
|
||||||
```bash
|
```bash
|
||||||
make up # поднять инфраструктуру
|
make generated-history-analytics
|
||||||
make ddl # создать базы и таблицы (в т.ч. слой ODS)
|
make up
|
||||||
LIMIT=50 make data # залить по 50 строк каждого файла в Kafka → STG
|
|
||||||
make transform # батч STG → ODS → DDS → DM (нас интересует первый шаг)
|
|
||||||
```
|
```
|
||||||
|
|
||||||
`make transform` прогоняет всю цепочку слоёв сразу, но прямо в консоли печатает то, что нам
|
Команда прогоняет всю цепочку слоёв и прямо в консоли печатает то, что нам нужно сейчас, —
|
||||||
нужно сейчас, — блок **«Статистика ODS»**. Это просто счётчики строк по всем восьми таблицам
|
блок **«Статистика ODS»**. Это просто счётчики строк по всем восьми таблицам слоя (четыре
|
||||||
слоя (четыре основных и четыре с ошибками):
|
основных и четыре с ошибками). Пример формы вывода:
|
||||||
|
|
||||||
```
|
```
|
||||||
Статистика ODS:
|
Статистика ODS:
|
||||||
┌─table──────────────────────┬─rows─┐
|
┌─table──────────────────────┬─rows─┐
|
||||||
│ ods.browser_event │ 50 │
|
│ ods.browser_event │ ... │
|
||||||
│ ods.location_event │ 50 │
|
│ ods.location_event │ ... │
|
||||||
│ ods.device_by_click │ 26 │
|
│ ods.device_by_click │ ... │
|
||||||
│ ods.geo_by_click │ 26 │
|
│ ods.geo_by_click │ ... │
|
||||||
│ ods.browser_event_errors │ 0 │
|
│ ods.browser_event_errors │ 0 │
|
||||||
│ ods.location_event_errors │ 0 │
|
│ ods.location_event_errors │ 0 │
|
||||||
│ ods.device_by_click_errors │ 0 │
|
│ ods.device_by_click_errors │ 0 │
|
||||||
@@ -114,35 +114,34 @@ make transform # батч STG → ODS → DDS → DM (нас инте
|
|||||||
|
|
||||||
Прочитаем эту табличку — в ней три вещи, которые стоит заметить.
|
Прочитаем эту табличку — в ней три вещи, которые стоит заметить.
|
||||||
|
|
||||||
**Все четыре `*_errors` — по нулям.** Значит, наш срез чистый: ни одна запись не дала ошибки
|
**Все четыре `*_errors` — по нулям.** Значит, стартовая история чистая: ни одна запись не дала
|
||||||
разбора, столбец `parse_errors` у всех пустой. Это нормально — данные в демо аккуратные.
|
ошибки разбора, столбец `parse_errors` у всех пустой. Это нормально — данные стенда аккуратные.
|
||||||
Ошибки мы увидим в секции 4, когда сами их устроим.
|
Ошибки мы увидим в секции 4, когда сами их устроим.
|
||||||
|
|
||||||
**`browser` и `location` дали 50 из 50.** Сколько событий пришло — столько и легло, один к
|
**`browser` и `location` идут в одном зерне события.** Сколько событий пришло, столько строк
|
||||||
одному.
|
и ожидаем увидеть после типизации, если ключи валидны.
|
||||||
|
|
||||||
**А `device` и `geo` — только 26 из 50.** Вот это уже интересно. Половина куда-то делась? Нет.
|
**А `device` и `geo` обычно меньше, чем событий.** Вот это уже интересно. Часть строк
|
||||||
И это важно понять, иначе дальше будет казаться, что данные текут.
|
куда-то делась? Нет. И это важно понять, иначе дальше будет казаться, что данные текут.
|
||||||
|
|
||||||
Дело в том, что эти две таблицы хранят не события, а **контекст клика**: с какого устройства
|
Дело в том, что эти две таблицы хранят не события, а **контекст клика**: с какого устройства
|
||||||
был клик и из какой точки на карте. Ключ у них — `click_id`. А в срезе на 50 событий разных
|
был клик и из какой точки на карте. Ключ у них — `click_id`. Разных кликов меньше, чем событий:
|
||||||
кликов всего 26: на один клик приходится несколько событий, и `click_id` у них повторяется.
|
на один клик приходится несколько событий, и `click_id` у них повторяется. Движок таблицы
|
||||||
Движок таблицы (про него — в секции 3) схлопывает повторы по ключу, оставляя по одной строке
|
(про него — в секции 3) схлопывает повторы по ключу, оставляя по одной строке на клик.
|
||||||
на клик. Отсюда и 26.
|
|
||||||
|
|
||||||
Проверь это сам, а не верь на слово. Открой SQL-консоль `http://localhost:9123/play`
|
Проверь это сам, а не верь на слово. Открой SQL-консоль `http://localhost:9123/play`
|
||||||
(пользователь `default`, пароль `123456`) и посчитай, сколько в срезе *различных* `click_id`:
|
(пользователь `default`, пароль `123456`) и посчитай, сколько в STG *различных* `click_id`:
|
||||||
|
|
||||||
```sql
|
```sql
|
||||||
-- Всего строк в STG — 50, но различных click_id среди них — ровно 26
|
-- Строк событий больше, чем различных click_id
|
||||||
SELECT count() AS stg_rows,
|
SELECT count() AS stg_rows,
|
||||||
uniqExact(toUUIDOrNull(JSONExtractString(raw, 'click_id'))) AS distinct_clicks
|
uniqExact(toUUIDOrNull(JSONExtractString(raw, 'click_id'))) AS distinct_clicks
|
||||||
FROM stg.geo_raw;
|
FROM stg.geo_raw;
|
||||||
```
|
```
|
||||||
|
|
||||||
Получишь `stg_rows = 50`, `distinct_clicks = 26` — ровно столько, сколько строк в
|
`distinct_clicks` должен быть меньше или равен `stg_rows` и совпадать с числом строк в
|
||||||
`ods.geo_by_click`. Значит, 26 — это схлопнутые повторы, а не пропавшие данные. Ничего не
|
`ods.geo_by_click`. Значит, это схлопнутые повторы по `click_id`, а не пропавшие данные.
|
||||||
потерялось молча.
|
Ничего не потерялось молча.
|
||||||
|
|
||||||
---
|
---
|
||||||
|
|
||||||
@@ -174,7 +173,7 @@ toFloat64OrNull(JSONExtractString(raw, 'geo_latitude')) AS
|
|||||||
Весь смысл — в суффиксе `OrNull`. Если значение **не** приводится к нужному типу (вместо UUID
|
Весь смысл — в суффиксе `OrNull`. Если значение **не** приводится к нужному типу (вместо UUID
|
||||||
пришёл мусор), функция не падает с ошибкой, а просто возвращает `NULL`. Это ровно то правило
|
пришёл мусор), функция не падает с ошибкой, а просто возвращает `NULL`. Это ровно то правило
|
||||||
стенда, что и в STG — «грязная запись не валит пайплайн», — только теперь на уровне типов.
|
стенда, что и в STG — «грязная запись не валит пайплайн», — только теперь на уровне типов.
|
||||||
Один кривой `event_id` станет `NULL` и будет помечен, а остальные 49 строк спокойно доедут.
|
Один кривой `event_id` станет `NULL` и будет помечен, а остальные строки спокойно доедут.
|
||||||
|
|
||||||
> Кстати, про `AS`: эти строки живут в блоке `WITH` в начале запроса. `WITH` — это просто
|
> Кстати, про `AS`: эти строки живут в блоке `WITH` в начале запроса. `WITH` — это просто
|
||||||
> способ заранее посчитать значение и дать ему имя, чтобы ниже по запросу ссылаться на него
|
> способ заранее посчитать значение и дать ему имя, чтобы ниже по запросу ссылаться на него
|
||||||
@@ -230,7 +229,7 @@ arrayFilter(x -> x != '', [
|
|||||||
> если поменять разбор только в одном из двух мест, они разойдутся. В секции 4 мы как раз этим
|
> если поменять разбор только в одном из двух мест, они разойдутся. В секции 4 мы как раз этим
|
||||||
> воспользуемся — и увидим, чем грозит такой рассинхрон.
|
> воспользуемся — и увидим, чем грозит такой рассинхрон.
|
||||||
|
|
||||||
### Движок: откуда взялись 26 строк
|
### Движок: почему строк контекста меньше
|
||||||
|
|
||||||
И последнее место — строчка про движок основных таблиц:
|
И последнее место — строчка про движок основных таблиц:
|
||||||
|
|
||||||
@@ -241,8 +240,8 @@ ORDER BY (click_id)
|
|||||||
|
|
||||||
`ReplacingMergeTree` — это таблица, которая схлопывает строки с одинаковым ключом (ключ берётся
|
`ReplacingMergeTree` — это таблица, которая схлопывает строки с одинаковым ключом (ключ берётся
|
||||||
из `ORDER BY`), оставляя самую свежую по `src_ingest_ts` — времени загрузки в ODS. Вот она,
|
из `ORDER BY`), оставляя самую свежую по `src_ingest_ts` — времени загрузки в ODS. Вот она,
|
||||||
причина «26 из 50» из секции 2: у `device` и `geo` много строк с одинаковым `click_id`, и
|
причина разницы из секции 2: у `device` и `geo` много строк с одинаковым `click_id`, и движок
|
||||||
движок оставляет по одной на клик.
|
оставляет по одной на клик.
|
||||||
|
|
||||||
---
|
---
|
||||||
|
|
||||||
@@ -276,8 +275,8 @@ make transform
|
|||||||
И смотрим на ту же «Статистику ODS». Таблица ошибок гео, которая была пустой, теперь полная:
|
И смотрим на ту же «Статистику ODS». Таблица ошибок гео, которая была пустой, теперь полная:
|
||||||
|
|
||||||
```
|
```
|
||||||
│ ods.geo_by_click │ 26 │
|
│ ods.geo_by_click │ ... │
|
||||||
│ ods.geo_by_click_errors │ 50 │ ← было 0
|
│ ods.geo_by_click_errors │ ... │ ← было 0
|
||||||
```
|
```
|
||||||
|
|
||||||
А в самой основной таблице широта пропала — но не молча, рядом стоит метка:
|
А в самой основной таблице широта пропала — но не молча, рядом стоит метка:
|
||||||
@@ -295,11 +294,10 @@ LIMIT 4;
|
|||||||
└──────────────┴──────────────┴───────────────┴──────────────────────┘
|
└──────────────┴──────────────┴───────────────┴──────────────────────┘
|
||||||
```
|
```
|
||||||
|
|
||||||
Вот теперь видно всё разом — и DQ-split, и «двойной учёт» из секции 3 вживую. 26 строк
|
Вот теперь видно всё разом — и DQ-split, и «двойной учёт» из секции 3 вживую. Строки с валидным
|
||||||
остались в основной таблице (ключ `click_id` цел) с пометкой `bad_geo_latitude`. И те же
|
`click_id` остались в основной таблице с пометкой `bad_geo_latitude`. И те же записи попали в
|
||||||
записи попали в число 50 строк `geo_by_click_errors`. Долгота на месте, а широты больше нет:
|
`geo_by_click_errors`. Долгота на месте, а широты больше нет: один неверный тип — и целое поле
|
||||||
один неверный тип — и целое поле потеряно по всему слою. Заметили это `parse_errors` и таблица
|
потеряно по всему слою. Заметили это `parse_errors` и таблица ошибок — для того DQ-split и нужен.
|
||||||
ошибок — для того DQ-split и нужен.
|
|
||||||
|
|
||||||
> **Бывает и хуже — тихо, совсем без метки.** Здесь нас спас суффикс `OrNull`: неверный тип
|
> **Бывает и хуже — тихо, совсем без метки.** Здесь нас спас суффикс `OrNull`: неверный тип
|
||||||
> дал `NULL`, а `NULL` мы умеем замечать (на него и сработал `parse_errors`). По-настоящему
|
> дал `NULL`, а `NULL` мы умеем замечать (на него и сработал `parse_errors`). По-настоящему
|
||||||
@@ -319,7 +317,7 @@ make transform
|
|||||||
```
|
```
|
||||||
|
|
||||||
После этого `geo_by_click_errors` снова `0`, широта на месте. А если стенд совсем «поплыл» —
|
После этого `geo_by_click_errors` снова `0`, широта на месте. А если стенд совсем «поплыл» —
|
||||||
всегда есть полный сброс: `make clean && make up && make ddl && LIMIT=50 make data && make transform`.
|
всегда есть полный чистый прогон: `make generated-history-analytics && make up`.
|
||||||
|
|
||||||
---
|
---
|
||||||
|
|
||||||
@@ -327,9 +325,9 @@ make transform
|
|||||||
|
|
||||||
| Действие | Где смотреть | Что ожидать |
|
| Действие | Где смотреть | Что ожидать |
|
||||||
|----------|--------------|-------------|
|
|----------|--------------|-------------|
|
||||||
| `make transform` (базовый прогон) | блок «Статистика ODS» | `browser`/`location` = 50, `device`/`geo` = 26, все `*_errors` = 0 |
|
| базовый прогон | блок «Статистика ODS» | основные таблицы не пустые, все `*_errors` = 0 |
|
||||||
| почему 26, а не 50 | запрос `uniqExact(click_id)` по `stg.geo_raw` | 26 различных `click_id` — это схлопывание повторов, а не потеря |
|
| почему `device`/`geo` меньше событий | запрос `uniqExact(click_id)` по `stg.geo_raw` | число различных `click_id` совпадает с `ods.geo_by_click` |
|
||||||
| правка из секции 4 | блок «Статистика ODS» | `ods.geo_by_click_errors` прыгнул `0 → 50` |
|
| правка из секции 4 | блок «Статистика ODS» | `ods.geo_by_click_errors` прыгнул с `0` на ненулевое число |
|
||||||
| правка из секции 4 | `SELECT geo_latitude, parse_errors FROM ods.geo_by_click` | широта `NULL`, в `parse_errors` — `bad_geo_latitude` |
|
| правка из секции 4 | `SELECT geo_latitude, parse_errors FROM ods.geo_by_click` | широта `NULL`, в `parse_errors` — `bad_geo_latitude` |
|
||||||
|
|
||||||
---
|
---
|
||||||
@@ -338,7 +336,8 @@ make transform
|
|||||||
|
|
||||||
После урока у тебя на руках — видимый результат (одно на выбор):
|
После урока у тебя на руках — видимый результат (одно на выбор):
|
||||||
|
|
||||||
- скрин блока «Статистика ODS», где после правки `ods.geo_by_click_errors` ушёл с `0` на `50`;
|
- скрин блока «Статистика ODS», где после правки `ods.geo_by_click_errors` ушёл с `0`
|
||||||
|
на ненулевое число;
|
||||||
- либо выборка из `ods.geo_by_click` с пустой широтой и меткой `bad_geo_latitude` рядом.
|
- либо выборка из `ods.geo_by_click` с пустой широтой и меткой `bad_geo_latitude` рядом.
|
||||||
|
|
||||||
И проверь себя на словах — примерно эти вопросы всплывут на еженедельном созвоне:
|
И проверь себя на словах — примерно эти вопросы всплывут на еженедельном созвоне:
|
||||||
|
|||||||
@@ -65,35 +65,35 @@
|
|||||||
|
|
||||||
## 2. Руки: смотрим базовый прогон
|
## 2. Руки: смотрим базовый прогон
|
||||||
|
|
||||||
Поднимаем стенд, создаём схему, заливаем **малый срез** (50 строк на топик) и запускаем
|
Поднимаем стенд и создаём стартовую историю. Это штатный путь курса: готовый источник
|
||||||
трансформацию:
|
данных стенда пишет события в Kafka, ClickHouse читает их в STG, затем batch строит
|
||||||
|
ODS, DDS и DM. Файлы `data/*.jsonl` пока остаются только кладовкой значений для этого
|
||||||
|
источника, а не источником аналитического контура.
|
||||||
|
|
||||||
```bash
|
```bash
|
||||||
make up # поднять инфраструктуру
|
make generated-history-analytics
|
||||||
make ddl # создать базы и таблицы (в т.ч. слой DDS)
|
make up
|
||||||
LIMIT=50 make data # залить по 50 строк каждого файла в Kafka → STG
|
|
||||||
make transform # батч STG → ODS → DDS → DM
|
|
||||||
```
|
```
|
||||||
|
|
||||||
`make transform` прогоняет всю цепочку слоёв и по дороге печатает в консоль блок **«Статистика
|
Команда прогоняет всю цепочку слоёв и по дороге печатает в консоль блок **«Статистика DDS»** —
|
||||||
DDS»** — счётчики строк по двум нашим сущностям:
|
счётчики строк по двум нашим сущностям:
|
||||||
|
|
||||||
```
|
```
|
||||||
Статистика DDS:
|
Статистика DDS:
|
||||||
┌─table─────┬─rows─┐
|
┌─table─────┬─rows─┐
|
||||||
│ dds.click │ 26 │
|
│ dds.click │ ... │
|
||||||
│ dds.event │ 50 │
|
│ dds.event │ ... │
|
||||||
└───────────┴──────┘
|
└───────────┴──────┘
|
||||||
```
|
```
|
||||||
|
|
||||||
Прочитаем эти две строки.
|
Прочитаем эти две строки.
|
||||||
|
|
||||||
**`dds.event` — 50.** Сколько событий пришло, столько карточек и собралось: одно событие — одна
|
**`dds.event` — события.** Сколько событий пришло с валидным ключом, столько карточек и
|
||||||
строка. Ровно как `ods.browser_event` из прошлого урока.
|
собралось: одно событие — одна строка. Ровно как `ods.browser_event` из прошлого урока.
|
||||||
|
|
||||||
**`dds.click` — 26, а не 50.** И это та же история, что мы уже разбирали в уроке 2. Карточка
|
**`dds.click` обычно меньше, чем `dds.event`.** И это та же история, что мы уже разбирали
|
||||||
клика — одна на клик, а в срезе на 50 событий разных кликов всего 26 (на один клик приходится
|
в уроке 2. Карточка клика — одна на клик, а событий на один клик может быть несколько.
|
||||||
несколько событий). Поэтому 50 событий ссылаются на 26 кликов — это нормально, так и должно быть.
|
Поэтому много событий ссылаются на меньшее число кликов — это нормально, так и должно быть.
|
||||||
|
|
||||||
Теперь — главный счётчик урока. Он печатается чуть ниже, в блоке **«Сводка по качеству данных»**
|
Теперь — главный счётчик урока. Он печатается чуть ниже, в блоке **«Сводка по качеству данных»**
|
||||||
(это таблица `dm.dq_summary`, куда стенд складывает метрики по всем слоям). Найди в ней строку
|
(это таблица `dm.dq_summary`, куда стенд складывает метрики по всем слоям). Найди в ней строку
|
||||||
@@ -108,9 +108,9 @@ DDS»** — счётчики строк по двум нашим сущност
|
|||||||
В колонке `check_date` стоит `today()` из кода витрины, так что у тебя там будет сегодняшняя
|
В колонке `check_date` стоит `today()` из кода витрины, так что у тебя там будет сегодняшняя
|
||||||
дата — не пугайся, если она не совпадёт с примером.
|
дата — не пугайся, если она не совпадёт с примером.
|
||||||
|
|
||||||
`orphan_events = 0` — ни одной сироты. Каждое из 50 событий нашло свой клик в `dds.click`. На
|
`orphan_events = 0` — ни одной сироты. Каждое событие нашло свой клик в `dds.click`. На
|
||||||
чистом демо-срезе так и должно быть: данные аккуратные, ничего не потерялось. В секции 4 мы
|
чистой стартовой истории так и должно быть: данные аккуратные, ничего не потерялось. В секции
|
||||||
сироту устроим сами — и эта строка оживёт.
|
4 мы сироту устроим сами — и эта строка оживёт.
|
||||||
|
|
||||||
Проверь нолик сам, не верь на слово. Открой SQL-консоль `http://localhost:9123/play`
|
Проверь нолик сам, не верь на слово. Открой SQL-консоль `http://localhost:9123/play`
|
||||||
(пользователь `default`, пароль `123456`) и посчитай сирот напрямую:
|
(пользователь `default`, пароль `123456`) и посчитай сирот напрямую:
|
||||||
@@ -158,9 +158,9 @@ SELECT click_id FROM ods.geo_by_click ...
|
|||||||
строим карточки: так не потеряется клик, который есть, например, в `geo`, но почему-то не доехал
|
строим карточки: так не потеряется клик, который есть, например, в `geo`, но почему-то не доехал
|
||||||
в `device`.
|
в `device`.
|
||||||
|
|
||||||
> На нашем срезе `device` и `geo` содержат один и тот же набор из 26 кликов, так что универсум
|
> На чистой стартовой истории `device` и `geo` должны содержать один и тот же набор кликов,
|
||||||
> тоже 26. Но код написан так, чтобы пережить случай, когда наборы **разойдутся**, — и это
|
> так что универсум совпадает с обоими источниками. Но код написан так, чтобы пережить случай,
|
||||||
> правильно: в проде они расходятся постоянно.
|
> когда наборы **разойдутся**, — и это правильно: в проде они расходятся постоянно.
|
||||||
|
|
||||||
### `argMax`: одна строка на клик, самая свежая
|
### `argMax`: одна строка на клик, самая свежая
|
||||||
|
|
||||||
@@ -221,7 +221,7 @@ LEFT JOIN ( ...снапшот geo... ) AS g ON g.click_id = c.click_id
|
|||||||
|
|
||||||
### Сироты: событие без клика
|
### Сироты: событие без клика
|
||||||
|
|
||||||
Мы собрали `dds.click` (26 карточек кликов) и `dds.event` (50 карточек событий). Внутри каждого
|
Мы собрали `dds.click` (карточки кликов) и `dds.event` (карточки событий). Внутри каждого
|
||||||
события лежит `click_id` — ссылка на клик. И вот тут возникает вопрос целостности из секции 1:
|
события лежит `click_id` — ссылка на клик. И вот тут возникает вопрос целостности из секции 1:
|
||||||
**а на каждую ли ссылку есть карточка клика?**
|
**а на каждую ли ссылку есть карточка клика?**
|
||||||
|
|
||||||
@@ -237,16 +237,16 @@ WHERE click_id IS NOT NULL
|
|||||||
Заметь разницу с предыдущим пунктом. Пустое гео — это когда у **клика** не подтянулся свой
|
Заметь разницу с предыдущим пунктом. Пустое гео — это когда у **клика** не подтянулся свой
|
||||||
контекст (внутренний пропуск в карточке, но сам клик есть). А сирота — это когда у **события**
|
контекст (внутренний пропуск в карточке, но сам клик есть). А сирота — это когда у **события**
|
||||||
нет вообще никакого клика (порвана связь между сущностями). Это разные дырки: первую видно по
|
нет вообще никакого клика (порвана связь между сущностями). Это разные дырки: первую видно по
|
||||||
пустым полям внутри карточки, вторую — отдельным счётчиком. На чистом срезе сирот ноль — сейчас
|
пустым полям внутри карточки, вторую — отдельным счётчиком. На чистой стартовой истории сирот
|
||||||
мы это изменим.
|
ноль — сейчас мы это изменим.
|
||||||
|
|
||||||
---
|
---
|
||||||
|
|
||||||
## 4. Управляемая правка: заведём сироту
|
## 4. Управляемая правка: заведём сироту
|
||||||
|
|
||||||
Сирота на чистом срезе не появится сама — данные слишком аккуратные. Поэтому **создадим её
|
Сирота на чистой стартовой истории не появится сама — данные слишком аккуратные. Поэтому
|
||||||
руками**: добавим в `dds.event` одно событие, которое ссылается на клик, которого в `dds.click`
|
**создадим её руками**: добавим в `dds.event` одно событие, которое ссылается на клик, которого
|
||||||
нет. И посмотрим, как оживёт счётчик сирот и как себя поведёт `LEFT JOIN`.
|
в `dds.click` нет. И посмотрим, как оживёт счётчик сирот и как себя поведёт `LEFT JOIN`.
|
||||||
|
|
||||||
Открой SQL-консоль `http://localhost:9123/play` и вставь придуманное событие:
|
Открой SQL-консоль `http://localhost:9123/play` и вставь придуманное событие:
|
||||||
|
|
||||||
@@ -309,7 +309,7 @@ make transform
|
|||||||
```
|
```
|
||||||
|
|
||||||
После этого `orphan_events` снова `0`, придуманное событие исчезло. А если стенд совсем «поплыл» —
|
После этого `orphan_events` снова `0`, придуманное событие исчезло. А если стенд совсем «поплыл» —
|
||||||
полный сброс: `make clean && make up && make ddl && LIMIT=50 make data && make transform`.
|
полный чистый прогон: `make generated-history-analytics && make up`.
|
||||||
|
|
||||||
---
|
---
|
||||||
|
|
||||||
@@ -317,9 +317,9 @@ make transform
|
|||||||
|
|
||||||
| Действие | Где смотреть | Что ожидать |
|
| Действие | Где смотреть | Что ожидать |
|
||||||
|----------|--------------|-------------|
|
|----------|--------------|-------------|
|
||||||
| `make transform` (базовый прогон) | блок «Статистика DDS» | `dds.click` = 26, `dds.event` = 50 |
|
| базовый прогон | блок «Статистика DDS» | `dds.click` и `dds.event` не пустые |
|
||||||
| `make transform` (базовый прогон) | блок «Сводка по качеству», строка `orphan_events` | `0` |
|
| `make transform` (базовый прогон) | блок «Сводка по качеству», строка `orphan_events` | `0` |
|
||||||
| почему `click` = 26, а `event` = 50 | запрос `count()` по `dds.click` и `dds.event` | 50 событий ссылаются на 26 кликов — норма |
|
| почему `click` меньше `event` | запрос `count()` по `dds.click` и `dds.event` | много событий ссылаются на меньшее число кликов — норма |
|
||||||
| правка из секции 4 (вставили сироту) | запрос `count()` сирот в play-консоли | `0 → 1` |
|
| правка из секции 4 (вставили сироту) | запрос `count()` сирот в play-консоли | `0 → 1` |
|
||||||
| та же сирота через `dm.v_events_enriched` | `SELECT device_type, geo_country ...` | поля клика пустые (`NULL`) — это `LEFT JOIN` |
|
| та же сирота через `dm.v_events_enriched` | `SELECT device_type, geo_country ...` | поля клика пустые (`NULL`) — это `LEFT JOIN` |
|
||||||
|
|
||||||
|
|||||||
@@ -56,12 +56,14 @@ Airflow. Главная единица Airflow — **DAG** (Directed Acyclic Gra
|
|||||||
|
|
||||||
## 2. Руки: запускаем DAG и смотрим зелёный прогон
|
## 2. Руки: запускаем DAG и смотрим зелёный прогон
|
||||||
|
|
||||||
Подними стенд, создай схему и залей малый срез:
|
Подними стенд и создай стартовую историю. Это штатный путь курса: готовый источник
|
||||||
|
данных стенда пишет события в Kafka, ClickHouse читает их в STG, затем batch строит
|
||||||
|
ODS, DDS и DM. Файлы `data/*.jsonl` пока остаются только кладовкой значений для этого
|
||||||
|
источника, а не источником аналитического контура.
|
||||||
|
|
||||||
```bash
|
```bash
|
||||||
|
make generated-history-analytics
|
||||||
make up
|
make up
|
||||||
make ddl
|
|
||||||
LIMIT=50 make data
|
|
||||||
```
|
```
|
||||||
|
|
||||||
Открой Airflow: `http://localhost:8080` (логин `admin`, пароль `admin`). Найди DAG
|
Открой Airflow: `http://localhost:8080` (логин `admin`, пароль `admin`). Найди DAG
|
||||||
@@ -75,13 +77,13 @@ LIMIT=50 make data
|
|||||||
заново наполнит их из ODS. Для учебного стенда это удобный чистый прогон: результат повторяемый,
|
заново наполнит их из ODS. Для учебного стенда это удобный чистый прогон: результат повторяемый,
|
||||||
старые эксперименты не мешают.
|
старые эксперименты не мешают.
|
||||||
|
|
||||||
Когда DAG завершится, открой его граф. На чистом срезе все задачи должны быть зелёными. Найди
|
Когда DAG завершится, открой его граф. На чистой стартовой истории все задачи должны быть
|
||||||
внутри группы `transform` две задачи подряд:
|
зелёными. Найди внутри группы `transform` две задачи подряд:
|
||||||
|
|
||||||
- `check_dds_integrity` — SQL-задача, которая считает сирот;
|
- `check_dds_integrity` — SQL-задача, которая считает сирот;
|
||||||
- `assert_dds_integrity` — Python-задача, которая решает, можно ли идти дальше.
|
- `assert_dds_integrity` — Python-задача, которая решает, можно ли идти дальше.
|
||||||
|
|
||||||
На чистом срезе `assert_dds_integrity` зелёная: сирот нет, пайплайн прошёл в DM.
|
На чистой стартовой истории `assert_dds_integrity` зелёная: сирот нет, пайплайн прошёл в DM.
|
||||||
|
|
||||||
Проверь то же число в ClickHouse play-консоли `http://localhost:9123/play`:
|
Проверь то же число в ClickHouse play-консоли `http://localhost:9123/play`:
|
||||||
|
|
||||||
@@ -305,14 +307,15 @@ WHERE click_id IS NOT NULL
|
|||||||
AND click_id NOT IN (SELECT click_id FROM dds.click);
|
AND click_id NOT IN (SELECT click_id FROM dds.click);
|
||||||
```
|
```
|
||||||
|
|
||||||
Снова должно быть `0`. Если стенд после экспериментов совсем запутался, полный сброс остаётся
|
Снова должно быть `0`. Если стенд после экспериментов совсем запутался, сделай штатный
|
||||||
тем же:
|
чистый прогон:
|
||||||
|
|
||||||
```bash
|
```bash
|
||||||
make clean && make up && make ddl && LIMIT=50 make data
|
make generated-history-analytics
|
||||||
|
make up
|
||||||
```
|
```
|
||||||
|
|
||||||
После этого запусти `etl_pipeline` с `{"full_refresh": true}`.
|
После этого при необходимости запусти `etl_pipeline` с `{"full_refresh": true}`.
|
||||||
|
|
||||||
---
|
---
|
||||||
|
|
||||||
|
|||||||
@@ -57,15 +57,18 @@
|
|||||||
|
|
||||||
## 2. Руки: открываем дашборды и targets
|
## 2. Руки: открываем дашборды и targets
|
||||||
|
|
||||||
Подними стенд и прогони маленький срез, если он ещё не поднят:
|
Подними стенд и создай стартовую историю, если он ещё не поднят. Это штатный путь курса:
|
||||||
|
готовый источник данных стенда пишет события в Kafka, ClickHouse читает их в STG, затем
|
||||||
|
batch строит ODS, DDS и DM. Файлы `data/*.jsonl` пока остаются только кладовкой значений
|
||||||
|
для этого источника, а не источником аналитического контура.
|
||||||
|
|
||||||
```bash
|
```bash
|
||||||
|
make generated-history-analytics
|
||||||
make up
|
make up
|
||||||
make ddl
|
|
||||||
LIMIT=50 make data
|
|
||||||
```
|
```
|
||||||
|
|
||||||
После этого запусти `etl_pipeline` в Airflow с конфигом:
|
Если хочешь отдельно посмотреть Airflow-run после готовой истории, запусти `etl_pipeline`
|
||||||
|
в Airflow с конфигом:
|
||||||
|
|
||||||
```json
|
```json
|
||||||
{"full_refresh": true}
|
{"full_refresh": true}
|
||||||
|
|||||||
@@ -63,42 +63,28 @@ DM-витрина — это SQL-объект в ClickHouse. Она задаёт
|
|||||||
|
|
||||||
## 2. Руки: запускаем Superset и смотрим дашборд
|
## 2. Руки: запускаем Superset и смотрим дашборд
|
||||||
|
|
||||||
Подними стенд и прогони полный демо-датасет:
|
Подними стенд и создай стартовую историю. Это штатный путь курса: готовый источник
|
||||||
|
данных стенда пишет события в Kafka, ClickHouse читает их в STG, затем batch строит
|
||||||
|
ODS, DDS, DM и обновляет Superset. Файлы `data/*.jsonl` пока остаются только кладовкой
|
||||||
|
значений для этого источника, а не источником аналитического контура.
|
||||||
|
|
||||||
```bash
|
```bash
|
||||||
|
make generated-history-analytics
|
||||||
make up
|
make up
|
||||||
make ddl
|
|
||||||
make data
|
|
||||||
make transform
|
|
||||||
```
|
```
|
||||||
|
|
||||||
`make transform` прогоняет цепочку STG → ODS → DDS → DM вне Airflow. Для этого урока так
|
Команда уже прогоняет цепочку STG → ODS → DDS → DM и создаёт metadata Superset:
|
||||||
быстрее: нам нужен готовый DM-слой, а не разбор DAG. Отладочный срез через
|
|
||||||
`LIMIT=50 make data` можно использовать для быстрых экспериментов, но эталонный dashboard
|
|
||||||
и числа урока рассчитаны на полном наборе данных.
|
|
||||||
|
|
||||||
Теперь инициализируй Superset:
|
|
||||||
|
|
||||||
```bash
|
|
||||||
make superset-init
|
|
||||||
```
|
|
||||||
|
|
||||||
Эта команда создаёт или обновляет:
|
|
||||||
|
|
||||||
- подключение `clickhouse_dwh`;
|
- подключение `clickhouse_dwh`;
|
||||||
- 6 datasets поверх `dm.*`.
|
- 6 datasets поверх `dm.*`.
|
||||||
|
|
||||||
После этого создай или обнови charts и сам dashboard:
|
|
||||||
|
|
||||||
```bash
|
|
||||||
make superset-dashboard
|
|
||||||
```
|
|
||||||
|
|
||||||
Эта команда создаёт или обновляет:
|
|
||||||
|
|
||||||
- 10 charts;
|
- 10 charts;
|
||||||
- dashboard `E-commerce Analytics Dashboard`.
|
- dashboard `E-commerce Analytics Dashboard`.
|
||||||
|
|
||||||
|
Для этого урока нам нужен готовый DM-слой и уже созданный dashboard, а не разбор DAG.
|
||||||
|
Если Superset metadata нужно пересоздать отдельно, используй `make superset-init`.
|
||||||
|
Если нужно обновить только charts и dashboard после правки в `superset/create_dashboard.py`,
|
||||||
|
используй `make superset-dashboard`.
|
||||||
|
|
||||||
Открой Superset: `http://localhost:8088` (логин `admin`, пароль `admin`).
|
Открой Superset: `http://localhost:8088` (логин `admin`, пароль `admin`).
|
||||||
|
|
||||||
Если Superset предлагает сменить пароль после первого входа, для учебного стенда можно нажать
|
Если Superset предлагает сменить пароль после первого входа, для учебного стенда можно нажать
|
||||||
@@ -133,33 +119,33 @@ http://localhost:8088/superset/dashboard/ecommerce-analytics/
|
|||||||
> **Зерно (grain), он же уровень гранулярности — это что считается одной строкой
|
> **Зерно (grain), он же уровень гранулярности — это что считается одной строкой
|
||||||
> таблицы.** У события зерно
|
> таблицы.** У события зерно
|
||||||
> «одно событие = одна строка» (ключ `event_id`), у визита — «один визит = одна
|
> «одно событие = одна строка» (ключ `event_id`), у визита — «один визит = одна
|
||||||
> строка» (ключ `click_id`). Это разные зёрна: событий 1000, а визитов 99, потому
|
> строка» (ключ `click_id`). Это разные зёрна: событий обычно больше, чем визитов,
|
||||||
> что в одном визите много событий. Складывать строки разного зерна в одно число
|
> потому что в одном визите много событий. Складывать строки разного зерна в одно
|
||||||
> бессмысленно — это всё равно что сложить «штуки яблок» и «корзины яблок».
|
> число бессмысленно — это всё равно что сложить «штуки яблок» и «корзины яблок».
|
||||||
|
|
||||||
Поэтому чарт держит **одно зерно — event**: берёт по одной канонической таблице
|
Поэтому чарт держит **одно зерно — event**: берёт по одной канонической таблице
|
||||||
событий на слой (`browser_raw → browser_event → event → v_events_enriched`), а не сумму
|
событий на слой (`browser_raw → browser_event → event → v_events_enriched`), а не сумму
|
||||||
по слою. Если просуммировать все таблицы слоя, в один столбец попадут таблицы разного
|
по слою. Если просуммировать все таблицы слоя, в один столбец попадут таблицы разного
|
||||||
зерна (события 1000 + визиты 99 + пустые error-таблицы) и получится «воронка потерь»,
|
зерна: события, визиты и технические таблицы ошибок. Получится «воронка потерь», которой
|
||||||
которой на самом деле нет.
|
на самом деле нет.
|
||||||
|
|
||||||
Шаг **1050 → 1000** на первом переходе — это не потеря данных, а дедупликация
|
Если на первом переходе число уменьшается, это не обязательно потеря данных. В ODS работает
|
||||||
at-least-once потока по `event_id` в ODS (`ReplacingMergeTree`): в STG приехало 1050 строк,
|
дедупликация at-least-once потока по `event_id` (`ReplacingMergeTree`): в STG могут приехать
|
||||||
но уникальных `event_id` среди них — 1000 (часть событий Kafka доставила повторно). Дальше
|
повторы, а дальше остаётся каноническое число уникальных событий. Настоящие проблемы качества
|
||||||
число стабильно. Настоящие проблемы качества (ошибки парсинга, осиротевшие события) на
|
(ошибки парсинга, осиротевшие события) на чистой стартовой истории равны нулю и лежат в
|
||||||
чистых демо-данных равны нулю и лежат в `dm.dq_summary` отдельными `check_name` — их
|
`dm.dq_summary` отдельными `check_name` — их разбирали уроки 3–4.
|
||||||
разбирали уроки 3–4.
|
|
||||||
|
|
||||||
### Фильтр даты
|
### Фильтр даты
|
||||||
|
|
||||||
В демо-данных события датированы `2022-11-28`. В текущей конфигурации dashboard фильтр даты
|
В стартовой истории события идут по модельному времени стенда. В текущей конфигурации
|
||||||
открывается как `No filter`. Если у тебя осталась старая metadata Superset и native filter
|
dashboard фильтр даты открывается как `No filter`. Если у тебя осталась старая metadata
|
||||||
**Date Range** стоит в значении `Last week`, часть графиков может быть пустой, хотя данные есть.
|
Superset и native filter **Date Range** стоит в значении `Last week`, часть графиков может
|
||||||
|
быть пустой, хотя данные есть.
|
||||||
|
|
||||||
Для этого урока поставь в фильтре даты одно из двух:
|
Для этого урока поставь в фильтре даты одно из двух:
|
||||||
|
|
||||||
- `No filter`;
|
- `No filter`;
|
||||||
- или ручной диапазон вокруг `2022-11-28`.
|
- или ручной диапазон вокруг дат, которые видны в `event_timestamp`.
|
||||||
|
|
||||||
После этого нажми **Apply filters**. Теперь смотри на dashboard как аналитик: какие графики
|
После этого нажми **Apply filters**. Теперь смотри на dashboard как аналитик: какие графики
|
||||||
отвечают на бизнес-вопросы, а какие только показывают техническое устройство конвейера.
|
отвечают на бизнес-вопросы, а какие только показывают техническое устройство конвейера.
|
||||||
@@ -415,7 +401,7 @@ metadata Superset.
|
|||||||
| `make superset-init` | Superset → **Settings → Database Connections** | есть подключение `clickhouse_dwh` |
|
| `make superset-init` | Superset → **Settings → Database Connections** | есть подключение `clickhouse_dwh` |
|
||||||
| открыть **Datasets** | Superset UI | есть datasets `v_events_enriched`, `v_top_pages_daily`, `dq_summary` |
|
| открыть **Datasets** | Superset UI | есть datasets `v_events_enriched`, `v_top_pages_daily`, `dq_summary` |
|
||||||
| открыть dashboard | Superset UI | видны KPI, маркетинг, география и прохождение строк по слоям |
|
| открыть dashboard | Superset UI | видны KPI, маркетинг, география и прохождение строк по слоям |
|
||||||
| поставить **Date Range → No filter** | dashboard filters | графики не скрываются из-за даты `2022-11-28` |
|
| поставить **Date Range → No filter** | dashboard filters | графики не скрываются из-за даты модельного времени |
|
||||||
| поменять `row_limit` у `Page Funnel` на `3` и запустить `make superset-dashboard` | chart `Page Funnel` | не больше трёх страниц |
|
| поменять `row_limit` у `Page Funnel` на `3` и запустить `make superset-dashboard` | chart `Page Funnel` | не больше трёх страниц |
|
||||||
| вернуть `row_limit` на `20` и запустить `make superset-dashboard` | chart `Page Funnel` | ограничение снова до 20 страниц |
|
| вернуть `row_limit` на `20` и запустить `make superset-dashboard` | chart `Page Funnel` | ограничение снова до 20 страниц |
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user