Files
ddadminandClaude Fable 5 0d9a83d6af docs(orders): устройство заказов переехало в живой набор architecture/orders
Зачем: собранная спека заказов стала базовой документацией сервиса, а
жанр спеки-события ей мал: дата в имени врёт, целиком в контекст агента
она не влезает, а трекер с резолюциями долговечным хранилищем не
считается. Решение владельца — держать детальное устройство компонента
связным набором живых документов (ADR 0011) и совместить переезд с
проходом на вычитание (#84).

Что: docs/architecture/orders/ — индекс README и файлы по частям
устройства: проекция, слепок и доставка, судьба, классы расхождений,
опись, мост к склейке, стартовый мир, правила кода; спека приёма
переехала в ingestion.md без содержательных правок. Резы вычитания по
итогам двух слепых линий: тела разделов о проводе и приёме сведены к
указателям на мастер-спеку, исследование формата и ADR (порядок строк
слепка — единственное правило, оставшееся на месте); замеры канонического
зерна и повторы-пояснения срезаны; списки отклонённых вариантов сохранены
как долговечная запись. Датированные файлы удалены, ссылки из мастер-спеки,
спеки генератора, ADR 0008/0010, исследования формата и storage.md
перенацелены; раздел «Структура» AGENTS.md дополнен правилом подпапки.

Проверка: grep по репозиторию не находит ссылок на удалённые файлы;
все относительные ссылки внутри набора разрешаются в существующие файлы;
впереди холодная сверка «ни одно решение не потеряно» и приёмка #88.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-08-16 23:50:13 +03:00

13 KiB

Приём заказов из Kafka в ODS

Учебный результат: менти различает версию бизнес-сущности, наблюдение источника и запуск загрузки, а затем читает физические версии через явную поверхность текущего состояния.

Проблема

Заказы приезжают полным слепком окна изменяемости, но Kafka передаёт его отдельными сообщениями и не сообщает потребителю, где слепок закончился. Прямое чтение Kafka Engine возвращает одну порцию. Поэтому прежняя публикация через REPLACE PARTITION snapshot_date приравнивала дату наблюдения к отсутствующей транспортной границе и могла заменить день неполным набором строк.

Одновременно типизированный ODS не должен принимать правдоподобные значения по умолчанию из грязного JSON или останавливать весь пакет из-за одной строки.

Цели

  • сохранить пакетный забор как контраст потоковому приёму событий;
  • один раз принять байты в STG и независимо разложить строки на годные и брак;
  • хранить в ODS типизированные версии заказов, не привязывая идемпотентность к snapshot_date;
  • дать следующим слоям один корректный способ прочитать текущее состояние;
  • оставить код проверки коротким и ограничить его контрактом провода.

Не входит

  • модель заказа в DDS: её зерно, связи, материализация и способ наполнения;
  • проверка полей внутри items, переходов статуса и равенств денежных сумм;
  • удаление заказа по отсутствию в следующем слепке;
  • маркер конца слепка, опись ожидаемых строк и транзакция между целями ODS.

Поток данных

После завершения генератора Airflow один раз читает байтовый Kafka-чтец и записывает полученную порцию в stg.orders_raw. Все строки получают _load_id, равный run_id Airflow. _load_ts вычисляется при этой записи и дальше переносится без пересчёта.

Один следующий task_id отвечает за весь переход STG → ODS. Внутри него два последовательных INSERT SELECT читают неизменный срез по _load_id: первый пишет годные строки в ods.order_snapshot, второй — брак в ods.order_snapshot_errors. Транзакции между запросами нет. При частичном сбое Airflow повторяет весь task_id; одинаковые исходные строки и служебные метки не вычисляются заново.

Условия запросов взаимодополняющие: один общий предикат определяет брак, а годная ветвь использует его буквальное отрицание. Все функции предиката возвращают результат без исключения, а сам предикат всегда заканчивается в true или false, не в NULL. Постоянный классификатор между STG и ODS для этого не нужен.

Граница строгого приёма

Единица решения — одна строка stg.orders_raw. Корень должен быть JSON-объектом с точным набором ключей: order_id, user_id, status, created_at, updated_at, items_total, discount, delivery, total, items, snapshot_date.

Скалярные поля проверяются по типу JSON. Деньги дополнительно обязаны быть строками с ровно двумя знаками после точки, времена — строками RFC 3339 в UTC с обязательными миллисекундами, дата слепка — строкой YYYY-MM-DD. items проверяется только как JSON-массив. Каноническая форма и основания выбора зафиксированы в исследовании формата.

Проверять все верхнеуровневые поля здесь уместно: их одиннадцать, и десять скалярных значений непосредственно образуют типизированную строку заказа. У события из 47 полей проверяются только пять опорных; переносить то сокращение на малый контракт заказа нет причины. Граница строгости заканчивается на форме провода: содержимое позиций и бизнес-инварианты намеренно остаются ниже.

Класс брака выбирается первым совпадением:

  1. not_an_object;
  2. keyset_mismatch;
  3. field_invalid.

Имя отдельного поля в класс не включается. В таблице ошибок остаются сырой текст, метаданные доставки и _load_id, поэтому единичный случай можно разобрать без постоянной детализации предиката.

Роль ODS

ods.order_snapshot_rep хранит физически принятые версии в ReplacingMergeTree(updated_at). Ключ сортировки — order_id, партиция — день неизменного created_at. ods.order_snapshot_dist шардирует по cityHash64(order_id): только так все версии заказа попадают на один шард и FINAL даёт корректный результат через распределённую таблицу.

Четыре координаты отвечают на разные вопросы:

  • updated_at — какая бизнес-версия заказа новее;
  • snapshot_date — в слепке какого модельного дня источник показал строку;
  • _load_id — какой запуск Airflow принял строку;
  • _load_ts — когда строка приехала в хранилище.

Обычное чтение _dist показывает физически сохранившиеся версии и нужно для диагностики. Их число зависит от фоновых слияний: ODS не служит архивом истории. ods.order_v сохраняет те же источник-ориентированные поля без обогащения данными модели, но возвращает одну актуальную версию на order_id. Сначала оно может быть простым представлением над _dist FINAL; способ выбора можно заменить, не меняя потребителей.

Это представление остаётся ответственностью ODS: оно скрывает механику чтения версий, но не строит бизнес-модель. DDS читает ods.order_v и отдельно решает, какие сущности, связи и производные признаки ему нужны. Нужен ли _load_id выше ODS, решается вместе с DDS, а не здесь.

Одно чтение Kafka

Стандартный слепок содержит около полутора тысяч строк, тогда как предел одной порции на стенде — десятки тысяч сообщений. Поэтому один запуск Airflow делает один прямой SELECT, без цикла до пустоты и без фиксации конечных офсетов.

Это допущение о размере стенда, а не доказательство полноты слепка. Если Kafka вернёт короткую порцию, непрочитанный хвост останется в топике и приедет в один из следующих запусков. После отказа от замены партиции это задержка, а не потеря или публикация неполного дня.

Отклонённые варианты

  • Партиционная идемпотентность — REPLACE PARTITION snapshot_date или ReplacingMergeTree по (snapshot_date, order_id): у потребителя нет признака полноты партиции, а дата наблюдения становится частью ключа сущности.
  • Обычный MergeTree в ODS с дедупликацией только в DDS: навсегда сохраняет технические повторы там, где семантика версии уже известна.
  • Маркер, опись, чтение до пустоты или конечные офсеты: добавляют протокол ради объёма, который с большим запасом помещается в одну порцию.
  • Два task_id или материализованный классификатор: дробят один короткий переход слоя, не добавляя транзакционности.
  • Представление текущего состояния в DDS и готовая схема dds.order в этой задаче: перекладывают механику ODS на следующий слой и преждевременно задают модель данных.

Риски и проверка

  • На стандартном мире сверить число отправленных заказов с числом строк, принятых одним прямым чтением. Это разовая приёмка допущения, не постоянный сторож.
  • На малой управляемой порции дать по одной строке каждого класса брака и две годные версии одного order_id. Две цели должны сохранить все непустые сообщения, а ods.order_v — вернуть новую версию независимо от фонового слияния.
  • Повторить переход с тем же _load_id: строка в ods.order_v и её _load_ts не должны измениться; версии одного заказа должны остаться на одном шарде и в одной партиции. Таблица ошибок может снова записать тот же брак: совпавшие _load_id и Kafka-координаты показывают повтор задачи.

Что проверено

MCP Context7 в сессии проектирования был недоступен. На локальном ClickHouse 26.3.17.56 проверено, что прямой SELECT Kafka Engine завершается после одной порции, а FINAL через Distributed исполняется на таблицах шардов. Поэтому версии одного order_id направляются на один шард. Фоновое схлопывание ReplacingMergeTree и необходимость точного чтения сверены с официальной документацией, поведение чтения — с исходниками той же версии StorageKafka.cpp и KafkaSource.cpp.

Связанные решения

  • ADR 0008 сохраняет выбор пакетного забора, RawBLOB, одного чтеца и одной партиции топика.
  • ADR 0010 заменяет публикацию слепка версионным ODS.
  • Вопрос _load_id выше ODS оставлен проектированию DDS в тикете #85.