feat(orders): добавить приём слепков в STG #101

Merged
ddmitry merged 1 commits from feat/93-orders-ingestion into main 2026-08-18 21:49:10 +03:00
10 changed files with 330 additions and 54 deletions
+6
View File
@@ -233,6 +233,12 @@ uv run --project generator python -m clickstream_generator batch \
| `world_live_day` | играет следующий день в темпе модельного времени, около двадцати четырёх минут | | `world_live_day` | играет следующий день в темпе модельного времени, около двадцати четырёх минут |
| `world_live` | выключатель: пока снят с паузы, дёргает `world_live_day` день за днём | | `world_live` | выключатель: пока снят с паузы, дёргает `world_live_day` день за днём |
Работник не только играет день: следом он отправляет слепок заказов за
сыгранное и дёргает отдельный даг, `orders_ingest`, — тот забирает приехавший
слепок из топика `orders` в сырьё, — и дожидается его конца. Так две половины
стенда идут в ногу: упал приём — краснеет и работник, а позиция остаётся на
месте ([ADR 0008](docs/adr/0008-order-ingestion.md)).
Позицию хранит переменная Airflow `world_position` — номер первого несыгранного Позицию хранит переменная Airflow `world_position` — номер первого несыгранного
дня. Ставит её работник, сыгравший день, и только по успеху: оборванный прогон дня. Ставит её работник, сыгравший день, и только по успеху: оборванный прогон
позицию не двигает, и следующий запуск играет тот же день заново. Нет позицию не двигает, и следующий запуск играет тот же день заново. Нет
+11
View File
@@ -68,6 +68,7 @@ x-airflow-common: &airflow-common
GENERATOR_IMAGE: *generator-image GENERATOR_IMAGE: *generator-image
STAND_NETWORK: ${COMPOSE_PROJECT_NAME}_default STAND_NETWORK: ${COMPOSE_PROJECT_NAME}_default
WORLD_STARTING_DAYS: *world-starting-days WORLD_STARTING_DAYS: *world-starting-days
KAFKA_ORDERS_TOPIC: orders
volumes: volumes:
- ./dags:/opt/airflow/dags:ro - ./dags:/opt/airflow/dags:ro
- ./infra/airflow/init.sh:/opt/airflow/init.sh:ro - ./infra/airflow/init.sh:/opt/airflow/init.sh:ro
@@ -238,6 +239,16 @@ services:
exit 1 exit 1
fi 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 по порядку имён файлов после готовности всего кластера. # DDL применяется с ноды 1 по порядку имён файлов после готовности всего кластера.
clickhouse-init: clickhouse-init:
image: clickhouse/clickhouse-server:26.3.17.56 image: clickhouse/clickhouse-server:26.3.17.56
+100
View File
@@ -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()
+55 -11
View File
@@ -7,6 +7,11 @@
конца. Снят с паузы — мир едет день за днём; поставлен на паузу — встал на конца. Снят с паузы — мир едет день за днём; поставлен на паузу — встал на
границе модельных суток. границе модельных суток.
Сыграть день — половина работы. Вторая половина: отправить слепок заказов и
дождаться, пока хранилище его заберёт. Поэтому у обоих работников за проигрышем
идут отправка слепка тем же контейнером генератора и ждущий триггер дага приёма
`orders_ingest`.
Разделение не косметическое. Расписание на самом работнике заставило бы кнопку Разделение не косметическое. Расписание на самом работнике заставило бы кнопку
паузы значить две вещи разом — «мир не едет сам» и «даг выключен», — а работник паузы значить две вещи разом — «мир не едет сам» и «даг выключен», — а работник
при этом выглядел бы в списке выключенным, хотя нажимают его каждый день. при этом выглядел бы в списке выключенным, хотя нажимают его каждый день.
@@ -35,6 +40,9 @@ GENERATOR_ENVIRONMENT = {
"KAFKA_BOOTSTRAP_SERVERS": os.environ["KAFKA_BOOTSTRAP_SERVERS"], "KAFKA_BOOTSTRAP_SERVERS": os.environ["KAFKA_BOOTSTRAP_SERVERS"],
"KAFKA_TOPIC": os.environ["KAFKA_TOPIC"], "KAFKA_TOPIC": os.environ["KAFKA_TOPIC"],
} }
# Топик слепка называется аргументом, а не окружением: KAFKA_TOPIC выше — топик
# событий, и промолчи мы, заказы уехали бы к ним.
ORDERS_TOPIC = os.environ["KAFKA_ORDERS_TOPIC"]
STARTING_DAYS = int(os.environ["WORLD_STARTING_DAYS"]) STARTING_DAYS = int(os.environ["WORLD_STARTING_DAYS"])
# Позиция на оси: номер первого несыгранного дня. Переменной нет — мир в # Позиция на оси: номер первого несыгранного дня. Переменной нет — мир в
@@ -50,6 +58,10 @@ WORLD_POSITION = "world_position"
# двадцати четырёх минут: тик только спрашивает «не пора ли снова». # двадцати четырёх минут: тик только спрашивает «не пора ли снова».
LIVE_TICK = datetime.timedelta(minutes=25) 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) START_DATE = datetime.datetime(2026, 1, 1, tzinfo=datetime.UTC)
TAGS = ["пульт мира"] TAGS = ["пульт мира"]
@@ -60,8 +72,8 @@ def first_unplayed_day() -> int:
return int(Variable.get(WORLD_POSITION, default=STARTING_DAYS)) 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, а Внутрь образа Airflow генератор не поставить: он требует Python 3.14, а
образ несёт 3.13. Да и обещание побайтовой воспроизводимости дано для образ несёт 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(
dag_id="world_next_day", dag_id="world_next_day",
schedule=None, schedule=None,
@@ -90,7 +118,11 @@ def _play(task_id: str, command: list[str]) -> DockerOperator:
is_paused_upon_creation=False, is_paused_upon_creation=False,
max_active_runs=1, max_active_runs=1,
tags=TAGS, 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(): def world_next_day():
"""Прожить следующие дни пачкой, без пауз. """Прожить следующие дни пачкой, без пауз.
@@ -106,18 +138,27 @@ def world_next_day():
Variable.set(WORLD_POSITION, str(first_day + days)) Variable.set(WORLD_POSITION, str(first_day + days))
first_day = first_unplayed_day() first_day = first_unplayed_day()
played = _play( played = _generator(
"play_days", "play_days",
["batch", "--day", PLAYED_DAY, "--days", "{{ params.days }}"],
)
# Разгон на N дней отправляет N слепков — по одному за сыгранный день, и
# каждый со своим сдвигом: прогон дня D везёт слепок дня D−1. Забор при
# этом остаётся один.
sent = _generator(
"send_snapshots",
[ [
"batch", "snapshot",
"--day", "--day",
"{{ ti.xcom_pull(task_ids='first_unplayed_day') }}", PLAYED_DAY,
"--days", "--days",
"{{ params.days }}", "{{ params.days }}",
"--topic",
ORDERS_TOPIC,
], ],
) )
first_day >> played >> remember_played(first_day) first_day >> played >> sent >> _ingest_orders() >> remember_played(first_day)
@dag( @dag(
@@ -142,12 +183,15 @@ def world_live_day():
Variable.set(WORLD_POSITION, str(first_day + 1)) Variable.set(WORLD_POSITION, str(first_day + 1))
first_day = first_unplayed_day() first_day = first_unplayed_day()
played = _play( played = _generator("play_day", ["live", "--day", PLAYED_DAY])
"play_day", # Слепок уезжает пачкой и после дня, а не в темпе: у выгрузки бэкенда темпа
["live", "--day", "{{ ti.xcom_pull(task_ids='first_unplayed_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( @dag(
+58 -24
View File
@@ -1,13 +1,15 @@
# ADR 0008. Приём заказов: пакетный забор слепка, инициируемый Airflow # ADR 0008. Приём заказов: пакетный забор слепка, инициируемый Airflow
Дата: 12 августа 2026 года. Статус: частично заменено Дата: 12 августа 2026 года. Статус: частично заменено
[ADR 0010](0010-order-versions-in-ods.md). Реализация — отдельным тикетом. [ADR 0010](0010-order-versions-in-ods.md); забор построен тикетом #93.
Сохраняются пакетный забор, `RawBLOB`, один чтец на `clickhouse-01`, одна Сохраняются пакетный забор, `RawBLOB`, один чтец на `clickhouse-01`, одна
партиция топика и отсутствие матвью. Отменены граница «одно чтение — полный партиция топика и отсутствие матвью. Отменены граница «одно чтение — полный
слепок», замена партиции `snapshot_date`, ODS без дедупликации и готовая схема слепок», замена партиции `snapshot_date`, ODS без дедупликации и готовая схема
`dds.order` с `argMax`. Действующая форма STG → ODS описана в `dds.order` с `argMax` — это ADR 0010. Отдельно от него, решением #93 от
[спецификации приёма заказов](../architecture/orders/ingestion.md). 17 августа 2026 года, производитель и потребитель слепка разъехались по разным
дагам; раздел «Решение» описывает построенное. Действующая форма STG → ODS
описана в [спецификации приёма заказов](../architecture/orders/ingestion.md).
## Решение ## Решение
@@ -15,13 +17,18 @@
у событий ([ADR 0005](0005-event-ingestion.md)). Матвью к ней не привязана. у событий ([ADR 0005](0005-event-ingestion.md)). Матвью к ней не привязана.
Сырьё забирает пакетный шаг, которым управляет Airflow: он вставляет Сырьё забирает пакетный шаг, которым управляет Airflow: он вставляет
прочитанное в `stg.orders_raw_dist`, разбирает его в типизированный слепок и прочитанное в `stg.orders_raw_dist`, разбирает его в типизированный слепок и
заменяет партицию дня в `ods.order_snapshot`. Тот же даг проигрывает модельный заменяет партицию дня в `ods.order_snapshot`.
день генератором, поэтому переливается ровно то, что он положил в топик.
Уточнение со сдвигом отправки, решённым позже (#71, [слепок и его Производитель и потребитель слепка — разные даги. Работник пульта мира играет
доставка](../architecture/orders/snapshot.md)): даг, играющий день D, в день и отправляет слепок, а забирает его отдельный даг `orders_ingest`; зовёт
штатном прогоне кладёт и забирает слепок дня D−1 — слепок предыдущего дня, а забор сам работник ждущим `TriggerDagRunOperator` и зеленеет только вслед за
не сыгранного. После падения между шагами в топике может ждать и хвост ним. Связывает две половины не общий даг, а этот дождавшийся успех: зелёный
прежних слепков; забор принимает всё приехавшее. работник значит, что забор отработал. Уточнение со сдвигом
отправки, решённым позже (#71, [слепок и его
доставка](../architecture/orders/snapshot.md)): работник, играющий день D, в
штатном прогоне кладёт слепок дня D−1 — слепок предыдущего дня, а не
сыгранного. После падения между шагами в топике может ждать и хвост прежних
слепков; забор принимает всё приехавшее.
Прямое чтение из Kafka-движка требует двух настроек, и вторая не очевидна: Прямое чтение из Kafka-движка требует двух настроек, и вторая не очевидна:
`stream_like_engine_allow_direct_select = 1` разрешает читать чтеца запросом, `stream_like_engine_allow_direct_select = 1` разрешает читать чтеца запросом,
@@ -47,8 +54,8 @@
партиции** `snapshot_date`, ровно как обещает раздел 2 мастер-спеки; партиции** `snapshot_date`, ровно как обещает раздел 2 мастер-спеки;
- `dds.order` — дедуп до последней версии через `argMax`, без изменений. - `dds.order` — дедуп до последней версии через `argMax`, без изменений.
Сенсора дневного батча нет. Ждать нечего: производитель и потребитель слепка Сенсора дневного батча нет. Ждать нечего: забор дёргает сам отправитель
живут в одном даге. слепка, когда отправил.
## Почему ## Почему
@@ -112,7 +119,7 @@
**Что отвергнуто ещё.** Сенсор дневного батча из раздела 9 мастер-спеки: у **Что отвергнуто ещё.** Сенсор дневного батча из раздела 9 мастер-спеки: у
топика нет сигнала «всё», и сенсор ловил бы момент, которого не существует, — топика нет сигнала «всё», и сенсор ловил бы момент, которого не существует, —
а раз генератор и переливка в одном даге, ждать нечего по построению. Лаба, а раз забор зовёт сам отправитель слепка, ждать нечего по построению. Лаба,
поднимающая типизированного чтеца во второй группе потребителей ради того же поднимающая типизированного чтеца во второй группе потребителей ради того же
сравнения: она понадобилась бы, останься сравнение невыполненным, но push сравнения: она понадобилась бы, останься сравнение невыполненным, но push
против pull даёт его живым и постоянным. против pull даёт его живым и постоянным.
@@ -124,10 +131,13 @@
один и тот же, и это честная разница двух режимов, а не потеря: при пулле один и тот же, и это честная разница двух режимов, а не потеря: при пулле
читает тот, кого спросили. читает тот, кого спросили.
Гарантий приёма у заказов не больше, чем у событий, но последствия мягче. Гарантий приёма у заказов не больше, чем у событий, и граница проходит по
Офсеты коммитятся при чтении, вставка идёт следом — упавшая вставка теряет живому. Успешное чтение короткой порции оставляет непрочитанный хвост в
пачку. У потока такая потеря невосстановима, у слепка её лечит следующий день: топике — он дождётся следующего забора. А отказ после чтения теряет саму
окно изменяемости K = 7 привезёт те же заказы заново. порцию: офсеты закоммичены в момент чтения, часть строк могла лечь на шарды,
и повторное чтение вернёт ноль. Позицию мира это не двигает, поэтому день
сыграется и уедет заново; в сырье он тогда окажется частичным дублем, а
старый хвост, ушедший той же порцией, не вернётся.
Этап 3 забирает у этапа 5 первый настоящий даг. Раздел 9 мастер-спеки отдавал Этап 3 забирает у этапа 5 первый настоящий даг. Раздел 9 мастер-спеки отдавал
даги этапу 5 целиком; приём заказов без дага не существует, поэтому порядок даги этапу 5 целиком; приём заказов без дага не существует, поэтому порядок
@@ -158,10 +168,34 @@ ClickHouse `26.3.17.56`.
- Про `Distributed` поверх Kafka документация не говорит ничего — ни - Про `Distributed` поверх Kafka документация не говорит ничего — ни
поддержки, ни запрета. поддержки, ни запрета.
Осталось проверить при исполнении, и это работа тикета реализации: что второй Забор проверен на живом стенде 18 августа 2026 года при исполнении #93
прогон подряд возвращает пусто, то есть офсеты действительно закоммичены; что ClickHouse `26.3.17.56`, Airflow 3.3.0. Три вопроса, оставленные этим ADR
одного чтения хватает на весь слепок дня; что виртуальные колонки доставки реализации, закрыты; заодно снят отказной путь. Обе настройки прямого чтения
(`_topic`, `_partition`, `_offset`, `_timestamp_ms`) доступны при прямом чтении доезжают до запроса, объявленные в конце `INSERT ... SELECT`: `system.query_log`
— на них стоят служебные колонки сырья, см. [доку показывает у каждой вставки единицу.
хранилища](../architecture/storage.md); что повторная заливка дня даёт в
`ods.order_snapshot` тот же счёт, а не удвоенный. - **Виртуальные колонки доставки при прямом чтении доступны все четыре.**
`_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),
раздел «Риски и проверка».
+8 -2
View File
@@ -8,8 +8,10 @@
**Работники** — без расписания и без паузы, запускаются руками: **Работники** — без расписания и без паузы, запускаются руками:
- `world_next_day` — день пачкой; параметр «сколько дней» (умолчание 1) даёт - `world_next_day` — день пачкой; параметр «сколько дней» (умолчание 1, потолок
разгон вперёд одним нажимом; 7) даёт разгон вперёд одним нажимом. Потолок — это и есть названное желание
«уедь на неделю сейчас», и ровно тот разгон, который забор заказов забирает
одним чтением ([ADR 0008](0008-order-ingestion.md));
- `world_live_day` — один день в темпе живого режима. - `world_live_day` — один день в темпе живого режима.
**Выключатель** — без своей работы: `world_live` несёт расписание около двадцати **Выключатель** — без своей работы: `world_live` несёт расписание около двадцати
@@ -99,6 +101,10 @@
которого стенд не поднимает; заводить службу ради одного ожидания дороже которого стенд не поднимает; заводить службу ради одного ожидания дороже
занятого слота. занятого слота.
Приросший к работникам забор заказов (#93, [ADR 0008](0008-order-ingestion.md))
делает из двух занятых слотов три: пока идёт забор, ждут выключатель и сам
работник, а третий слот занимает задача `orders_ingest`.
**Пакетная автоматика отложена, потому что своего желания у неё пока одно.** **Пакетная автоматика отложена, потому что своего желания у неё пока одно.**
«Шагни и стой» закрывает кнопка, «уедь на неделю сейчас» — параметр «сколько «Шагни и стой» закрывает кнопка, «уедь на неделю сейчас» — параметр «сколько
дней», «живи, пока я смотрю» — выключатель. Ей остаётся «едь сам быстрее, чем дней», «живи, пока я смотрю» — выключатель. Ей остаётся «едь сам быстрее, чем
+3 -2
View File
@@ -57,8 +57,9 @@
доли классов, стоимость доставки — калибровка при реализации; финальная доли классов, стоимость доставки — калибровка при реализации; финальная
фиксация чисел — пересборка эталонного мира, этап 7. При пересборке правки фиксация чисел — пересборка эталонного мира, этап 7. При пересборке правки
потребуют только числа, не устройство. потребуют только числа, не устройство.
- **Проверки приёма**разовая приёмка допущения «один запуск — одно чтение» - **Проверки приёма**опыты из [«Рисков и проверки»](ingestion.md) про брак
и опыты из [«Рисков и проверки»](ingestion.md) — тикет реализации приёма. и версии в ODS — тикет перехода STG → ODS (#94). Допущение «один запуск —
одно чтение» принято живым прогоном при исполнении #93.
- **`_load_id` выше ODS** — вместе с устройством `dds.order` (#85). - **`_load_id` выше ODS** — вместе с устройством `dds.order` (#85).
- **Контур проверок качества для расхождений** (даг DQ, `dm.dq_summary`) — - **Контур проверок качества для расхождений** (даг DQ, `dm.dq_summary`) —
остаётся в тумане карты #69; естественное место разговора — этап 4. остаётся в тумане карты #69; естественное место разговора — этап 4.
+16 -6
View File
@@ -34,9 +34,10 @@
## Поток данных ## Поток данных
После завершения генератора Airflow один раз читает байтовый Kafka-чтец и После завершения генератора Airflow один раз читает байтовый Kafka-чтец и
записывает полученную порцию в `stg.orders_raw`. Все строки получают `_load_id`, записывает полученную порцию в `stg.orders_raw`. Даг зовётся `orders_ingest`, а
равный `run_id` Airflow. `_load_ts` вычисляется при этой записи и дальше дёргает его тот, кто положил слепок в топик, — работник пульта мира, — и ждёт
переносится без пересчёта. конца прогона. Все строки получают `_load_id`, равный `run_id` Airflow.
`_load_ts` вычисляется при этой записи и дальше переносится без пересчёта.
Один следующий `task_id` отвечает за весь переход STG → ODS. Внутри него два Один следующий `task_id` отвечает за весь переход STG → ODS. Внутри него два
последовательных `INSERT SELECT` читают неизменный срез по `_load_id`: первый последовательных `INSERT SELECT` читают неизменный срез по `_load_id`: первый
@@ -119,6 +120,13 @@ JSON-объектом с точным набором ключей: `order_id`, `
из следующих запусков. После отказа от замены партиции это задержка, а не потеря из следующих запусков. После отказа от замены партиции это задержка, а не потеря
или публикация неполного дня. или публикация неполного дня.
Отказ — случай другой, и «хвост дождётся» на него не распространяется. Офсеты
порции коммитятся в момент чтения, поэтому упавшая вставка уносит прочитанное с
собой: повторное чтение вернёт ноль, а часть строк может уже лежать на шарде.
Позиция мира при этом не двигается, и тот же день уедет заново — в сырье он
окажется частичным дублем, а старый хвост, ушедший той же порцией, не вернётся
([ADR 0008](../../adr/0008-order-ingestion.md), «Следствия»).
## Отклонённые варианты ## Отклонённые варианты
- Партиционная идемпотентность — `REPLACE PARTITION snapshot_date` или - Партиционная идемпотентность — `REPLACE PARTITION snapshot_date` или
@@ -137,9 +145,6 @@ JSON-объектом с точным набором ключей: `order_id`, `
## Риски и проверка ## Риски и проверка
- На стандартном мире сверить число отправленных заказов с числом строк,
принятых одним прямым чтением. Это разовая приёмка допущения, не постоянный
сторож.
- На малой управляемой порции дать по одной строке каждого класса брака и две - На малой управляемой порции дать по одной строке каждого класса брака и две
годные версии одного `order_id`. Две цели должны сохранить все непустые годные версии одного `order_id`. Две цели должны сохранить все непустые
сообщения, а `ods.order_v` — вернуть новую версию независимо от фонового сообщения, а `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 MCP Context7 в сессии проектирования был недоступен. На локальном ClickHouse
`26.3.17.56` проверено, что прямой `SELECT` Kafka Engine завершается после `26.3.17.56` проверено, что прямой `SELECT` Kafka Engine завершается после
одной порции, а `FINAL` через `Distributed` исполняется на таблицах шардов. одной порции, а `FINAL` через `Distributed` исполняется на таблицах шардов.
+19 -8
View File
@@ -9,8 +9,10 @@
keeper, Kafka, каркас сервисов. Этап 2 идёт: в `sql/ddl/` лежит вся цепочка keeper, Kafka, каркас сервисов. Этап 2 идёт: в `sql/ddl/` лежит вся цепочка
`Kafka → STG → ODS` — чтец топика `hits`, таблицы сырья, типизированное `Kafka → STG → ODS` — чтец топика `hits`, таблицы сырья, типизированное
событие с таблицей ошибок, поверхность актуального состояния и три матвью. событие с таблицей ошибок, поверхность актуального состояния и три матвью.
Дальше по тексту устройство описано так, как оно проектируется; построенное от Этап 3 добавил вход второго источника: топик `orders`, свой чтец и своё сырьё,
заложенного отличает карта таблиц в конце. которое наполняет даг `orders_ingest`, а не матвью. Дальше по тексту устройство
описано так, как оно проектируется; построенное от заложенного отличает карта
таблиц в конце.
Зона ответственности у документа одна — хранилище. Генератор описан отдельно: Зона ответственности у документа одна — хранилище. Генератор описан отдельно:
его замысел — в [спеке генератора](../specs/2026-08-01-generator.md), формат его замысел — в [спеке генератора](../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 → матвью разбора → событие и Цепочка одна: чтец топика → матвью → сырьё STG → матвью разбора → событие и
таблица ошибок ODS. таблица ошибок ODS.
@@ -414,7 +421,7 @@ kafka_offset)`: смотрят такую таблицу от класса, а
| Файл | Что в нём | | Файл | Что в нём |
|---|---| |---|---|
| `00-databases.sql` | базы слоёв | | `00-databases.sql` | базы слоёв |
| `10-stg-tables.sql` | Kafka-таблица, локальная и распределённая таблицы сырья | | `10-stg-tables.sql` | чтецы топиков `hits` и `orders`, локальные и распределённые таблицы сырья обоих источников |
| `20-ods-tables.sql` | типизированное событие и таблица ошибок | | `20-ods-tables.sql` | типизированное событие и таблица ошибок |
| `30-ods-views.sql` | актуальные события и матвью разбора в ODS | | `30-ods-views.sql` | актуальные события и матвью разбора в ODS |
| `40-stg-views.sql` | матвью приёма: чтец в сырьё | | `40-stg-views.sql` | матвью приёма: чтец в сырьё |
@@ -430,10 +437,12 @@ ODS. Второе: матвью приёма создаётся последне
отладке. отладке.
Применение — двумя одноразовыми сервисами при `make up`, по образцу уже Применение — двумя одноразовыми сервисами при `make up`, по образцу уже
работающих `airflow-init` и `superset-init`. Сначала `kafka-init` создаёт топик работающих `airflow-init` и `superset-init`. Сначала `kafka-init` создаёт топики
`hits` с двумя партициями, затем `clickhouse-init` дожидается его завершения и `hits` с двумя партициями и `orders` с одной ([ADR
применяет файлы с ноды 1, `ON CLUSTER`. Этот порядок страхует от автосоздания 0008](../adr/0008-order-ingestion.md)), — затем `clickhouse-init` дожидается его
топика с одной партицией. Переключателей тут два, и путать их не надо: брокер завершения и применяет файлы с ноды 1: всё `ON CLUSTER`, кроме чтеца заказов —
он объявлен только на этой ноде. Порядок страхует `hits` от автосоздания с
одной партицией. Переключателей тут два, и путать их не надо: брокер
автосоздание разрешает, а потребитель librdkafka внутри ClickHouse его не автосоздание разрешает, а потребитель 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_kafka` | чтец топика `hits`, формат `RawBLOB` |
| STG | `stg.hits_raw_rep` / `_dist` | сырая строка сообщения плюс метаданные доставки | | STG | `stg.hits_raw_rep` / `_dist` | сырая строка сообщения плюс метаданные доставки |
| STG | `stg.hits_raw_mv` | наполняет сырьё из чтеца | | 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_rep` / `_dist` | типизированное широкое событие |
| ODS | `ods.event_v` | актуальная версия события с полями источника | | ODS | `ods.event_v` | актуальная версия события с полями источника |
| ODS | `ods.event_errors_rep` / `_dist` | строки, не прошедшие строгий приём | | ODS | `ods.event_errors_rep` / `_dist` | строки, не прошедшие строгий приём |
+54 -1
View File
@@ -1,4 +1,4 @@
-- STG: чтец топика hits и таблицы сырья. -- STG: чтецы топиков hits и orders, таблицы сырья обоих источников.
-- --
-- Слой сырья ничего не интерпретирует: сообщение ложится строкой, как пришло, -- Слой сырья ничего не интерпретирует: сообщение ложится строкой, как пришло,
-- рядом с метаданными доставки. Довод целиком — ADR 0005, конвенции колонок и -- рядом с метаданными доставки. Довод целиком — 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 CREATE TABLE IF NOT EXISTS stg.hits_raw_dist ON CLUSTER clickstream_cluster
AS stg.hits_raw_rep AS stg.hits_raw_rep
ENGINE = Distributed('clickstream_cluster', 'stg', 'hits_raw_rep', cityHash64(raw)); 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));