Брак в слепке: куда уходит и по каким классам #80

Closed
opened 2026-08-12 21:52:36 +03:00 by ddmitry · 1 comment
Owner

Part of #69.

Вопрос

Куда уходит непрошедшая строка слепка и по каким классам она делится?

У событий это решено: две матвью делят поток без зазора и нахлёста, три класса
брака с объявленным приоритетом, сырой текст рядом с классом
(ADR 0005).
Раздел 7 мастер-спеки обещает заказам такую же пару — ods.order_snapshot
(+_errors), — но механики за обещанием нет, а приём у заказов теперь другого
рода: не две матвью, а пакетный шаг
(ADR 0008).
Значит форму «годное сюда, брак туда» надо выбирать заново, а не копировать.

Что решать:

  • классы брака слепка и нужен ли им приоритет — у события их три, и порядок
    обязателен, потому что классы пересекаются;
  • где живёт разделение при пакетном шаге: два запроса, один с INSERT SELECT
    в каждую цель, или иная форма;
  • что считать опорными полями слепка по образцу пяти опорных колонок события;
  • держится ли правило «грязные записи не валят пайплайн», когда разбор идёт
    запросом, а не матвью: у матвью исключение останавливает потребление, у
    пакетного шага — роняет задачу дага, и это другая цена.

Правило репозитория «ошибки разбора уходят в *_errors» — вход в разговор, а
не ответ на него.

Part of #69. ## Вопрос Куда уходит непрошедшая строка слепка и по каким классам она делится? У событий это решено: две матвью делят поток без зазора и нахлёста, три класса брака с объявленным приоритетом, сырой текст рядом с классом ([ADR 0005](https://git.dementev.space/ddmitry/clickstream-data-platform/src/branch/main/docs/adr/0005-event-ingestion.md)). Раздел 7 мастер-спеки обещает заказам такую же пару — `ods.order_snapshot` (+`_errors`), — но механики за обещанием нет, а приём у заказов теперь другого рода: не две матвью, а пакетный шаг ([ADR 0008](https://git.dementev.space/ddmitry/clickstream-data-platform/src/branch/main/docs/adr/0008-order-ingestion.md)). Значит форму «годное сюда, брак туда» надо выбирать заново, а не копировать. Что решать: - классы брака слепка и нужен ли им приоритет — у события их три, и порядок обязателен, потому что классы пересекаются; - где живёт разделение при пакетном шаге: два запроса, один с `INSERT SELECT` в каждую цель, или иная форма; - что считать опорными полями слепка по образцу пяти опорных колонок события; - держится ли правило «грязные записи не валят пайплайн», когда разбор идёт запросом, а не матвью: у матвью исключение останавливает потребление, у пакетного шага — роняет задачу дага, и это другая цена. Правило репозитория «ошибки разбора уходят в `*_errors`» — вход в разговор, а не ответ на него.
ddmitry added the wayfinder:grilling label 2026-08-12 21:52:36 +03:00
ddmitry added a new dependency 2026-08-12 22:51:02 +03:00
ddmitry self-assigned this 2026-08-16 18:03:19 +03:00
Author
Owner

Решение

Учебный результат: менти различает бизнес-версию сущности (updated_at),
происхождение строки (snapshot_date) и запуск загрузки (_load_id), а также
видит цену отложенного схлопывания версий — FINAL на границе точного чтения.

Граница брака

Единица приёма — одна строка stg.orders_raw. Годные строки продолжают путь в
ODS, негодные независимо уходят в ods.order_snapshot_errors; одна грязная
строка не валит задачу Airflow.

Строгий приём проверяет только контракт провода:

  • корень — JSON-объект с точным набором 11 верхнеуровневых ключей;
  • целые, строки денег с двумя знаками, времена UTC до миллисекунд и дата имеют
    каноническую форму, принятую в
    «Форме записи слепка на проводе»;
  • items на этой границе проверяется только как JSON-массив.

Внутренние поля items, допустимые переходы статуса, равенства сумм и прочие
бизнес-инварианты сюда не входят. Это уже проверка содержания, а не способность
безопасно разобрать строку.

Классы взаимоисключающи за счёт приоритета:

  1. not_an_object;
  2. keyset_mismatch;
  3. field_invalid.

Отдельной детализации до имени поля нет: сырой текст остаётся рядом и служит
разбору единичного случая.

