Брак в слепке: куда уходит и по каким классам #80
Notifications
Due Date
No due date set.
Blocks
Depends on
#84 Вычитание: какая механика этапа 3 не окупает себя эффектом
ddmitry/clickstream-data-platform
#85 Трассировка загрузки: нужен ли _load_id выше ODS
ddmitry/clickstream-data-platform
#71 Генератор слепков: где живёт и чем связан с событиями
ddmitry/clickstream-data-platform
#81 Форма записи слепка на проводе
ddmitry/clickstream-data-platform
Reference: ddmitry/clickstream-data-platform#80
Reference in New Issue
Block a user
Part of #69.
Вопрос
Куда уходит непрошедшая строка слепка и по каким классам она делится?
У событий это решено: две матвью делят поток без зазора и нахлёста, три класса
брака с объявленным приоритетом, сырой текст рядом с классом
(ADR 0005).
Раздел 7 мастер-спеки обещает заказам такую же пару —
ods.order_snapshot(+
_errors), — но механики за обещанием нет, а приём у заказов теперь другогорода: не две матвью, а пакетный шаг
(ADR 0008).
Значит форму «годное сюда, брак туда» надо выбирать заново, а не копировать.
Что решать:
обязателен, потому что классы пересекаются;
INSERT SELECTв каждую цель, или иная форма;
запросом, а не матвью: у матвью исключение останавливает потребление, у
пакетного шага — роняет задачу дага, и это другая цена.
Правило репозитория «ошибки разбора уходят в
*_errors» — вход в разговор, ане ответ на него.
Решение
Учебный результат: менти различает бизнес-версию сущности (
updated_at),происхождение строки (
snapshot_date) и запуск загрузки (_load_id), а такжевидит цену отложенного схлопывания версий —
FINALна границе точного чтения.Граница брака
Единица приёма — одна строка
stg.orders_raw. Годные строки продолжают путь вODS, негодные независимо уходят в
ods.order_snapshot_errors; одна грязнаястрока не валит задачу Airflow.
Строгий приём проверяет только контракт провода:
каноническую форму, принятую в
«Форме записи слепка на проводе»;
itemsна этой границе проверяется только как JSON-массив.Внутренние поля
items, допустимые переходы статуса, равенства сумм и прочиебизнес-инварианты сюда не входят. Это уже проверка содержания, а не способность
безопасно разобрать строку.
Классы взаимоисключающи за счёт приоритета:
not_an_object;keyset_mismatch;field_invalid.Отдельной детализации до имени поля нет: сырой текст остаётся рядом и служит
разбору единичного случая.
Форма пакетного шага
Сначала прямое чтение Kafka один раз записывает порцию в
stg.orders_rawсодним
_load_id. Затем дваINSERT SELECTчитают один и тот же срез STG:предикат брака пишет в
_errors, его буквальное отрицание — в типизированныйODS. Постоянная таблица или представление классификатора не нужны.
_load_id—run_idAirflow и граница одного запуска приёма, а не номерслепка. Он хранится в 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— подтвердили: прямой
SELECTKafka Engine заканчивается после однойполученной порции, а не осушает топик до конца.
Отдельно вынесен вопрос
«Трассировка загрузки: нужен ли
_load_idвыше ODS»:его следует решать вместе с устройством DDS, а не здесь.