Зачем: датированная спека — событие истории, а устройство компонента — живой документ; большое полотно плохо грузится и агентом, и человеком (ADR 0011). Что: вычитание #84 слито с переустройством формы: набор docs/architecture/orders/ — индекс README и семь файлов по частям устройства (нарезка по правилу «семь плюс-минус два»); датированные файлы удалены, ссылки перенацелены, AGENTS.md дополнен правилом подпапки. Приёмка владельцем #88 пройдена, черновой статус снят из README. Проверка: холодная сверка миграции свежим тредом — потерь решений нет; обход относительных ссылок набора — битых нет. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
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 полей проверяются только пять опорных; переносить то сокращение на малый контракт заказа нет причины. Граница строгости заканчивается на форме провода: содержимое позиций и бизнес-инварианты намеренно остаются ниже.
Класс брака выбирается первым совпадением:
not_an_object;keyset_mismatch;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.