From 6bd91901423f1614ff838e99ea8ebeb2c60282e2 Mon Sep 17 00:00:00 2001 From: Dmitry Dementiev Date: Tue, 18 Aug 2026 21:33:13 +0300 Subject: [PATCH] =?UTF-8?q?feat(orders):=20=D0=B4=D0=BE=D0=B1=D0=B0=D0=B2?= =?UTF-8?q?=D0=B8=D1=82=D1=8C=20=D0=BF=D1=80=D0=B8=D1=91=D0=BC=20=D1=81?= =?UTF-8?q?=D0=BB=D0=B5=D0=BF=D0=BA=D0=BE=D0=B2=20=D0=B2=20STG?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Зачем: - связать проигрывание модельных дней с пакетным приёмом заказов - показать на одном стенде различие потокового push и пакетного pull Что: - добавлен топик, Kafka-чтец и реплицированное сырьё заказов - добавлен даг orders_ingest с одним прямым чтением и идентификатором загрузки - работники мира отправляют слепок, ждут приём и только затем двигают позицию - решения, границы отказа и проверки отражены в ADR и архитектурных документах Проверка: - make lint - make config-test - make smoke - make check-clickhouse - make check-services --- README.md | 6 ++ compose.yaml | 11 +++ dags/orders_ingest.py | 100 ++++++++++++++++++++++++++ dags/world_control.py | 66 ++++++++++++++--- docs/adr/0008-order-ingestion.md | 82 ++++++++++++++------- docs/adr/0009-world-control.md | 10 ++- docs/architecture/orders/README.md | 5 +- docs/architecture/orders/ingestion.md | 22 ++++-- docs/architecture/storage.md | 27 ++++--- sql/ddl/10-stg-tables.sql | 55 +++++++++++++- 10 files changed, 330 insertions(+), 54 deletions(-) create mode 100644 dags/orders_ingest.py diff --git a/README.md b/README.md index ad7ce26..d0368ee 100644 --- a/README.md +++ b/README.md @@ -233,6 +233,12 @@ uv run --project generator python -m clickstream_generator batch \ | `world_live_day` | играет следующий день в темпе модельного времени, около двадцати четырёх минут | | `world_live` | выключатель: пока снят с паузы, дёргает `world_live_day` день за днём | +Работник не только играет день: следом он отправляет слепок заказов за +сыгранное и дёргает отдельный даг, `orders_ingest`, — тот забирает приехавший +слепок из топика `orders` в сырьё, — и дожидается его конца. Так две половины +стенда идут в ногу: упал приём — краснеет и работник, а позиция остаётся на +месте ([ADR 0008](docs/adr/0008-order-ingestion.md)). + Позицию хранит переменная Airflow `world_position` — номер первого несыгранного дня. Ставит её работник, сыгравший день, и только по успеху: оборванный прогон позицию не двигает, и следующий запуск играет тот же день заново. Нет diff --git a/compose.yaml b/compose.yaml index 916aee4..0e0687d 100644 --- a/compose.yaml +++ b/compose.yaml @@ -68,6 +68,7 @@ x-airflow-common: &airflow-common GENERATOR_IMAGE: *generator-image STAND_NETWORK: ${COMPOSE_PROJECT_NAME}_default WORLD_STARTING_DAYS: *world-starting-days + KAFKA_ORDERS_TOPIC: orders volumes: - ./dags:/opt/airflow/dags:ro - ./infra/airflow/init.sh:/opt/airflow/init.sh:ro @@ -238,6 +239,16 @@ services: exit 1 fi + # Топик заказов: одна партиция, читает его один чтец на ноде 1 — + # довод в ADR 0008. + /opt/kafka/bin/kafka-topics.sh \ + --bootstrap-server kafka:9092 \ + --create \ + --if-not-exists \ + --topic orders \ + --partitions 1 \ + --replication-factor 1 + # DDL применяется с ноды 1 по порядку имён файлов после готовности всего кластера. clickhouse-init: image: clickhouse/clickhouse-server:26.3.17.56 diff --git a/dags/orders_ingest.py b/dags/orders_ingest.py new file mode 100644 index 0000000..49da20b --- /dev/null +++ b/dags/orders_ingest.py @@ -0,0 +1,100 @@ +"""Приём заказов: пакетный забор слепка из Kafka в хранилище. + +Путь Kafka → STG → ODS у заказов один и живёт одним дагом; здесь его первый +шаг — забор. Расписания у дага нет: его дёргает тот, кто положил слепок в +топик, — работник пульта мира после проигрыша дня, — и ждёт конца. + +Контраст с приёмом событий и есть урок. События тянет матвью, навсегда +подписанная на чтеца: приём идёт сам, пока идёт поток. Слепок заказов +приезжает раз в модельный день целой выгрузкой, у которой есть начало и конец, +и забирает её запрос по команде. Push против pull — два режима на одном стенде, +каждый там, где ему место по природе источника. + +Решения и доводы целиком — ADR 0008; форма перехода STG → ODS, который +прирастёт сюда вторым шагом, — docs/architecture/orders/ingestion.md. +""" + +from __future__ import annotations + +import datetime +import logging + +from airflow.sdk import Connection, dag, get_current_context, task + +START_DATE = datetime.datetime(2026, 1, 1, tzinfo=datetime.UTC) + +# Забор — один прямой SELECT, без цикла до пустоты: одна порция ClickHouse +# берёт десятки тысяч сообщений, а слепок дня — порядка полутора тысяч строк. +# Короткая порция оставит хвост до следующего прогона, а отказ после чтения +# унесёт прочитанное с собой: офсеты коммитятся в момент чтения. Граница +# целиком — ADR 0008, «Следствия». +# +# stream_like_engine_allow_direct_select разрешает читать чтеца запросом; вторая +# половина пары объявлена на самой таблице (sql/ddl/10-stg-tables.sql). +# distributed_foreground_insert = 1 — конвенция ETL-вставок стенда: задача не +# должна зеленеть раньше, чем строки легли на шарды. +TAKE_ONE_BATCH = """ +INSERT INTO stg.orders_raw_dist +SELECT + raw, + _topic AS kafka_topic, + _partition AS kafka_partition, + _offset AS kafka_offset, + _timestamp_ms AS kafka_timestamp, + hostName() AS consumer_host, + {load_id:String} AS _load_id, + now64(3) AS _load_ts +FROM stg.orders_raw_kafka +SETTINGS + stream_like_engine_allow_direct_select = 1, + distributed_foreground_insert = 1 +""" + + +@dag( + dag_id="orders_ingest", + schedule=None, + start_date=START_DATE, + is_paused_upon_creation=False, + # Чтец у топика один, и группа потребителей у него одна. Два прогона разом + # дрались бы за неё, а слепок разъехался бы по двум _load_id; второй + # прогон подождёт своей очереди. + max_active_runs=1, + tags=["заказы"], +) +def orders_ingest(): + """Забрать приехавший слепок заказов из топика в сырьё.""" + + @task + def pull_batch() -> None: + # Импорт внутри задачи: обработчик DAG разбирает этот файл снова и + # снова, и импорт наверху оплачивался бы каждым разбором. + import clickhouse_connect + + load_id = get_current_context()["run_id"] + connection = Connection.get("clickhouse_default") + client = clickhouse_connect.get_client( + host=connection.host, + port=connection.port, + username=connection.login, + password=connection.password, + database=connection.schema or "default", + connect_timeout=5, + send_receive_timeout=30, + ) + try: + summary = client.command(TAKE_ONE_BATCH, parameters={"load_id": load_id}) + finally: + client.close() + # Размер порции — read_rows: written_rows у вставки в Distributed + # считает не приехавшее. + logging.info( + "порция принята: строк %s, _load_id %s", + summary.summary["read_rows"], + load_id, + ) + + pull_batch() + + +orders_ingest() diff --git a/dags/world_control.py b/dags/world_control.py index fcad3a4..0ecad54 100644 --- a/dags/world_control.py +++ b/dags/world_control.py @@ -7,6 +7,11 @@ конца. Снят с паузы — мир едет день за днём; поставлен на паузу — встал на границе модельных суток. +Сыграть день — половина работы. Вторая половина: отправить слепок заказов и +дождаться, пока хранилище его заберёт. Поэтому у обоих работников за проигрышем +идут отправка слепка тем же контейнером генератора и ждущий триггер дага приёма +`orders_ingest`. + Разделение не косметическое. Расписание на самом работнике заставило бы кнопку паузы значить две вещи разом — «мир не едет сам» и «даг выключен», — а работник при этом выглядел бы в списке выключенным, хотя нажимают его каждый день. @@ -35,6 +40,9 @@ GENERATOR_ENVIRONMENT = { "KAFKA_BOOTSTRAP_SERVERS": os.environ["KAFKA_BOOTSTRAP_SERVERS"], "KAFKA_TOPIC": os.environ["KAFKA_TOPIC"], } +# Топик слепка называется аргументом, а не окружением: KAFKA_TOPIC выше — топик +# событий, и промолчи мы, заказы уехали бы к ним. +ORDERS_TOPIC = os.environ["KAFKA_ORDERS_TOPIC"] STARTING_DAYS = int(os.environ["WORLD_STARTING_DAYS"]) # Позиция на оси: номер первого несыгранного дня. Переменной нет — мир в @@ -50,6 +58,10 @@ WORLD_POSITION = "world_position" # двадцати четырёх минут: тик только спрашивает «не пора ли снова». LIVE_TICK = datetime.timedelta(minutes=25) +# Какой день играть, работник узнаёт у первой задачи. Спрашивают её все +# генераторные шаги: слепок снимается с той же позиции, что и проигрыш. +PLAYED_DAY = "{{ ti.xcom_pull(task_ids='first_unplayed_day') }}" + START_DATE = datetime.datetime(2026, 1, 1, tzinfo=datetime.UTC) TAGS = ["пульт мира"] @@ -60,8 +72,8 @@ def first_unplayed_day() -> int: return int(Variable.get(WORLD_POSITION, default=STARTING_DAYS)) -def _play(task_id: str, command: list[str]) -> DockerOperator: - """Задача, играющая дни в каноническом контейнере генератора. +def _generator(task_id: str, command: list[str]) -> DockerOperator: + """Задача, зовущая генератор в его каноническом контейнере. Внутрь образа Airflow генератор не поставить: он требует Python 3.14, а образ несёт 3.13. Да и обещание побайтовой воспроизводимости дано для @@ -83,6 +95,22 @@ def _play(task_id: str, command: list[str]) -> DockerOperator: ) +def _ingest_orders() -> TriggerDagRunOperator: + """Задача забора: дёрнуть даг приёма и дождаться, чем он кончился. + + Ожидание здесь несущее. Без него работник позеленел бы, не узнав, доехал + ли слепок, и позиция мира ушла бы вперёд хранилища — а зелёный конец графа + не должен переживать отказ выше (ADR 0003). + """ + return TriggerDagRunOperator( + task_id="trigger_orders_ingest", + trigger_dag_id="orders_ingest", + wait_for_completion=True, + # Умолчание — минута опроса, а весь прогон работника пачкой короче. + poke_interval=10, + ) + + @dag( dag_id="world_next_day", schedule=None, @@ -90,7 +118,11 @@ def _play(task_id: str, command: list[str]) -> DockerOperator: is_paused_upon_creation=False, max_active_runs=1, tags=TAGS, - params={"days": Param(1, type="integer", minimum=1, title="Сколько дней прожить")}, + params={ + "days": Param( + 1, type="integer", minimum=1, maximum=7, title="Сколько дней прожить" + ) + }, ) def world_next_day(): """Прожить следующие дни пачкой, без пауз. @@ -106,18 +138,27 @@ def world_next_day(): Variable.set(WORLD_POSITION, str(first_day + days)) first_day = first_unplayed_day() - played = _play( + played = _generator( "play_days", + ["batch", "--day", PLAYED_DAY, "--days", "{{ params.days }}"], + ) + # Разгон на N дней отправляет N слепков — по одному за сыгранный день, и + # каждый со своим сдвигом: прогон дня D везёт слепок дня D−1. Забор при + # этом остаётся один. + sent = _generator( + "send_snapshots", [ - "batch", + "snapshot", "--day", - "{{ ti.xcom_pull(task_ids='first_unplayed_day') }}", + PLAYED_DAY, "--days", "{{ params.days }}", + "--topic", + ORDERS_TOPIC, ], ) - first_day >> played >> remember_played(first_day) + first_day >> played >> sent >> _ingest_orders() >> remember_played(first_day) @dag( @@ -142,12 +183,15 @@ def world_live_day(): Variable.set(WORLD_POSITION, str(first_day + 1)) first_day = first_unplayed_day() - played = _play( - "play_day", - ["live", "--day", "{{ ti.xcom_pull(task_ids='first_unplayed_day') }}"], + played = _generator("play_day", ["live", "--day", PLAYED_DAY]) + # Слепок уезжает пачкой и после дня, а не в темпе: у выгрузки бэкенда темпа + # нет вовсе — она снимается на границе суток целиком. + sent = _generator( + "send_snapshot", + ["snapshot", "--day", PLAYED_DAY, "--topic", ORDERS_TOPIC], ) - first_day >> played >> remember_played(first_day) + first_day >> played >> sent >> _ingest_orders() >> remember_played(first_day) @dag( diff --git a/docs/adr/0008-order-ingestion.md b/docs/adr/0008-order-ingestion.md index 2f960bc..0f2a936 100644 --- a/docs/adr/0008-order-ingestion.md +++ b/docs/adr/0008-order-ingestion.md @@ -1,13 +1,15 @@ # ADR 0008. Приём заказов: пакетный забор слепка, инициируемый Airflow Дата: 12 августа 2026 года. Статус: частично заменено -[ADR 0010](0010-order-versions-in-ods.md). Реализация — отдельным тикетом. +[ADR 0010](0010-order-versions-in-ods.md); забор построен тикетом #93. Сохраняются пакетный забор, `RawBLOB`, один чтец на `clickhouse-01`, одна партиция топика и отсутствие матвью. Отменены граница «одно чтение — полный слепок», замена партиции `snapshot_date`, ODS без дедупликации и готовая схема -`dds.order` с `argMax`. Действующая форма STG → ODS описана в -[спецификации приёма заказов](../architecture/orders/ingestion.md). +`dds.order` с `argMax` — это ADR 0010. Отдельно от него, решением #93 от +17 августа 2026 года, производитель и потребитель слепка разъехались по разным +дагам; раздел «Решение» описывает построенное. Действующая форма STG → ODS +описана в [спецификации приёма заказов](../architecture/orders/ingestion.md). ## Решение @@ -15,13 +17,18 @@ у событий ([ADR 0005](0005-event-ingestion.md)). Матвью к ней не привязана. Сырьё забирает пакетный шаг, которым управляет Airflow: он вставляет прочитанное в `stg.orders_raw_dist`, разбирает его в типизированный слепок и -заменяет партицию дня в `ods.order_snapshot`. Тот же даг проигрывает модельный -день генератором, поэтому переливается ровно то, что он положил в топик. -Уточнение со сдвигом отправки, решённым позже (#71, [слепок и его -доставка](../architecture/orders/snapshot.md)): даг, играющий день D, в -штатном прогоне кладёт и забирает слепок дня D−1 — слепок предыдущего дня, а -не сыгранного. После падения между шагами в топике может ждать и хвост -прежних слепков; забор принимает всё приехавшее. +заменяет партицию дня в `ods.order_snapshot`. + +Производитель и потребитель слепка — разные даги. Работник пульта мира играет +день и отправляет слепок, а забирает его отдельный даг `orders_ingest`; зовёт +забор сам работник ждущим `TriggerDagRunOperator` и зеленеет только вслед за +ним. Связывает две половины не общий даг, а этот дождавшийся успех: зелёный +работник значит, что забор отработал. Уточнение со сдвигом +отправки, решённым позже (#71, [слепок и его +доставка](../architecture/orders/snapshot.md)): работник, играющий день D, в +штатном прогоне кладёт слепок дня D−1 — слепок предыдущего дня, а не +сыгранного. После падения между шагами в топике может ждать и хвост прежних +слепков; забор принимает всё приехавшее. Прямое чтение из Kafka-движка требует двух настроек, и вторая не очевидна: `stream_like_engine_allow_direct_select = 1` разрешает читать чтеца запросом, @@ -47,8 +54,8 @@ партиции** `snapshot_date`, ровно как обещает раздел 2 мастер-спеки; - `dds.order` — дедуп до последней версии через `argMax`, без изменений. -Сенсора дневного батча нет. Ждать нечего: производитель и потребитель слепка -живут в одном даге. +Сенсора дневного батча нет. Ждать нечего: забор дёргает сам отправитель +слепка, когда отправил. ## Почему @@ -112,7 +119,7 @@ **Что отвергнуто ещё.** Сенсор дневного батча из раздела 9 мастер-спеки: у топика нет сигнала «всё», и сенсор ловил бы момент, которого не существует, — -а раз генератор и переливка в одном даге, ждать нечего по построению. Лаба, +а раз забор зовёт сам отправитель слепка, ждать нечего по построению. Лаба, поднимающая типизированного чтеца во второй группе потребителей ради того же сравнения: она понадобилась бы, останься сравнение невыполненным, но push против pull даёт его живым и постоянным. @@ -124,10 +131,13 @@ один и тот же, и это честная разница двух режимов, а не потеря: при пулле читает тот, кого спросили. -Гарантий приёма у заказов не больше, чем у событий, но последствия мягче. -Офсеты коммитятся при чтении, вставка идёт следом — упавшая вставка теряет -пачку. У потока такая потеря невосстановима, у слепка её лечит следующий день: -окно изменяемости K = 7 привезёт те же заказы заново. +Гарантий приёма у заказов не больше, чем у событий, и граница проходит по +живому. Успешное чтение короткой порции оставляет непрочитанный хвост в +топике — он дождётся следующего забора. А отказ после чтения теряет саму +порцию: офсеты закоммичены в момент чтения, часть строк могла лечь на шарды, +и повторное чтение вернёт ноль. Позицию мира это не двигает, поэтому день +сыграется и уедет заново; в сырье он тогда окажется частичным дублем, а +старый хвост, ушедший той же порцией, не вернётся. Этап 3 забирает у этапа 5 первый настоящий даг. Раздел 9 мастер-спеки отдавал даги этапу 5 целиком; приём заказов без дага не существует, поэтому порядок @@ -158,10 +168,34 @@ ClickHouse `26.3.17.56`. - Про `Distributed` поверх Kafka документация не говорит ничего — ни поддержки, ни запрета. -Осталось проверить при исполнении, и это работа тикета реализации: что второй -прогон подряд возвращает пусто, то есть офсеты действительно закоммичены; что -одного чтения хватает на весь слепок дня; что виртуальные колонки доставки -(`_topic`, `_partition`, `_offset`, `_timestamp_ms`) доступны при прямом чтении -— на них стоят служебные колонки сырья, см. [доку -хранилища](../architecture/storage.md); что повторная заливка дня даёт в -`ods.order_snapshot` тот же счёт, а не удвоенный. +Забор проверен на живом стенде 18 августа 2026 года при исполнении #93 — +ClickHouse `26.3.17.56`, Airflow 3.3.0. Три вопроса, оставленные этим ADR +реализации, закрыты; заодно снят отказной путь. Обе настройки прямого чтения +доезжают до запроса, объявленные в конце `INSERT ... SELECT`: `system.query_log` +показывает у каждой вставки единицу. + +- **Виртуальные колонки доставки при прямом чтении доступны все четыре.** + `_topic`, `_partition`, `_offset` и `_timestamp_ms` читаются тем же + выражением, что в матвью приёма событий, и метаданные в `stg.orders_raw` + заполнены: 1694 строки слепка дня 7 приехали с `orders`, партицией 0, + сплошными офсетами 0…1693 и непустой меткой брокера. Отдельного механизма + пуллу не понадобилось. +- **Одного чтения хватает и на слепок, и на разгон.** Пять заборов подряд + взяли 1694, 1730, 3446, 1694 и 5097 строк — каждый раз ровно столько, + сколько напечатал генератор. Числа сверх слепка объясняются сами: 3446 — это + свежий слепок плюс хвост, оставшийся от упавшего прогона, а 5097 — три + слепка разгона `days = 3`, уехавшие одной порцией. +- **Офсеты действительно коммитятся.** Второй забор подряд по пустому топику + вернул ноль строк, а офсеты следующего продолжились с 1694 — то есть + `kafka_commit_on_select` работает, и слепок не забирается заново. +- **Упавший забор не двигает мир.** Чтец снесли, и прогон работника покраснел + на триггере забора: `remember_played` ушла в `upstream_failed`, позиция + осталась прежней, а после починки тот же день сыгран заново, и ждавший хвост + уехал одной порцией вместе со свежим слепком. Чтения в этом опыте не было + вовсе — падать было нечему, — поэтому он показывает только несдвиг позиции. + Судьба уже прочитанной порции другая, и она описана в «Следствиях». + +Осталось проверить при исполнении #94, на стороне ODS: что повторный разбор +того же `_load_id` не двоит версии заказа и не меняет `_load_ts` — критерии +целиком в [спецификации приёма заказов](../architecture/orders/ingestion.md), +раздел «Риски и проверка». diff --git a/docs/adr/0009-world-control.md b/docs/adr/0009-world-control.md index b3092a5..2cdbbc7 100644 --- a/docs/adr/0009-world-control.md +++ b/docs/adr/0009-world-control.md @@ -8,8 +8,10 @@ **Работники** — без расписания и без паузы, запускаются руками: -- `world_next_day` — день пачкой; параметр «сколько дней» (умолчание 1) даёт - разгон вперёд одним нажимом; +- `world_next_day` — день пачкой; параметр «сколько дней» (умолчание 1, потолок + 7) даёт разгон вперёд одним нажимом. Потолок — это и есть названное желание + «уедь на неделю сейчас», и ровно тот разгон, который забор заказов забирает + одним чтением ([ADR 0008](0008-order-ingestion.md)); - `world_live_day` — один день в темпе живого режима. **Выключатель** — без своей работы: `world_live` несёт расписание около двадцати @@ -99,6 +101,10 @@ которого стенд не поднимает; заводить службу ради одного ожидания дороже занятого слота. +Приросший к работникам забор заказов (#93, [ADR 0008](0008-order-ingestion.md)) +делает из двух занятых слотов три: пока идёт забор, ждут выключатель и сам +работник, а третий слот занимает задача `orders_ingest`. + **Пакетная автоматика отложена, потому что своего желания у неё пока одно.** «Шагни и стой» закрывает кнопка, «уедь на неделю сейчас» — параметр «сколько дней», «живи, пока я смотрю» — выключатель. Ей остаётся «едь сам быстрее, чем diff --git a/docs/architecture/orders/README.md b/docs/architecture/orders/README.md index a35849d..46bf515 100644 --- a/docs/architecture/orders/README.md +++ b/docs/architecture/orders/README.md @@ -57,8 +57,9 @@ доли классов, стоимость доставки — калибровка при реализации; финальная фиксация чисел — пересборка эталонного мира, этап 7. При пересборке правки потребуют только числа, не устройство. -- **Проверки приёма** — разовая приёмка допущения «один запуск — одно чтение» - и опыты из [«Рисков и проверки»](ingestion.md) — тикет реализации приёма. +- **Проверки приёма** — опыты из [«Рисков и проверки»](ingestion.md) про брак + и версии в ODS — тикет перехода STG → ODS (#94). Допущение «один запуск — + одно чтение» принято живым прогоном при исполнении #93. - **`_load_id` выше ODS** — вместе с устройством `dds.order` (#85). - **Контур проверок качества для расхождений** (даг DQ, `dm.dq_summary`) — остаётся в тумане карты #69; естественное место разговора — этап 4. diff --git a/docs/architecture/orders/ingestion.md b/docs/architecture/orders/ingestion.md index 72795c2..787abcb 100644 --- a/docs/architecture/orders/ingestion.md +++ b/docs/architecture/orders/ingestion.md @@ -34,9 +34,10 @@ ## Поток данных После завершения генератора Airflow один раз читает байтовый Kafka-чтец и -записывает полученную порцию в `stg.orders_raw`. Все строки получают `_load_id`, -равный `run_id` Airflow. `_load_ts` вычисляется при этой записи и дальше -переносится без пересчёта. +записывает полученную порцию в `stg.orders_raw`. Даг зовётся `orders_ingest`, а +дёргает его тот, кто положил слепок в топик, — работник пульта мира, — и ждёт +конца прогона. Все строки получают `_load_id`, равный `run_id` Airflow. +`_load_ts` вычисляется при этой записи и дальше переносится без пересчёта. Один следующий `task_id` отвечает за весь переход STG → ODS. Внутри него два последовательных `INSERT SELECT` читают неизменный срез по `_load_id`: первый @@ -119,6 +120,13 @@ JSON-объектом с точным набором ключей: `order_id`, ` из следующих запусков. После отказа от замены партиции это задержка, а не потеря или публикация неполного дня. +Отказ — случай другой, и «хвост дождётся» на него не распространяется. Офсеты +порции коммитятся в момент чтения, поэтому упавшая вставка уносит прочитанное с +собой: повторное чтение вернёт ноль, а часть строк может уже лежать на шарде. +Позиция мира при этом не двигается, и тот же день уедет заново — в сырье он +окажется частичным дублем, а старый хвост, ушедший той же порцией, не вернётся +([ADR 0008](../../adr/0008-order-ingestion.md), «Следствия»). + ## Отклонённые варианты - Партиционная идемпотентность — `REPLACE PARTITION snapshot_date` или @@ -137,9 +145,6 @@ JSON-объектом с точным набором ключей: `order_id`, ` ## Риски и проверка -- На стандартном мире сверить число отправленных заказов с числом строк, - принятых одним прямым чтением. Это разовая приёмка допущения, не постоянный - сторож. - На малой управляемой порции дать по одной строке каждого класса брака и две годные версии одного `order_id`. Две цели должны сохранить все непустые сообщения, а `ods.order_v` — вернуть новую версию независимо от фонового @@ -151,6 +156,11 @@ JSON-объектом с точным набором ключей: `order_id`, ` ## Что проверено +Забор из Kafka в STG снят на живом стенде 18 августа 2026 года при исполнении +#93: одно прямое чтение приносит весь слепок дня, метаданные доставки доступны, +офсеты коммитятся, а сбой забора не двигает позицию мира. Числа — [ADR +0008](../../adr/0008-order-ingestion.md), раздел «Что проверено». + MCP Context7 в сессии проектирования был недоступен. На локальном ClickHouse `26.3.17.56` проверено, что прямой `SELECT` Kafka Engine завершается после одной порции, а `FINAL` через `Distributed` исполняется на таблицах шардов. diff --git a/docs/architecture/storage.md b/docs/architecture/storage.md index 680765b..283fe47 100644 --- a/docs/architecture/storage.md +++ b/docs/architecture/storage.md @@ -9,8 +9,10 @@ keeper, Kafka, каркас сервисов. Этап 2 идёт: в `sql/ddl/` лежит вся цепочка `Kafka → STG → ODS` — чтец топика `hits`, таблицы сырья, типизированное событие с таблицей ошибок, поверхность актуального состояния и три матвью. -Дальше по тексту устройство описано так, как оно проектируется; построенное от -заложенного отличает карта таблиц в конце. +Этап 3 добавил вход второго источника: топик `orders`, свой чтец и своё сырьё, +которое наполняет даг `orders_ingest`, а не матвью. Дальше по тексту устройство +описано так, как оно проектируется; построенное от заложенного отличает карта +таблиц в конце. Зона ответственности у документа одна — хранилище. Генератор описан отдельно: его замысел — в [спеке генератора](../specs/2026-08-01-generator.md), формат @@ -248,6 +250,11 @@ UTC+4), и пересчёт идёт один раз при наполнении ## Приём: поток и его свойства +Раздел — про события. У заказов приём устроен иначе: слепок забирает по команде +даг `orders_ingest`, а не подписанная навсегда матвью — [ADR +0008](../adr/0008-order-ingestion.md) и [спецификация приёма +заказов](orders/ingestion.md). + Цепочка одна: чтец топика → матвью → сырьё STG → матвью разбора → событие и таблица ошибок ODS. @@ -414,7 +421,7 @@ kafka_offset)`: смотрят такую таблицу от класса, а | Файл | Что в нём | |---|---| | `00-databases.sql` | базы слоёв | -| `10-stg-tables.sql` | Kafka-таблица, локальная и распределённая таблицы сырья | +| `10-stg-tables.sql` | чтецы топиков `hits` и `orders`, локальные и распределённые таблицы сырья обоих источников | | `20-ods-tables.sql` | типизированное событие и таблица ошибок | | `30-ods-views.sql` | актуальные события и матвью разбора в ODS | | `40-stg-views.sql` | матвью приёма: чтец в сырьё | @@ -430,10 +437,12 @@ ODS. Второе: матвью приёма создаётся последне отладке. Применение — двумя одноразовыми сервисами при `make up`, по образцу уже -работающих `airflow-init` и `superset-init`. Сначала `kafka-init` создаёт топик -`hits` с двумя партициями, затем `clickhouse-init` дожидается его завершения и -применяет файлы с ноды 1, `ON CLUSTER`. Этот порядок страхует от автосоздания -топика с одной партицией. Переключателей тут два, и путать их не надо: брокер +работающих `airflow-init` и `superset-init`. Сначала `kafka-init` создаёт топики +— `hits` с двумя партициями и `orders` с одной ([ADR +0008](../adr/0008-order-ingestion.md)), — затем `clickhouse-init` дожидается его +завершения и применяет файлы с ноды 1: всё `ON CLUSTER`, кроме чтеца заказов — +он объявлен только на этой ноде. Порядок страхует `hits` от автосоздания с +одной партицией. Переключателей тут два, и путать их не надо: брокер автосоздание разрешает, а потребитель librdkafka внутри ClickHouse его не просит — оба конца измерены, см. «Что проверено». То есть стенд держится на умолчании клиента, а урок «обе ноды читают топик» умирает тихо, поэтому топик и @@ -460,13 +469,15 @@ ODS. Второе: матвью приёма создаётся последне ## Карта таблиц -Ниже — то, что закладывает этап 2; всё перечисленное лежит в `sql/ddl/`. +Ниже — то, что закладывают этапы 2 и 3; всё перечисленное лежит в `sql/ddl/`. | Слой | Объект | Что это | |---|---|---| | STG | `stg.hits_raw_kafka` | чтец топика `hits`, формат `RawBLOB` | | STG | `stg.hits_raw_rep` / `_dist` | сырая строка сообщения плюс метаданные доставки | | STG | `stg.hits_raw_mv` | наполняет сырьё из чтеца | +| STG | `stg.orders_raw_kafka` | чтец топика `orders`, формат `RawBLOB`, только на ноде 1 и без матвью | +| STG | `stg.orders_raw_rep` / `_dist` | сырое сообщение слепка, метаданные доставки и `_load_id` | | ODS | `ods.event_rep` / `_dist` | типизированное широкое событие | | ODS | `ods.event_v` | актуальная версия события с полями источника | | ODS | `ods.event_errors_rep` / `_dist` | строки, не прошедшие строгий приём | diff --git a/sql/ddl/10-stg-tables.sql b/sql/ddl/10-stg-tables.sql index 9a163a1..6836cae 100644 --- a/sql/ddl/10-stg-tables.sql +++ b/sql/ddl/10-stg-tables.sql @@ -1,4 +1,4 @@ --- STG: чтец топика hits и таблицы сырья. +-- STG: чтецы топиков hits и orders, таблицы сырья обоих источников. -- -- Слой сырья ничего не интерпретирует: сообщение ложится строкой, как пришло, -- рядом с метаданными доставки. Довод целиком — ADR 0005, конвенции колонок и @@ -91,3 +91,56 @@ SETTINGS ttl_only_drop_parts = 1; CREATE TABLE IF NOT EXISTS stg.hits_raw_dist ON CLUSTER clickstream_cluster AS stg.hits_raw_rep ENGINE = Distributed('clickstream_cluster', 'stg', 'hits_raw_rep', cityHash64(raw)); + +-- Чтец топика orders. Матвью к нему не привязана: слепок забирает прямым +-- SELECT даг orders_ingest. Почему пулл, почему без матвью и почему у топика +-- одна партиция — ADR 0008; сам топик создаёт kafka-init в compose.yaml. +-- +-- Без ON CLUSTER: таблица нужна только на clickhouse-01 — той ноде, к которой +-- у Airflow подключение, и она же одна читает топик. +-- +-- kafka_commit_on_select — вторая настройка прямого чтения: без неё офсеты не +-- коммитятся и каждый запуск забирает один и тот же слепок заново. Первая, +-- stream_like_engine_allow_direct_select, живёт на уровне запроса и стоит в +-- самом заборе (dags/orders_ingest.py). +CREATE TABLE IF NOT EXISTS stg.orders_raw_kafka +( + raw String +) +ENGINE = Kafka +SETTINGS + kafka_broker_list = 'kafka:9092', + kafka_topic_list = 'orders', + kafka_group_name = 'clickstream_orders', + kafka_format = 'RawBLOB', + kafka_commit_on_select = 1; + +-- Локальная таблица сырья заказов. Колонки, типы, ключ, нарезка и срок жизни — +-- те же, что у сырья событий, и по тем же доводам: docs/architecture/storage.md. +-- +-- Своя колонка здесь одна — _load_id, идентификатор пачки загрузки: он равен +-- run_id прогона Airflow, который забрал порцию, и по нему разбор в ODS читает +-- неизменный срез. +CREATE TABLE IF NOT EXISTS stg.orders_raw_rep ON CLUSTER clickstream_cluster +( + raw String, + kafka_topic LowCardinality(String), + kafka_partition UInt64, + kafka_offset UInt64, + kafka_timestamp Nullable(DateTime64(3, 'UTC')), + consumer_host LowCardinality(String), + _load_id String, + _load_ts DateTime64(3, 'UTC') +) +ENGINE = ReplicatedMergeTree('/clickhouse/tables/{shard}/{database}/{table}', '{replica}') +PARTITION BY toDate(_load_ts) +ORDER BY (kafka_partition, kafka_offset) +TTL toDateTime(_load_ts) + INTERVAL 3 DAY +SETTINGS ttl_only_drop_parts = 1; + +-- Лицо слоя: пакетный забор пишет сюда, а не в локальную таблицу. Раскладку по +-- шардам обязан решать ключ шардирования, то есть свойство данных, а не то, +-- какая нода выполняла запрос, — а при пулле она всегда одна и та же. +CREATE TABLE IF NOT EXISTS stg.orders_raw_dist ON CLUSTER clickstream_cluster +AS stg.orders_raw_rep +ENGINE = Distributed('clickstream_cluster', 'stg', 'orders_raw_rep', cityHash64(raw)); -- 2.54.0