Files
clickstream-data-platform/docs/adr/0005-event-ingestion.md
T
ddadmin e6f300a66e docs(storage): правки конвенции по двум холодным ревью
- Зачем:
  - раздел «Часовые пояса» прошёл два холодных ревью — по дефектам и по
    уместности. Первое поймало ложный замер и три расхождения с живым
    стендом, второе — материал не своей зоны и дубли (#63).
- Что:
  - замер «расхождение живёт по HTTP» отозван: мерил toString(UTCEventTime)
    в родном клиенте против голой колонки по HTTP, а это разные вещи.
    Перемерено — клиенты ведут себя одинаково; записан верный факт: вывод
    колонки идёт по поясу сессии, функция — по поясу типа.
  - «по поясу сервера» заменено на «по поясу сессии, а тот по умолчанию
    серверный» — в разделе и в ADR 0005; утверждение в ледгере переписано
    под измеренный раскол вывода и типа.
  - PARTITION BY toDate(_load_ts) больше не выдаётся за уже соблюдённое
    правило: _load_ts сегодня DateTime64(3) без пояса.
  - абзац ADR 0005 больше не спорит с цитатой вызова строкой выше.
  - вырезано: веер отклонённых вариантов под заголовком (живые отказы
    разложены прозой по своим абзацам, как принято в этом документе),
    ссылка на несуществующую связку в world.py, осиротевшая строка про
    Grafana, абзац про пояс показа — он уехал комментарием в #63.
- Проверка:
  - make lint
  - замеры повторены на живом стенде 8 августа 2026 года
  - DDL к конвенции по-прежнему не приведён: документы описывают цель
2026-08-08 19:05:39 +03:00

24 KiB
Raw Blame History

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 спеки предупреждает, что, повторяясь примерно в семи местах, они расходятся молча.

И сверка ключей эту цену покрывает не всю: она смотрит на ключи сообщения, а не на выражения матвью. Опечатка в имени внутри JSONExtract даёт умолчание типа — ноль, пустую строку, пустой массив, — и молчит она у сорока двух обычных колонок ровно так же, как у двенадцати массивов. Громко ломаются только пять опорных: у них разбор Nullable стоит в предикате, и опечатка уводит в таблицу ошибок все строки до единой. Остальные сорок две сторожит сверка разобранного события против сырого текста — разовый опыт при исполнении #43, а не постоянная проверка; сила его в том, что выражения сверки собираются из контракта, а выражения матвью написаны руками по описанию выгрузки, и одна опечатка в двух местах не повторяется.

Второе — разбор функциями дороже разбора форматом. На объёмах стенда это несущественно; если станет заметно, сорок семь вызовов сворачиваются в один 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 срабатывает на вставку именно в эту распределённую таблицу — блок она видит до разрезания по шардам. Проверено владельцем на рабочих проектах; на стенде проверяется вместе с матвью разбора, то есть при исполнении #43.

Пункт про RawBLOB, стоявший в списке ниже первым, закрыт при исполнении #37, на стенде 5 августа 2026 года: формат даёт ровно одну строку на каждое непустое сообщение, пачка продюсера границы сообщений не стирает, запасной LineAsString не понадобился. Тогда же нашлась и граница обещания — запись с пустым значением и запись-надгробие не дают строки вовсе; из-за неё в «Решении» и в «Почему» приписано слово «непустое». Замер целиком — в доке хранилища, раздел «Что проверено».

Остальные четыре закрыты при исполнении #43, на стенде 7 августа 2026 года, ClickHouse 26.3.17.56. Все четыре ответили так, как ждала постановка:

  • isValidJSON('123') возвращает единицу — скаляр законный JSON, и первый класс брака поэтому проверяет именно объект, а не валидность;
  • JSONExtractKeys('123') возвращает пустой массив, и он же приходит от вовсе не-JSON. Значит скаляр проваливает и сверку ключей — отсюда обязательный порядок классов. Раз isValidJSON объект от скаляра не отличает, первый класс держит JSONType: она возвращает Object у объекта, Int64 у скаляра 123 и Null у вовсе не-JSON и у пустой строки — исключения не бросает ни в одном случае, то есть годится в предикат;
  • JSONAsString на некорректном вводе падает, а не пропускает строку: код 117 INCORRECT_DATA, «JSON object must begin with '{'». Падает и на скаляре 123. На этом стоит отказ от него в пользу RawBLOB;
  • форма именованного кортежа с Nullable-членами работает и годится в запасной вариант: JSONExtract(raw, 'Tuple(WatchID Nullable(UInt64), …)') даёт NULL в тех членах, что не разобрались, обращение по имени члена доступно, а на не-объекте кортеж выходит целиком из NULL и исключения нет.

Тем же заходом нашлось то, о чём никто не спрашивал, и оно оказалось блокирующим. JSONExtract с типом DateTime не разбирает ISO-8601 с суффиксом зоны. На проводе UTCEventTime уезжает как 2026-06-01T12:34:56Z (спека генератора, раздел 4), а JSONExtract(raw, 'UTCEventTime', 'Nullable(DateTime)') отдаёт на такой строке NULL — то есть все события до единого уходили бы в брак с классом key_field_unparsed. Ни DateTime64, ни DateTime('UTC') суффикс тоже не берут; без Z та же строка разбирается. Спека генератора обещала обратное («принимает ISO без плясок») — обещание было ошибочным и исправлено тем же PR. Настройка cast_string_to_date_time_mode = 'best_effort' дела не меняет: с ней суффикс берёт CAST, а JSONExtract по-прежнему отдаёт NULL — то есть внутренний разбор JSONExtract её не слушает. EventDate уезжает как 2026-06-01 и разбирается JSONExtract без оговорок.

Разбор метки времени идёт parseDateTimeOrNull(JSONExtractString(raw, 'UTCEventTime'), '%Y-%m-%dT%H:%i:%SZ') — по буквально названному формату, а не через parseDateTimeBestEffort. Обе функции ISO-8601 понимают и обе в варианте *OrNull отдают NULL вместо исключения, то есть годятся в предикат. Выбран точный формат потому, что широта здесь работает против строгого приёма: parseDateTimeBestEffort на непонятной строке не краснеет, а достраивает недостающее. Измерено 7 августа 2026 года — обрезанное 20:00:21 он превращает в 2026-01-01 20:00:21, подставив текущий год и первое января; голая дата 2026-05-31 становится полуночью; строка цифр читается числом эпохи. Такое сообщение прошло бы строгий приём с тихо неверным временем — ровно с той порчей, ради которой класс key_field_unparsed и заведён. Разбор по названному формату отдаёт на всех трёх NULL. Источник у топика один и шлёт одну запись, так что широта не нужна вовсе, а платится за неё отключённой проверкой.

Пояс разбору при #63 добавлен третьим аргументом — 'UTC'; вызов выше приведён без него, каким он был до этого решения. Без имени пояса функция трактует показания часов по поясу сессии, а тот по умолчанию серверный. Правило целиком и его довод — конвенция часовых поясов.

Цена выбора измерена на настоящих данных: по всем 101 252 строкам сырья модельного дня (день залит дважды) точный формат разобрал метку у каждой, и ни на одной не разошёлся с parseDateTimeBestEffort. Различаются они только на порче.

Оговорка к обнуляемому разбору даты, измеренная там же: Nullable(Date) даёт NULL на строке, которая датой не является вовсе («мусор»), и на числе, но невозможную дату 2026-13-99 молча приводит к 1970-01-01. То есть класс key_field_unparsed ловит порчу типа, а не порчу значения внутри типа.