Форма пакетного шага

Сначала прямое чтение Kafka один раз записывает порцию в stg.orders_raw с
одним _load_id. Затем два INSERT SELECT читают один и тот же срез STG:
предикат брака пишет в _errors, его буквальное отрицание — в типизированный
ODS. Постоянная таблица или представление классификатора не нужны.

_load_idrun_id Airflow и граница одного запуска приёма, а не номер
слепка. Он хранится в STG и обеих целях ODS. _load_ts остаётся временем
прибытия строки и не участвует в выборе актуальной бизнес-версии.

ODS как поток версий

Горячее ревью отменило прежнюю подмену партиций из
«Приёма заказов: нужен ли слепку слой сырья».
Одно прямое чтение Kafka отдаёт одну произвольную порцию, а протокол не несёт
признака конца слепка. Делать snapshot_date границей целостности означало бы
обучать свойству стенда, которого нет у источника.

ods.order_snapshot принимает версии заказа на ReplacingMergeTree(updated_at)
с бизнес-ключом order_id; партиционирование строится от неизменного
created_at, чтобы поздние версии одного заказа попадали в одну партицию.
snapshot_date сохраняется как дата наблюдения, но не управляет публикацией и
идемпотентностью. Физическое удаление отсутствием строки не моделируется:
жизненный цикл заказа выражен статусом.

Обычный SELECT показывает физически поступившие версии и годится для
диагностики. Потребитель, которому нужно точное актуальное состояние, читает с
FINAL. Нужен ли при сборке DDS именно FINAL, argMax или другой
инкрементальный приём, этот тикет не решает.

Что отвергнуто

  • REPLACE PARTITION snapshot_date — требует отсутствующего признака полноты;
    чтение до временной пустоты было бы искусственным контрактом стенда.
  • ReplacingMergeTree по (snapshot_date, order_id) — оставляет синтетическую
    дату техническим ключом и превращает ODS в архив копий выгрузки.
  • Обычный MergeTree с дедупликацией только в DDS — сохраняет повторы доставки
    навсегда и откладывает уже известную семантику ODS.
  • Маркер конца, манифест или ожидаемый счётчик — новый протокол без нужного
    наблюдаемого эффекта.
  • argMax по всем полям прямо в контракте ODS — повторяет движок многословным
    запросом; выбор чтения DDS оставлен проектированию DDS.
  • Полная проверка бизнес-контракта и вложенных позиций — оборонительный код, за
    который менти платит вниманием, не получая урока приёма.

Проверка спорного поведения ClickHouse

MCP Context7 в текущей сессии был недоступен. Поведение проверено на локальном
ClickHouse 26.3.17.56 и по
официальной документации ReplacingMergeTree.
JSONExtract* допускает преобразования типов и значения по умолчанию, поэтому
строгий предикат сначала проверяет JSON-тип и каноническую строковую форму.
ReplacingMergeTree схлопывает версии в фоновых слияниях, поэтому точное чтение
требует FINAL. Исходники той же версии —
StorageKafka.cpp
и
KafkaSource.cpp
— подтвердили: прямой SELECT Kafka Engine заканчивается после одной
полученной порции, а не осушает топик до конца.

Отдельно вынесен вопрос
«Трассировка загрузки: нужен ли _load_id выше ODS»:
его следует решать вместе с устройством DDS, а не здесь.

