- Зачем:
- холодное ревью связности нашло девять мест, где вставленный текст спорит с
соседним; отдельно вскрылось, что представление дат в JSON не зафиксировано
нигде, а #43 обязан его знать раньше, чем #41 напишет сериализатор.
- Что:
- гарантия приёма переписана: после снятия синхронной вставки «хотя бы один
раз» стало неправдой — есть и окно потери, и окно дубля.
- критерий выбора пяти опорных колонок приведён к списку, который он
порождает; `CounterID` оговорён отдельно.
- «переобработки у ODS нет вовсе» смягчено до пакетной: ручная вставка из
сырья в пределах окна возможна.
- в спеку генератора добавлена форма дат на проводе — ISO-8601, с доводом от
читаемости слоя сырья.
- убраны осиротевшая фраза про порядок сервисов, дубль порядка классов брака,
устаревшая датировка сверки и ещё три следа вставок.
- Проверка:
- make config-test
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
18 KiB
ADR 0005. Приём событий: сырьё в STG байтами, разбор функциями в ODS
Дата: 3 августа 2026 года. Статус: принято.
Решение
Топик hits читает одна Kafka-таблица формата RawBLOB: сообщение ложится в
stg.hits_raw_dist строкой, как пришло, рядом с метаданными доставки. Ни
типизации, ни проверки на этом шаге нет — слой сырья ничего не интерпретирует.
У формата есть следствие для DDL: он читает вход в одно значение и рассчитан на
таблицу с единственной колонкой, поэтому у чтеца она ровно одна — raw, а
метаданные доставки берутся только из виртуальных колонок и добавить к чтецу
что-либо своё нельзя.
Типизированный слой наполняют две матвью, привязанные к stg.hits_raw_dist.
Поля достаются JSONExtract. В таблицу ошибок уходят три класса брака:
сообщение, не являющееся объектом JSON; объект, чей набор ключей разошёлся с
контрактным; объект, у которого не разобрался ключевой идентификатор или метка
времени. Остальное — в событие. Первый класс проверяется именно на объект, а не
на валидность: isValidJSON('123') возвращает единицу, скаляр — тоже законный
JSON.
Классы пересекаются: скаляр проваливает заодно и сверку ключей, потому что
JSONExtractKeys от него даёт пустой массив. Поэтому они проверяются по порядку,
а в колонку error_class пишется первый совпавший — not_an_object,
keyset_mismatch, key_field_unparsed. Приём тот же, что у mismatch_class в
витрине сверки: пересекающиеся классы плюс объявленный приоритет.
Присутствие полей целиком держит сверка набора ключей — одно сравнение
arraySort(JSONExtractKeys(raw)) с контрактным списком, завёрнутым в тот же
arraySort. Обёртка с обеих сторон стоит ноль и снимает ошибку, которая иначе
увела бы в брак вообще всё: сорок семь CamelCase-имён, выписанных руками ровно в
байтовом порядке. Обязательны все сорок
семь полей: генератор шлёт их все в каждом событии, а «пусто» по контракту —
пустое значение, а не отсутствие ключа. Этим же закрыт критерий #43 про опечатку
в имени.
Тип проверяется не у всех колонок, а у пяти: WatchID, VisitID, ClientID,
EventDate, UTCEventTime разбираются в Nullable и дают NULL, если значение
не той природы. Остальные сорок две достаются обычными типами. Соотношение
цены и пользы: единственный производитель топика — собственный генератор,
сериализующий из контракта по объявленным типам, поэтому неверный тип может
прийти только из руки, а сорок семь проверок на NULL превратили бы матвью в
простыню. Пять выбраны по последствию: это идентификаторы события, визита и
посетителя, дата партиции и метка времени, по которой события упорядочиваются
внутри сессии, — порча любой отравляет всё ниже по течению. CounterID
формально тоже входит в ключ сортировки, но на стенде он константа, и NULL там
взяться неоткуда.
Присутствие иначе и не проверить. Nullable
в ClickHouse не оборачивает составные типы: Nullable(Array) запрещён, а
Array(Nullable(T)) при пропавшем ключе даёт пустой массив, неотличимый от
пустого по смыслу. Таких колонок в контракте двенадцать из сорока семи.
Типизированная Kafka-таблица не используется. Режим kafka_handle_error_mode = 'stream' не используется тоже: при чтении байтами на входе нечему ломаться, и
ошибке разбора взяться неоткуда.
Отсюда ограничение на форму выражений разбора: они собираются только из функций,
которые не бросают исключений. JSONExtract и родственные возвращают значение по
умолчанию или NULL, но не падают, — и правило репозитория «грязные записи не
валят пайплайн» держится теперь именно на этом. Исключение в матвью не ошибка
формата, его не перехватит никакой режим Kafka-движка: вставка упадёт, офсеты не
закоммитятся, блок пойдёт читаться снова. Ломается это громко и чинится без
потерь — поправил матвью, потребление продолжилось с некоммиченного офсета, — но
пока не починено, топик стоит. Поэтому приведение Nullable к необнуляемому типу
и любая арифметика в этих выражениях живут за предикатом, который NULL уже отсёк.
Этим решение снимает ограничение, записанное в постановке #37: «ошибки разбора должны рождаться на шаге Kafka-движка». Оно ставилось как условие достижимости критерия #43 про громкую ошибку на опечатку в имени поля — критерий достижим и без него, средствами матвью.
Требование строгого приёма из раздела 6 спеки остаётся в силе, меняется его
механизм: настройка формата input_format_skip_unknown_fields при чтении
байтами беспредметна, раскладки по полям на этом шаге нет вовсе.
Почему
Типизированный чтец не отдаёт сырьё. При kafka_handle_error_mode = 'stream'
виртуальные колонки _raw_message и _error заполняются только при ошибке
разбора и пусты у разобранных сообщений. Один типизированный чтец оставил бы в
STG сырьё исключительно от брака: слой сырых строк хранил бы всё, кроме того,
что доехало.
Слои идут цепочкой. Kafka → STG → ODS → DDS → DM, каждый читает предыдущий. Отвергнутая схема с двумя чтецами одного топика — сырым для STG и типизированным для ODS — в бою обычна и дала бы строгий приём даром, средствами самого движка. Но ODS перестал бы быть надстройкой над STG и стал бы вторым входом с шины, а на стенде, где слои и есть предмет изучения, это дороже сэкономленного. Побочная выгода скромнее, чем кажется на первый взгляд: слой сырья при двух чтецах существовал бы точно так же, но SQL разбора не существовал бы нигде, и для переобработки дня X его пришлось бы писать заново. В цепочке он уже есть в матвью, и ручная вставка получается его копией.
Слой сырья ничего не проверяет, и формат чтеца выбран под это. Соседний вариант,
JSONAsString, требует, чтобы каждое сообщение было корректным объектом JSON, и
на некорректном падает: потребление встаёт, хотя правило репозитория гласит, что
грязные записи не должны валить пайплайн. Лечится это двумя способами — включить
на чтеце режим обработки ошибок или не проверять на входе вовсе. Второе честнее:
частичная проверка в слое, чья работа — не проверять, спорит сама с собой, а
настоящий разбор всё равно идёт ниже. При RawBLOB ломаться нечему, доезжают
любые байты, и оба класса брака разбираются в одном месте.
Архив нужен буквальный и читаемый глазами. RawBLOB — это про способ чтения, а
не про тип колонки: на диске лежит обычный String с текстом события, и менти
открывает его в обычном клиенте и читает. В этом и смысл слоя — увидеть, что
реально пришло. Мусор при этом лежит в той же колонке, что и целые события, а не
в отдельной. Отвергнутый боевой вариант — собрать слой сырья на движке Null,
тогда данные не хранятся вовсе, а матвью работает чистым триггером. Отлаживать
так неудобно, поэтому здесь окно в трое суток, а не ноль.
Что при этом потеряно и чем закрыто. Настройка input_format_skip_unknown_fields = 0 ловила два случая: ожидаемое поле пропало и появилось лишнее, неизвестное.
Оба берёт на себя сверка набора ключей: опечатка в имени — это одновременно
пропавшее ожидаемое и появившееся лишнее, и сравнение множеств видит и то, и
другое. Она же единственная защита у колонок-массивов и она же ловит молчаливое
расширение контракта, когда в сообщении появляется поле, о котором хранилище не
знает. Nullable-разбор к этой работе отношения не имеет — он про тип, а не про
присутствие, и потому сужен до пяти колонок.
Цена решения. Контракт получает ещё два места: типы сорока семи колонок в
выражениях матвью и список тех же имён для сверки ключей. Раздел 1.4 спеки
предупреждает, что, повторяясь примерно в семи местах, они расходятся молча, а
contract-тест из #43 сюда не дотягивается — он сравнивает system.columns
целевой таблицы со схемой генератора и о выражениях матвью ничего не знает.
Смягчение работает не везде: опечатка в имени скалярного поля уводит строки в
таблицу ошибок пачкой и видна сразу, а опечатка в имени массива даёт пустой
массив тихо — сверка ключей проверяет ключи сообщения, а не выражения матвью.
Эти двенадцать колонок сторожит smoke: известное событие с товарами обязано
доезжать с непустыми массивами.
Второе — разбор функциями дороже разбора форматом. На объёмах стенда это
несущественно; если станет заметно, сорок семь вызовов сворачиваются в один
JSONExtract в именованный кортеж, и строка разбирается однократно.
Что проверено
По документации ClickHouse через MCP Context7: основная сверка — 3 августа
2026 года, перепроверка после правок — 5 августа. Датировка важна: 5 августа
утверждение про Nullable(Tuple) развернулось на противоположное.
При режиме stream движок отдаёт _raw_message и _error только для
сообщений, которые не разобрались, и оставляет их пустыми для разобранных.
Отсюда весь довод о том, что типизированный чтец не может наполнить слой сырья.
Составные типы Nullable не оборачивает: Nullable(Array) и Nullable(Map) не
поддерживаются, а Nullable внутри них — да. Отсюда слепота Nullable-разбора
на двенадцати колонках и необходимость сверки ключей. Оговорка про кортеж:
Nullable(Tuple) в ClickHouse всё-таки есть, но за настройкой
enable_nullable_tuple_type и в статусе беты; на довод это не влияет — кортежей
в контракте нет. Nullable внутри кортежа разрешён без всяких флагов, поэтому
запасной вариант со свёрткой в именованный кортеж жив.
Оттуда же: у Kafka-движка есть третий режим обработки ошибок —
dead_letter_queue с записью в системную таблицу. Он отвергнут независимо от
остального: спека требует свои *_errors, системная таблица их не заменяет.
Форматы RawBLOB и LineAsString в ClickHouse есть, и Kafka-движок
поддерживает все форматы. RawBLOB читает вход в одно значение и рассчитан на
таблицу с единственным полем String — отсюда ограничение на форму чтеца.
Матвью с источником-Distributed срабатывает на вставку именно в эту
распределённую таблицу — блок она видит до разрезания по шардам. Проверено
владельцем на рабочих проектах; на стенде подтверждается заодно с приёмкой #37.
На живом стенде проверяется при исполнении #37. Первые два внесены в раздел 11
спеки; остальные — однострочные SELECT, их довольно прогнать заодно:
RawBLOBв Kafka-движке даёт ровно одну строку на сообщение. Проверять это нужно первым и до написания DDL: формулировка «читает вход в одно значение» про файл понятна, а про пачку сообщений из топика — нет, и если сообщения склеятся, переделывать придётся решение целиком, а не DDL. Опыт стоит трёх сообщений и одногоcount(). Запасной вариант —LineAsString: он режет по переводу строки, а события у нас однострочные; цена запасного — сообщение с переводом строки внутри даст две строки вместо одной;- форма именованного кортежа в
JSONExtractсNullable-членами — нужна для свёртки сорока семи вызовов в один, если разбор окажется дорогим; isValidJSON('123')возвращает единицу, аJSONExtractKeysот скаляра — пустой массив. На обоих стоят классы брака и их приоритет, а документация поведение на не-объекте не описывает: дваSELECTзакрывают вопрос;JSONAsStringдействительно падает на некорректном JSON, а не пропускает строку. На этом стоит отказ от него в пользуRawBLOB; документация про ошибочный ввод молчит.