Merge pull request 'Добавить каноническое чтение актуальных событий через ods.event_v' (#97) from feat/86-event-v into main
Reviewed-on: #97
This commit was merged in pull request #97.
This commit is contained in:
@@ -8,9 +8,9 @@
|
|||||||
**Что здесь описано и чего ещё нет.** Собран этап 1: кластер из двух шардов,
|
**Что здесь описано и чего ещё нет.** Собран этап 1: кластер из двух шардов,
|
||||||
keeper, Kafka, каркас сервисов. Этап 2 идёт: в `sql/ddl/` лежит вся цепочка
|
keeper, Kafka, каркас сервисов. Этап 2 идёт: в `sql/ddl/` лежит вся цепочка
|
||||||
`Kafka → STG → ODS` — чтец топика `hits`, таблицы сырья, типизированное
|
`Kafka → STG → ODS` — чтец топика `hits`, таблицы сырья, типизированное
|
||||||
событие с таблицей ошибок и три матвью. Дальше по тексту устройство описано
|
событие с таблицей ошибок, поверхность актуального состояния и три матвью.
|
||||||
так, как оно проектируется; построенное от заложенного отличает карта таблиц в
|
Дальше по тексту устройство описано так, как оно проектируется; построенное от
|
||||||
конце.
|
заложенного отличает карта таблиц в конце.
|
||||||
|
|
||||||
Зона ответственности у документа одна — хранилище. Генератор описан отдельно:
|
Зона ответственности у документа одна — хранилище. Генератор описан отдельно:
|
||||||
его замысел — в [спеке генератора](../specs/2026-08-01-generator.md), формат
|
его замысел — в [спеке генератора](../specs/2026-08-01-generator.md), формат
|
||||||
@@ -157,6 +157,19 @@ Greenplum, чтобы словарь был общим у двух хранил
|
|||||||
назван: `_load_id` равен `run_id` Airflow и переносится из STG в годную строку
|
назван: `_load_id` равен `run_id` Airflow и переносится из STG в годную строку
|
||||||
ODS и в таблицу ошибок. Нужен ли он выше ODS, решается вместе с моделью DDS.
|
ODS и в таблицу ошибок. Нужен ли он выше ODS, решается вместе с моделью DDS.
|
||||||
|
|
||||||
|
## Актуальные события
|
||||||
|
|
||||||
|
`ods.event_rep` хранит физически принятые версии события в
|
||||||
|
`ReplacingMergeTree(_load_ts)`, а `ods.event_dist` служит для записи и
|
||||||
|
диагностики доставки. Обычное чтение этой пары может показать несколько строк
|
||||||
|
одного `WatchID`, пока фоновые слияния их не схлопнули.
|
||||||
|
|
||||||
|
`ods.event_v` — обычное представление над `ods.event_dist FINAL`. Оно сохраняет
|
||||||
|
имена, типы и зерно события, но всегда возвращает актуальную версию. Актуальное
|
||||||
|
состояние читают через него, а не повторяют `FINAL` в каждом запросе. Будущее
|
||||||
|
`dds.event_v` отвечает за другое: даёт событию язык бизнес-модели, в том числе
|
||||||
|
snake_case-имена и расшифровку кодов.
|
||||||
|
|
||||||
## Версии заказов
|
## Версии заказов
|
||||||
|
|
||||||
`ods.order_snapshot_rep` хранит принятые версии в
|
`ods.order_snapshot_rep` хранит принятые версии в
|
||||||
@@ -315,15 +328,15 @@ Airflow она стоит — там это обычный `SETTINGS` у зап
|
|||||||
получили N строк сырья и N в сумме событий и ошибок. Постоянной целью такой
|
получили N строк сырья и N в сумме событий и ошибок. Постоянной целью такой
|
||||||
опыт не становится, и почему — в [карте
|
опыт не становится, и почему — в [карте
|
||||||
проверок](testing.md), раздел «Интеграционная проверка постоянной целью не
|
проверок](testing.md), раздел «Интеграционная проверка постоянной целью не
|
||||||
становится». Счёт по ODS идёт через `FINAL`: голый `count()` по
|
становится». Счёт актуальных событий идёт через `ods.event_v`, которое скрывает
|
||||||
`ReplacingMergeTree` зависит от того, сколько мержей успело пройти, и спека это
|
`FINAL`: голый `count()` по `ReplacingMergeTree` зависит от того, сколько
|
||||||
прямо запрещает (раздел 6).
|
мержей успело пройти, и спека это прямо запрещает (раздел 6).
|
||||||
|
|
||||||
**Постоянная сверка одна, и она другого рода** — не слой против слоя, а ODS
|
**Постоянная сверка одна, и она другого рода** — не слой против слоя, а ODS
|
||||||
против [описи мира](../../data/world-inventory.json), лежащей в git
|
против [описи мира](../../data/world-inventory.json), лежащей в git
|
||||||
(`make check-clickhouse`, пришла с #42). Разъехаться сама она не может, и в
|
(`make check-clickhouse`, пришла с #42). Разъехаться сама она не может, и в
|
||||||
этом вся разница: опись описывает восемь дней стартового мира, счёт обрамлён
|
этом вся разница: опись описывает восемь дней стартового мира, счёт обрамлён
|
||||||
их датами, повтор заливки схлопывается под `FINAL`, а срок хранения ей не
|
их датами, повтор заливки схлопывается в `ods.event_v`, а срок хранения ей не
|
||||||
помеха — ODS хранит всё. Обе стороны равенства зафиксированы: одна кодом
|
помеха — ODS хранит всё. Обе стороны равенства зафиксированы: одна кодом
|
||||||
генератора, другая файлом в git. У межслойной сверки такой опоры нет ни с
|
генератора, другая файлом в git. У межслойной сверки такой опоры нет ни с
|
||||||
одной стороны.
|
одной стороны.
|
||||||
@@ -403,7 +416,7 @@ kafka_offset)`: смотрят такую таблицу от класса, а
|
|||||||
| `00-databases.sql` | базы слоёв |
|
| `00-databases.sql` | базы слоёв |
|
||||||
| `10-stg-tables.sql` | Kafka-таблица, локальная и распределённая таблицы сырья |
|
| `10-stg-tables.sql` | Kafka-таблица, локальная и распределённая таблицы сырья |
|
||||||
| `20-ods-tables.sql` | типизированное событие и таблица ошибок |
|
| `20-ods-tables.sql` | типизированное событие и таблица ошибок |
|
||||||
| `30-ods-views.sql` | матвью разбора: сырьё в событие и в ошибки |
|
| `30-ods-views.sql` | актуальные события и матвью разбора в ODS |
|
||||||
| `40-stg-views.sql` | матвью приёма: чтец в сырьё |
|
| `40-stg-views.sql` | матвью приёма: чтец в сырьё |
|
||||||
|
|
||||||
Порядок задают два правила. Первое: матвью принадлежит слою своей цели, а не
|
Порядок задают два правила. Первое: матвью принадлежит слою своей цели, а не
|
||||||
@@ -455,6 +468,7 @@ ODS. Второе: матвью приёма создаётся последне
|
|||||||
| STG | `stg.hits_raw_rep` / `_dist` | сырая строка сообщения плюс метаданные доставки |
|
| STG | `stg.hits_raw_rep` / `_dist` | сырая строка сообщения плюс метаданные доставки |
|
||||||
| STG | `stg.hits_raw_mv` | наполняет сырьё из чтеца |
|
| STG | `stg.hits_raw_mv` | наполняет сырьё из чтеца |
|
||||||
| ODS | `ods.event_rep` / `_dist` | типизированное широкое событие |
|
| ODS | `ods.event_rep` / `_dist` | типизированное широкое событие |
|
||||||
|
| ODS | `ods.event_v` | актуальная версия события с полями источника |
|
||||||
| ODS | `ods.event_errors_rep` / `_dist` | строки, не прошедшие строгий приём |
|
| ODS | `ods.event_errors_rep` / `_dist` | строки, не прошедшие строгий приём |
|
||||||
| ODS | `ods.event_mv`, `ods.event_errors_mv` | разбор сырья в событие и в ошибки |
|
| ODS | `ods.event_mv`, `ods.event_errors_mv` | разбор сырья в событие и в ошибки |
|
||||||
|
|
||||||
@@ -470,6 +484,13 @@ ODS. Второе: матвью приёма создаётся последне
|
|||||||
документацией — через MCP Context7, 5 августа 2026 года; то же разведение для
|
документацией — через MCP Context7, 5 августа 2026 года; то же разведение для
|
||||||
механики приёма — в [ADR 0005](../adr/0005-event-ingestion.md).
|
механики приёма — в [ADR 0005](../adr/0005-event-ingestion.md).
|
||||||
|
|
||||||
|
**Проверка `ods.event_v` 18 августа 2026 года.** Документация ClickHouse через
|
||||||
|
MCP Context7 подтвердила обычное представление с `FINAL` как способ скрыть
|
||||||
|
схлопывание от потребителя. На закреплённом ClickHouse 26.3 в ODS добавили
|
||||||
|
вторую физическую версию одного события: обычный счёт вырос на один, а
|
||||||
|
`ods.event_v` остался равен `ods.event_dist FINAL`. Имена и типы всех колонок
|
||||||
|
представления совпали с распределённой таблицей.
|
||||||
|
|
||||||
**Сверено с документацией.** Собственная колонка с именем виртуальной делает
|
**Сверено с документацией.** Собственная колонка с именем виртуальной делает
|
||||||
виртуальную недоступной. При вставке в `Distributed` шард выбирается по ключу
|
виртуальную недоступной. При вставке в `Distributed` шард выбирается по ключу
|
||||||
шардирования; фоновый режим — умолчание, а `distributed_foreground_insert = 1`
|
шардирования; фоновый режим — умолчание, а `distributed_foreground_insert = 1`
|
||||||
|
|||||||
@@ -407,7 +407,7 @@ README.
|
|||||||
| Kafka | `hits`, `orders` | два топика: `hits` — 2 партиции, `orders` — одна |
|
| Kafka | `hits`, `orders` | два топика: `hits` — 2 партиции, `orders` — одна |
|
||||||
| STG | `stg.hits_raw_kafka`, `stg.hits_raw` + MV | сырые строки событий, Kafka Engine на обеих нодах |
|
| STG | `stg.hits_raw_kafka`, `stg.hits_raw` + MV | сырые строки событий, Kafka Engine на обеих нодах |
|
||||||
| STG | `stg.orders_raw_kafka`, `stg.orders_raw`, без MV | сырые строки слепка; чтец на ноде 1, забирает пакетный шаг |
|
| STG | `stg.orders_raw_kafka`, `stg.orders_raw`, без MV | сырые строки слепка; чтец на ноде 1, забирает пакетный шаг |
|
||||||
| ODS | `ods.event` (+`_errors`) | типизированное широкое событие, ReplacingMergeTree |
|
| ODS | `ods.event` (+`_errors`), `ods.event_v` | типизированные версии событий, брак и точное актуальное состояние |
|
||||||
| ODS | `ods.order_snapshot` (+`_errors`), `ods.order_v` | типизированные версии заказов, брак и точное текущее состояние |
|
| ODS | `ods.order_snapshot` (+`_errors`), `ods.order_v` | типизированные версии заказов, брак и точное текущее состояние |
|
||||||
| DDS | `dds.session` | сборка сессий из событий (наследник `dds.click`) |
|
| DDS | `dds.session` | сборка сессий из событий (наследник `dds.click`) |
|
||||||
| DDS | `dds.event_v` | представление над `ods.event`: snake_case-имена, расшифровка кодов `DeviceCategory`; витрины DM читают его, а не ODS напрямую |
|
| DDS | `dds.event_v` | представление над `ods.event`: snake_case-имена, расшифровка кодов `DeviceCategory`; витрины DM читают его, а не ODS напрямую |
|
||||||
@@ -445,8 +445,11 @@ README.
|
|||||||
Kafka день переигрывается генератором заново: снимок — кэш чистой функции,
|
Kafka день переигрывается генератором заново: снимок — кэш чистой функции,
|
||||||
см. [спеку генератора](2026-08-01-generator.md).
|
см. [спеку генератора](2026-08-01-generator.md).
|
||||||
|
|
||||||
`dds.event_v` — первый на стенде пример правила «слой — это контракт, а не
|
`ods.event_v` сохраняет имена, типы и зерно ODS, но скрывает `FINAL` и возвращает
|
||||||
обязательно копия данных».
|
актуальную версию события. Будущее `dds.event_v` — другой контракт: оно
|
||||||
|
переводит событие на язык бизнес-модели, в том числе даёт snake_case-имена и
|
||||||
|
расшифровывает коды. Это первый на стенде пример правила «слой — это контракт,
|
||||||
|
а не обязательно копия данных».
|
||||||
|
|
||||||
Событие в DDS не дублируется: склейки четырёх источников больше нет, ODS уже
|
Событие в DDS не дублируется: склейки четырёх источников больше нет, ODS уже
|
||||||
широкий и типизированный; DDS хранит бизнес-сущности (сессия, заказ,
|
широкий и типизированный; DDS хранит бизнес-сущности (сессия, заказ,
|
||||||
|
|||||||
@@ -98,8 +98,9 @@ on_signal() {
|
|||||||
# этап 5 добавит следующий, — и голый счёт по таблице разойдётся с описью
|
# этап 5 добавит следующий, — и голый счёт по таблице разойдётся с описью
|
||||||
# законно, без всякой поломки.
|
# законно, без всякой поломки.
|
||||||
#
|
#
|
||||||
# **Счёт идёт через FINAL** — правило репозитория для ODS, довод в storage.md:
|
# **Счёт идёт через ods.event_v** — каноническую поверхность актуального
|
||||||
# голый `count()` по ReplacingMergeTree зависит от числа прошедших мержей.
|
# состояния ODS. Представление скрывает FINAL: потребителю не приходится
|
||||||
|
# помнить, что голый `count()` по ReplacingMergeTree зависит от числа мержей.
|
||||||
#
|
#
|
||||||
# Таблица брака здесь не второе утверждение, а объяснение первого. Утверждай мы
|
# Таблица брака здесь не второе утверждение, а объяснение первого. Утверждай мы
|
||||||
# «брака нет», проверка краснела бы навсегда после первого же урока, где менти
|
# «брака нет», проверка краснела бы навсегда после первого же урока, где менти
|
||||||
@@ -118,7 +119,7 @@ check_starting_world() {
|
|||||||
|
|
||||||
actual="$(query clickhouse-01 "
|
actual="$(query clickhouse-01 "
|
||||||
SELECT EventDate, count()
|
SELECT EventDate, count()
|
||||||
FROM ods.event_dist FINAL
|
FROM ods.event_v
|
||||||
WHERE EventDate BETWEEN '${first_date}' AND '${last_date}'
|
WHERE EventDate BETWEEN '${first_date}' AND '${last_date}'
|
||||||
GROUP BY EventDate
|
GROUP BY EventDate
|
||||||
ORDER BY EventDate
|
ORDER BY EventDate
|
||||||
|
|||||||
@@ -1,3 +1,13 @@
|
|||||||
|
-- ODS: каноническое чтение актуальных событий.
|
||||||
|
--
|
||||||
|
-- Представление сохраняет поля и зерно ods.event_dist, но скрывает FINAL от
|
||||||
|
-- потребителя. Поэтому физическая таблица остаётся местом диагностики
|
||||||
|
-- доставленных версий, а точное чтение получает одно имя.
|
||||||
|
CREATE VIEW IF NOT EXISTS ods.event_v ON CLUSTER clickstream_cluster
|
||||||
|
AS
|
||||||
|
SELECT *
|
||||||
|
FROM ods.event_dist FINAL;
|
||||||
|
|
||||||
-- ODS: разбор сырья в событие и в таблицу ошибок.
|
-- ODS: разбор сырья в событие и в таблицу ошибок.
|
||||||
--
|
--
|
||||||
-- Матвью две, и вместе они обязаны делить поток без зазора и без нахлёста:
|
-- Матвью две, и вместе они обязаны делить поток без зазора и без нахлёста:
|
||||||
|
|||||||
Reference in New Issue
Block a user