## Решение Учебный результат: менти различает бизнес-версию сущности (`updated_at`), происхождение строки (`snapshot_date`) и запуск загрузки (`_load_id`), а также видит цену отложенного схлопывания версий — `FINAL` на границе точного чтения. ### Граница брака Единица приёма — одна строка `stg.orders_raw`. Годные строки продолжают путь в ODS, негодные независимо уходят в `ods.order_snapshot_errors`; одна грязная строка не валит задачу Airflow. Строгий приём проверяет только контракт провода: - корень — JSON-объект с точным набором 11 верхнеуровневых ключей; - целые, строки денег с двумя знаками, времена UTC до миллисекунд и дата имеют каноническую форму, принятую в [«Форме записи слепка на проводе»](https://git.dementev.space/ddmitry/clickstream-data-platform/issues/81); - `items` на этой границе проверяется только как JSON-массив. Внутренние поля `items`, допустимые переходы статуса, равенства сумм и прочие бизнес-инварианты сюда не входят. Это уже проверка содержания, а не способность безопасно разобрать строку. Классы взаимоисключающи за счёт приоритета: 1. `not_an_object`; 2. `keyset_mismatch`; 3. `field_invalid`. Отдельной детализации до имени поля нет: сырой текст остаётся рядом и служит разбору единичного случая. ### Форма пакетного шага Сначала прямое чтение Kafka один раз записывает порцию в `stg.orders_raw` с одним `_load_id`. Затем два `INSERT SELECT` читают один и тот же срез STG: предикат брака пишет в `_errors`, его буквальное отрицание — в типизированный ODS. Постоянная таблица или представление классификатора не нужны. `_load_id` — `run_id` Airflow и граница одного запуска приёма, а не номер слепка. Он хранится в STG и обеих целях ODS. `_load_ts` остаётся временем прибытия строки и не участвует в выборе актуальной бизнес-версии. ### ODS как поток версий Горячее ревью отменило прежнюю подмену партиций из [«Приёма заказов: нужен ли слепку слой сырья»](https://git.dementev.space/ddmitry/clickstream-data-platform/issues/70). Одно прямое чтение Kafka отдаёт одну произвольную порцию, а протокол не несёт признака конца слепка. Делать `snapshot_date` границей целостности означало бы обучать свойству стенда, которого нет у источника. `ods.order_snapshot` принимает версии заказа на `ReplacingMergeTree(updated_at)` с бизнес-ключом `order_id`; партиционирование строится от неизменного `created_at`, чтобы поздние версии одного заказа попадали в одну партицию. `snapshot_date` сохраняется как дата наблюдения, но не управляет публикацией и идемпотентностью. Физическое удаление отсутствием строки не моделируется: жизненный цикл заказа выражен статусом. Обычный `SELECT` показывает физически поступившие версии и годится для диагностики. Потребитель, которому нужно точное актуальное состояние, читает с `FINAL`. Нужен ли при сборке DDS именно `FINAL`, `argMax` или другой инкрементальный приём, этот тикет не решает. ### Что отвергнуто - `REPLACE PARTITION snapshot_date` — требует отсутствующего признака полноты; чтение до временной пустоты было бы искусственным контрактом стенда. - `ReplacingMergeTree` по `(snapshot_date, order_id)` — оставляет синтетическую дату техническим ключом и превращает ODS в архив копий выгрузки. - Обычный `MergeTree` с дедупликацией только в DDS — сохраняет повторы доставки навсегда и откладывает уже известную семантику ODS. - Маркер конца, манифест или ожидаемый счётчик — новый протокол без нужного наблюдаемого эффекта. - `argMax` по всем полям прямо в контракте ODS — повторяет движок многословным запросом; выбор чтения DDS оставлен проектированию DDS. - Полная проверка бизнес-контракта и вложенных позиций — оборонительный код, за который менти платит вниманием, не получая урока приёма. ### Проверка спорного поведения ClickHouse MCP Context7 в текущей сессии был недоступен. Поведение проверено на локальном ClickHouse `26.3.17.56` и по [официальной документации `ReplacingMergeTree`](https://clickhouse.com/docs/reference/engines/table-engines/mergetree-family/replacingmergetree). `JSONExtract*` допускает преобразования типов и значения по умолчанию, поэтому строгий предикат сначала проверяет JSON-тип и каноническую строковую форму. `ReplacingMergeTree` схлопывает версии в фоновых слияниях, поэтому точное чтение требует `FINAL`. Исходники той же версии — [`StorageKafka.cpp`](https://github.com/ClickHouse/ClickHouse/blob/v26.3.17.56-lts/src/Storages/Kafka/StorageKafka.cpp) и [`KafkaSource.cpp`](https://github.com/ClickHouse/ClickHouse/blob/v26.3.17.56-lts/src/Storages/Kafka/KafkaSource.cpp) — подтвердили: прямой `SELECT` Kafka Engine заканчивается после одной полученной порции, а не осушает топик до конца. Отдельно вынесен вопрос [«Трассировка загрузки: нужен ли `_load_id` выше ODS»](https://git.dementev.space/ddmitry/clickstream-data-platform/issues/85): его следует решать вместе с устройством DDS, а не здесь.
Sign in to join this conversation.