Files
ddadminandClaude Opus 5 1fb219f836 docs(adr): решён приём заказов — пакетный забор слепка
- Зачем:
  - развилка этапа 3 стояла нерешённой прямо в разделе 7 мастер-спеки: нужен
    ли слепку слой сырья и как заказы попадают из топика в хранилище (#70).
- Что:
  - заведён ADR 0008 — байтовый чтец без матвью, слой сырья у заказов
    остаётся, в ods.order_snapshot пишет шаг Airflow заменой партиции; топик
    orders в одну партицию, чтец на clickhouse-01 без ON CLUSTER.
  - мастер-спека приведена в соответствие, разделы 6, 7, 9, 11, 12: сравнение
    двух приёмов переписано на «поток против слепка», сенсор дневного батча
    снят, первый даг переехал с этапа 5 на этап 3.
  - в CONTEXT.md заведены «слепок», «окно изменяемости», «пакетный забор».
- Проверка:
  - решение сверено по документации ClickHouse через MCP Context7 12 августа
    2026 года; что осталось замерить на стенде — списком в конце ADR 0008.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-12 21:57:37 +03:00

14 KiB
Raw Permalink Blame History

ADR 0008. Приём заказов: пакетный забор слепка, инициируемый Airflow

Дата: 12 августа 2026 года. Статус: принято. Реализация — отдельным тикетом.

Решение

Топик orders читает Kafka-таблица формата RawBLOB — та же форма чтеца, что у событий (ADR 0005). Матвью к ней не привязана. Сырьё забирает пакетный шаг, которым управляет Airflow: он вставляет прочитанное в stg.orders_raw_dist, разбирает его в типизированный слепок и заменяет партицию дня в ods.order_snapshot. Тот же даг проигрывает модельный день генератором, поэтому переливается ровно то, что он положил в топик.

Прямое чтение из 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, RabbitMQ, FileLog) запрещено начиная с версии 21.12 и открывается настройкой stream_like_engine_allow_direct_select.
  • При привязанной матвью прямое чтение остаётся запрещённым и с этой настройкой. Отсюда вывод, что два способа приёма взаимоисключающи.
  • Прямое чтение офсеты по умолчанию не коммитит; коммит включается настройкой kafka_commit_on_select на самой таблице. Это тот подводный камень, который молчит на первом прогоне и вылезает на втором.
  • Сколько строк отдаёт одно чтение, задаёт kafka_max_block_size. При слепке порядка полутора тысяч строк это один блок с запасом.
  • Про Distributed поверх Kafka документация не говорит ничего — ни поддержки, ни запрета.

Осталось проверить при исполнении, и это работа тикета реализации: что второй прогон подряд возвращает пусто, то есть офсеты действительно закоммичены; что одного чтения хватает на весь слепок дня; что виртуальные колонки доставки (_topic, _partition, _offset, _timestamp_ms) доступны при прямом чтении — на них стоят служебные колонки сырья, см. доку хранилища; что повторная заливка дня даёт в ods.order_snapshot тот же счёт, а не удвоенный.