diff --git a/CONTEXT.md b/CONTEXT.md index d18bde1..3f726a4 100644 --- a/CONTEXT.md +++ b/CONTEXT.md @@ -147,3 +147,21 @@ _Избегать_: описание схемы, документация кон **Проигрыватель**: Компонент доставки готового потока дня в приёмник. Два режима: пакетный (пачкой, без темпа) и живой день. + +**Слепок**: +Полная выгрузка заказов окна изменяемости, снятая бэкендом на границе суток: +состояние заказов на этот момент, а не поток их изменений. Один заказ +приезжает в стольких слепках, сколько дней окна он прожил. + +**Окно изменяемости**: +Сколько модельных дней заказ ещё может измениться и потому продолжает ездить +в слепках. Константа мира: семь дней. За окном заказ замёрз, выручка дня +перестала «дышать». + +**Пакетный забор**: +Способ приёма топика, при котором данные тянет запрос по команде, а не +матвью непрерывно. Так стенд принимает заказы: у слепка есть начало и конец, +и забирать его уместно тогда, когда он собран целиком. Противоположность — +потоковый приём событий. +_Избегать_: путать с приёмником и пакетным режимом — те про сторону +генератора, про отправку; забор — про сторону хранилища, про чтение. diff --git a/docs/adr/0008-order-ingestion.md b/docs/adr/0008-order-ingestion.md new file mode 100644 index 0000000..1c93f50 --- /dev/null +++ b/docs/adr/0008-order-ingestion.md @@ -0,0 +1,152 @@ +# ADR 0008. Приём заказов: пакетный забор слепка, инициируемый Airflow + +Дата: 12 августа 2026 года. Статус: принято. Реализация — отдельным тикетом. + +## Решение + +Топик `orders` читает Kafka-таблица формата `RawBLOB` — та же форма чтеца, что +у событий ([ADR 0005](0005-event-ingestion.md)). Матвью к ней не привязана. +Сырьё забирает пакетный шаг, которым управляет Airflow: он вставляет +прочитанное в `stg.orders_raw_dist`, разбирает его в типизированный слепок и +заменяет партицию дня в `ods.order_snapshot`. Тот же даг проигрывает модельный +день генератором, поэтому переливается ровно то, что он положил в топик. + +Прямое чтение из Kafka-движка требует двух настроек, и вторая не очевидна: +`stream_like_engine_allow_direct_select = 1` разрешает читать чтеца запросом, +а `kafka_commit_on_select = 1` на самой таблице заставляет запрос коммитить +офсеты. По умолчанию прямое чтение их не коммитит — без этой настройки каждый +прогон забирает один и тот же слепок заново. + +Матвью к чтецу не привязана не по вкусу, а по устройству: при привязанной +матвью прямое чтение остаётся запрещённым независимо от первой настройки. Два +способа приёма взаимоисключающи по построению, и выбрать половину каждого +нельзя. + +У топика **одна партиция**, а чтец объявлен **только на `clickhouse-01` и без +`ON CLUSTER`**. Пишет пакетный шаг по-прежнему в распределённую таблицу, так +что раскладка по шардам не меняется: ключ у сырья прежний, `cityHash64` сырой +строки. + +Слои от этого получают разные роли, и каждая своя: + +- `stg.orders_raw` — байты как приехали, срок жизни и метаданные доставки как + у сырья событий; повторная заливка дня честно удваивает строки, как и там; +- `ods.order_snapshot` — типизированный слепок дня, идемпотентный **заменой + партиции** `snapshot_date`, ровно как обещает раздел 2 мастер-спеки; +- `dds.order` — дедуп до последней версии через `argMax`, без изменений. + +Сенсора дневного батча нет. Ждать нечего: производитель и потребитель слепка +живут в одном даге. + +## Почему + +**Слепок — не поток, и приём обязан это признать.** События приезжают +непрерывно, и матвью, тянущая их на лету, — честная форма для непрерывного. +Заказы приезжают раз в модельный день целой выгрузкой окна изменяемости; у +такой доставки есть начало и конец, и забирать её уместно по команде, а не +подписью на бесконечность. Стенд от этого получает не два оттенка одного +приёма, а два разных режима — push и pull, — и каждый стоит там, где ему место +по природе источника. Именно это сравнение раздел 12 мастер-спеки и заказывал +опорной точкой; прежняя его формулировка противопоставляла байтового чтеца +типизированному, то есть две формы одного и того же приёма, и переписана. + +**Почему не типизированный чтец прямо в `ods.order_snapshot`.** Он давал бы +строгий приём средствами движка и живое сравнение двух чтецов даром, но ODS +перестал бы быть надстройкой над STG и стал бы вторым входом с шины. Ровно эту +схему ADR 0005 отверг для событий, и повод здесь тот же: на стенде слои и есть +предмет изучения. + +**Почему не матвью из сырья в ODS, как у событий.** Тогда `ods.order_snapshot` +получил бы те же свойства, что `stg.orders_raw`: «как приехало», без дедупа, +дубли при повторной заливке законны. Два слоя подряд с одной ролью — один +лишний. Хуже того, обещание идемпотентности из раздела 2 повисло бы ни на чём: +матвью партиций не заменяет, а у ODS, в отличие от сырья, срока жизни нет, и +удвоение слепка жило бы вечно. При пакетном шаге каждый слой отрабатывает +своё, а замену партиции делает тот, кто данные и принёс. + +**Почему чтец остался байтовым.** Пулл снял возражение про второй вход с шины, +и типизированный чтец снова стал допустим — но тогда между двумя приёмами +стенда менялись бы сразу две переменные, и сравнение перестало бы читаться. +Меняется одна: push против pull. Вдобавок `RawBLOB` сохраняет за слепком то же, +что даёт событиям, — колонку, которую менти открывает в обычном клиенте и +читает глазами, и возможность переразобрать сырьё, не переигрывая день. +Свойства этого формата уже сняты живыми запросами при исполнении #37 и #43; +менять его здесь значило бы платить второй раз за уже купленное. + +**Почему одна партиция и один чтец.** Вторая партиция ничего не покупает: +слепок дня — порядка полутора тысяч строк, параллелизм не нужен. А урок «какая +нода читала топик» при пулле мёртв в любом случае — читает та нода, которую +спросили. То есть вторая партиция продаёт единственный настоящий риск схемы: +половина слепка застревает у ноды, к которой запрос не пришёл. Одна партиция +риск смягчает, но не снимает — брокер отдаст её любому из двух потребителей; +снимает его единственный чтец. У Airflow подготовлено одно подключение, к +`clickhouse-01`, — там чтецу и место. + +Объявление без `ON CLUSTER` — не оговорка к конвенции, а её первое осознанное +исключение: у пулла один тянущий по определению. Заодно это контрпример к +рефлексу «везде `ON CLUSTER`»: приставка не ритуал, а решение. + +**Почему не `Distributed` поверх Kafka.** Один запрос к распределённой таблице +над двумя чтецами осушил бы обе ноды разом и снял бы вопрос о партициях. +Механически это, скорее всего, работает — `Distributed` просто просит каждый +шард выполнить локальное чтение по имени таблицы, — но документация о такой +связке молчит: ни поддержки, ни запрета. Цена молчания высока. Офсеты +коммитятся на каждом шарде в момент чтения, а вставка идёт следом на +инициаторе: упала вставка — потеряны обе половины, а не одна. Недоступный шард +даёт худший из возможных исходов — тихо приехавшую половину слепка. Чинить +такое пришлось бы чтением исходников. И учебная цена своя: менти обязан +расшифровать конструкцию, которой нет ни в документации, ни в бою, и получает +за это трюк. При одной партиции и одном чтеце она не нужна вовсе. + +**Что отвергнуто ещё.** Сенсор дневного батча из раздела 9 мастер-спеки: у +топика нет сигнала «всё», и сенсор ловил бы момент, которого не существует, — +а раз генератор и переливка в одном даге, ждать нечего по построению. Лаба, +поднимающая типизированного чтеца во второй группе потребителей ради того же +сравнения: она понадобилась бы, останься сравнение невыполненным, но push +против pull даёт его живым и постоянным. + +## Следствия + +Урок про виртуальные колонки — «какая нода читала топик, меняется между +прогонами» — остаётся целиком за событиями. У заказов `consumer_host` всегда +один и тот же, и это честная разница двух режимов, а не потеря: при пулле +читает тот, кого спросили. + +Гарантий приёма у заказов не больше, чем у событий, но последствия мягче. +Офсеты коммитятся при чтении, вставка идёт следом — упавшая вставка теряет +пачку. У потока такая потеря невосстановима, у слепка её лечит следующий день: +окно изменяемости K = 7 привезёт те же заказы заново. + +Этап 3 забирает у этапа 5 первый настоящий даг. Раздел 9 мастер-спеки отдавал +даги этапу 5 целиком; приём заказов без дага не существует, поэтому порядок +меняется. `etl_pipeline` остаётся за этапом 5. + +Расхождения с мастер-спекой, внесённые тем же коммитом: раздел 6 (чтец заказов +живёт на одной ноде и без матвью), раздел 7 (развилка закрыта, у `orders` одна +партиция), раздел 9 (сенсор снят, даг переехал на этап 3), раздел 12 +(сравнение приёмов переформулировано). + +## Что проверено + +По документации ClickHouse через MCP Context7, 12 августа 2026 года. + +- Прямое чтение из движков-очередей (Kafka, RabbitMQ, FileLog) запрещено + начиная с версии 21.12 и открывается настройкой + `stream_like_engine_allow_direct_select`. +- При привязанной матвью прямое чтение остаётся запрещённым и с этой + настройкой. Отсюда вывод, что два способа приёма взаимоисключающи. +- Прямое чтение офсеты по умолчанию **не** коммитит; коммит включается + настройкой `kafka_commit_on_select` на самой таблице. Это тот подводный + камень, который молчит на первом прогоне и вылезает на втором. +- Сколько строк отдаёт одно чтение, задаёт `kafka_max_block_size`. При слепке + порядка полутора тысяч строк это один блок с запасом. +- Про `Distributed` поверх Kafka документация не говорит ничего — ни + поддержки, ни запрета. + +Осталось проверить при исполнении, и это работа тикета реализации: что второй +прогон подряд возвращает пусто, то есть офсеты действительно закоммичены; что +одного чтения хватает на весь слепок дня; что виртуальные колонки доставки +(`_topic`, `_partition`, `_offset`, `_timestamp_ms`) доступны при прямом чтении +— на них стоят служебные колонки сырья, см. [доку +хранилища](../architecture/storage.md); что повторная заливка дня даёт в +`ods.order_snapshot` тот же счёт, а не удвоенный. diff --git a/docs/specs/2026-07-30-stand-v2-realism.md b/docs/specs/2026-07-30-stand-v2-realism.md index 4159545..4546e80 100644 --- a/docs/specs/2026-07-30-stand-v2-realism.md +++ b/docs/specs/2026-07-30-stand-v2-realism.md @@ -307,12 +307,17 @@ CSV в репозитории (`data/catalog/products.csv`: `sku`, `name`, `cate - Роли нод: нода 1 — инициатор DDL и подключение Airflow; **Superset — на ноду 2**. Это осознанная ловушка правильных ошибок: забытый ON CLUSTER или VIEW поверх локальной таблицы проявляются в дашборде сами. -- **Приём Kafka**: Kafka-таблицы и MV — на обеих нодах, одна consumer group, - 2 партиции на топик; MV пишут в Distributed-цели. Раскладку решает ключ: +- **Приём Kafka**: у событий — Kafka-таблицы и MV на обеих нодах, одна + consumer group, 2 партиции на топик `hits`; MV пишут в Distributed-цели. + У заказов приём другого рода — пакетный забор без MV, чтец на одной ноде, + топик `orders` в одну партицию ([ADR 0008](../adr/0008-order-ingestion.md)). + Раскладку в обоих случаях решает ключ: события — по `cityHash64(ClientID)` (см. 1.3), заказы — `cityHash64(order_id)`, сырьё STG — `cityHash64(сырой строки)`; полный - список и доводы — в [доке хранилища](../architecture/storage.md). Урок: - «какая нода читала топик — меняется между прогонами, куда легли данные — нет». + список и доводы — в [доке хранилища](../architecture/storage.md). Урок + «какая нода читала топик — меняется между прогонами, куда легли данные — + нет» живёт на событиях: при пакетном заборе читает та нода, которую + спросили. - **Приём строгий**: пять опорных колонок — `WatchID`, `VisitID`, `ClientID`, `EventDate`, `UTCEventTime` — разбираются как `Nullable`, а набор ключей сообщения сверяется с контрактным; строка с NULL среди опорных колонок @@ -382,8 +387,9 @@ README. | Слой | Объект | Что это | |---|---|---| -| Kafka | `hits`, `orders` | два топика, по 2 партиции | -| STG | `stg.hits_raw_kafka`, `stg.hits_raw` + MV; для orders — развилка этапа 3, не решена (ниже) | сырые строки, Kafka Engine на обеих нодах | +| Kafka | `hits`, `orders` | два топика: `hits` — 2 партиции, `orders` — одна | +| STG | `stg.hits_raw_kafka`, `stg.hits_raw` + MV | сырые строки событий, Kafka Engine на обеих нодах | +| STG | `stg.orders_raw_kafka`, `stg.orders_raw`, без MV | сырые строки слепка; чтец на ноде 1, забирает пакетный шаг | | ODS | `ods.event` (+`_errors`) | типизированное широкое событие, ReplacingMergeTree | | ODS | `ods.order_snapshot` (+`_errors`) | слепки заказов как приехали, партиция по `snapshot_date`, без дедупа | | DDS | `dds.session` | сборка сессий из событий (наследник `dds.click`) | @@ -397,14 +403,16 @@ README. суффиксов; она же задаёт служебные колонки, нарезку и срок хранения сырья — см. [доку хранилища](../architecture/storage.md). -Как принимаются заказы — развилка этапа 3, и она не решена. Событиям выбран -приём сырья байтами с разбором функциями ([ADR 0005](../adr/0005-event-ingestion.md)); -заказам этот же способ идёт только вместе с ответом на вопрос, нужен ли им слой -сырья вообще — у них слепок, а не поток. Нужен — и типизированный чтец даст двух -чтецов на один топик, а такую схему ADR 0005 отверг; не нужен — и слои -перестают быть единообразными. Разбирать грилингом, когда дойдём до заказов; -как учебное сравнение двух способов приёма это записано и в опорных точках -раздела 12. +Заказы принимаются **пакетным забором** ([ADR +0008](../adr/0008-order-ingestion.md)): чтец топика байтовый, как у событий, но +матвью к нему не привязана, и сырьё забирает шаг, которым управляет Airflow — +он же вставляет прочитанное в `stg.orders_raw`, разбирает в типизированный +слепок и заменяет партицию дня в `ods.order_snapshot`. Тот же даг проигрывает +модельный день генератором, поэтому переливается ровно то, что он положил в +топик. Слой сырья у заказов остаётся: без него `ods.order_snapshot` повторил бы +его роль, а обещание идемпотентности повисло бы — матвью партиций не заменяет. +Стенд получает от этого два режима приёма рядом, поток и слепок, и сравнение +из опорных точек раздела 12 переформулировано под них. Состав служебных колонок задаёт дока хранилища. Спеке важны два следствия: `ods.event` и `ods.order_snapshot` получают метку загрузки `_load_ts`, и у @@ -495,7 +503,7 @@ v2 стартует пустым, поэтому объём ниже — это | Генератор | с нуля: модель v1 не переносится (другая модель данных, плюс известные проблемы производительности v1); широкое событие, таксономия, анонимность, N:1, заказы слепками, расхождения A–D, каталог; масштаб — ~4–5 тыс. строк с тестами | L | | Инфраструктура | compose: 2 ноды CH + keeper + остальной стенд; конфиги кластера, макросы; make/скрипты | M — ~10–12 файлов | | SQL | DDL по слоям и ролям (ON CLUSTER, Replicated*, Distributed; раскладка файлов — в доке хранилища) + трансформации событий, заказов, identity, сверки + словарь | L — ~12–15 файлов, главная сложность | -| Airflow | DAG'и по образцу v1: etl_pipeline (партиционная переобработка, ожидание дневного батча заказов — сенсор/Datasets), world_init/next_day, helpers | M — ~5–6 файлов | +| Airflow | DAG'и по образцу v1: etl_pipeline (партиционная переобработка), world_init/next_day, helpers; первый даг — переливка слепка — приходит раньше, этапом 3 | M — ~5–6 файлов | | Superset | датасеты + дашборд с тремя новыми сюжетами | M — 2 файла | | Эталонный мир | пересборка снимка на месте, счётчики описи, чек-скрипты | M–L | | Мониторинг | дашборды Grafana «данные», «кластер», «запросы»; ClickHouse источником данных, панели на SQL; Prometheus тонким полом (ADR 0002) | M — конфиги и дашборды | @@ -537,10 +545,11 @@ v2 стартует пустым, поэтому объём ниже — это целиком). В конце этапа фиксируется маленький стартовый мир для стабильных приёмок следующих этапов (полная пересборка эталонного мира — отдельный этап 7). -3. Заказы и каталог: генератор слепков, STG/ODS/DDS заказа, словарь. +3. Заказы и каталог: генератор слепков, STG/ODS/DDS заказа, словарь и первый + настоящий даг — проигрыш модельного дня плюс переливка слепка ([ADR + 0008](../adr/0008-order-ingestion.md)). 4. Трансформации и витрины: сессии, identity_map, выручка, сверка A+C. -5. Airflow: `etl_pipeline` (партиционная переобработка, ожидание дневного - батча заказов — сенсор/Datasets). +5. Airflow: `etl_pipeline` (партиционная переобработка). 6. Расхождения B+D и опоздания; счётчики описи. 7. Эталонный мир: опись и пересборка снимка, чек-скрипты; CI-генерация на amd64 и arm64. @@ -603,9 +612,13 @@ Kafka Engine на двух нодах снят с этого списка при - Размер артефакта эталонного мира после пересборки. - Спорные API (Airflow Datasets/сенсоры, ClickHouse DDL) — перед кодом сверять через MCP Context7 (правило AGENTS.md). -- Airflow 3.x: версия фиксируется на этапе 1 (каркас); DAG'и этапа 5 пишутся - под API третьей версии (Datasets → Assets) — актуальные операторы и сенсоры - проверить через Context7. +- Airflow 3.x: версия фиксируется на этапе 1 (каркас); DAG'и пишутся под API + третьей версии (Datasets → Assets) — актуальные операторы проверить через + Context7. Первый даг приходит этапом 3, а не 5 ([ADR + 0008](../adr/0008-order-ingestion.md)). +- Пакетный забор слепка из Kafka: коммит офсетов прямым чтением, хватает ли + одного чтения на слепок дня, доступны ли при нём виртуальные колонки + доставки — список и ответы в [ADR 0008](../adr/0008-order-ingestion.md). - Генератор: рабочее решение — Python с производительной архитектурой (батчевая генерация вместо посточной, быстрая JSON-сериализация, распараллеливание по модельным дням). Читаемость генератора для менти — @@ -653,14 +666,14 @@ v2, этап 0). синтетическая постановка — осознанный приём; - лекция «`Sign` и CollapsingMergeTree»: почему на стенде `sum(Sign)` = `count()`, а в бою — нет; частый вопрос на собеседованиях; -- два способа принять топик, рядом на одном стенде: сырьё байтами с разбором - функциями (`hits`, [ADR 0005](../adr/0005-event-ingestion.md)) против - типизированного чтеца с `kafka_handle_error_mode` — сравнение цены и - наблюдаемости как задание. **Развилка этапа 3, не решена**: типизированный - чтец идёт заказам только вместе с ответом на вопрос, нужен ли им слой сырья. - Нужен — и чтецов на один топик станет два, а эту схему ADR 0005 отверг; не - нужен — и слои перестают быть единообразными. Разбирать грилингом, когда - дойдём до заказов; +- два способа принять топик, рядом на одном стенде: **поток против слепка, + push против pull**. События тянет матвью на лету ([ADR + 0005](../adr/0005-event-ingestion.md)), слепок заказов забирает пакетный шаг + по команде дага ([ADR 0008](../adr/0008-order-ingestion.md)); чтец в обоих + случаях байтовый, так что различает их ровно режим. Сравнение цены, + наблюдаемости и того, чем платит каждый режим, — как задание. Третий способ, + типизированный чтец с `kafka_handle_error_mode`, на стенде не живёт: он + отвергнут обоими ADR, и остаётся материалом для рассказа; - матвью как рабочий механизм, а не диковина: их видно на приёме и на сборке ODS, а пакетная работа начинается выше. Отдельным заданием — как читать из ODS последние версии, через `FINAL` или оконной функцией: что нагляднее, решаем на