Зачем: собранная спека заказов стала базовой документацией сервиса, а жанр спеки-события ей мал: дата в имени врёт, целиком в контекст агента она не влезает, а трекер с резолюциями долговечным хранилищем не считается. Решение владельца — держать детальное устройство компонента связным набором живых документов (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>
16 KiB
ADR 0008. Приём заказов: пакетный забор слепка, инициируемый Airflow
Дата: 12 августа 2026 года. Статус: частично заменено ADR 0010. Реализация — отдельным тикетом.
Сохраняются пакетный забор, RawBLOB, один чтец на clickhouse-01, одна
партиция топика и отсутствие матвью. Отменены граница «одно чтение — полный
слепок», замена партиции snapshot_date, ODS без дедупликации и готовая схема
dds.order с argMax. Действующая форма STG → ODS описана в
спецификации приёма заказов.
Решение
Топик orders читает Kafka-таблица формата RawBLOB — та же форма чтеца, что
у событий (ADR 0005). Матвью к ней не привязана.
Сырьё забирает пакетный шаг, которым управляет Airflow: он вставляет
прочитанное в stg.orders_raw_dist, разбирает его в типизированный слепок и
заменяет партицию дня в ods.order_snapshot. Тот же даг проигрывает модельный
день генератором, поэтому переливается ровно то, что он положил в топик.
Уточнение со сдвигом отправки, решённым позже (#71, слепок и его
доставка): даг, играющий день D, в
штатном прогоне кладёт и забирает слепок дня D−1 — слепок предыдущего дня, а
не сыгранного. После падения между шагами в топике может ждать и хвост
прежних слепков; забор принимает всё приехавшее.
Прямое чтение из 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 дополнительно снят 16 августа на локальном
ClickHouse 26.3.17.56.
- Прямое чтение из движков-очередей (Kafka, RabbitMQ, FileLog) запрещено
начиная с версии 21.12 и открывается настройкой
stream_like_engine_allow_direct_select. - При привязанной матвью прямое чтение остаётся запрещённым и с этой настройкой. Отсюда вывод, что два способа приёма взаимоисключающи.
- Прямое чтение офсеты по умолчанию не коммитит; коммит включается
настройкой
kafka_commit_on_selectна самой таблице. Это тот подводный камень, который молчит на первом прогоне и вылезает на втором. - Прямое чтение возвращает одну порцию, полученную одним опросом Kafka. При настройках стенда её предел — 65 409 сообщений, поэтому слепок порядка полутора тысяч строк помещается с запасом.
- Про
Distributedповерх Kafka документация не говорит ничего — ни поддержки, ни запрета.
Осталось проверить при исполнении, и это работа тикета реализации: что второй
прогон подряд возвращает пусто, то есть офсеты действительно закоммичены; что
одного чтения хватает на весь слепок дня; что виртуальные колонки доставки
(_topic, _partition, _offset, _timestamp_ms) доступны при прямом чтении
— на них стоят служебные колонки сырья, см. доку
хранилища; что повторная заливка дня даёт в
ods.order_snapshot тот же счёт, а не удвоенный.