Топик orders, байтовый чтец и даг приёма: забор в STG #93
Notifications
Due Date
No due date set.
Blocks
Depends on
#94 Версионный ODS заказов: строгий переход, брак и ods.order_v
ddmitry/clickstream-data-platform
#92 Команда snapshot: окно, граница суток, байты и опись
ddmitry/clickstream-data-platform
Reference: ddmitry/clickstream-data-platform#93
Reference in New Issue
Block a user
Part of #5.
Цель
Транспорт второго источника: топик
orders, байтовый чтец без матвью иотдельный даг приёма заказов, который одной порцией кладёт слепок в
stg.orders_raw. Работники пульта играют день и триггерят приём; тот же дагразово дёргает заливка стартового мира (#96). Учебный результат: пакетный
забор, управляемый Airflow, — контраст потоковому приёму
hits.Что войдёт
kafka-init: топикorders, одна партиция (ADR 0008).RawBLOBна ноде 1, безON CLUSTERи без матвью;stg.orders_raw— сырой текст, метаданные доставки из виртуальныхколонок,
_load_id,_load_ts. Пары таблиц, сроки хранения и типыслужебных колонок — по конвенциям
docs/architecture/storage.md;расхождение с конвенцией — вопрос владельцу.
2026 года): путь Kafka → STG → ODS живёт одним дагом. Этот тикет даёт ему
шаг забора: один прямой
SELECTчтеца →INSERTвstg.orders_raw;_load_id=run_idэтого дага,_load_tsвычисляется при записи.Настройки прямого чтения —
stream_like_engine_allow_direct_select = 1иkafka_commit_on_select = 1; вторая неочевидна: без неё каждый прогонзабирает один и тот же слепок заново (ADR 0008).
world_next_day,world_live_day) прирастают: послепроигрыша дня — вызов команды
snapshotсыгранного диапазона(
DockerOperator, какbatch; какой слепок уезжает — правило сдвига впроигрывателе, #92), затем триггер дага приёма с ожиданием завершения;
упавший приём красит работника (правило конца графа — ADR 0003).
допущение «один запуск — одно чтение» покрывает разгоны порядка недели
(~1,5 тыс. строк на слепок при пределе порции в десятки тысяч). Короткая
порция — задержка, не потеря.
доставки при прямом чтении — опыт здесь, результат — в ADR тем же PR.
Границы
ods.order_v— #94.строк, чтение до пустоты, фиксация офсетов — отклонены (
ingestion.md).из #92.
Кому что
Форму дага и DDL с комментариями пишет Opus — учебный код, менти его читает. Кодексу — обвязка и опыты: настройки чтеца, эксперимент с виртуальными колонками, прогоны работников. Бриф Кодексу обязан требовать: проверять свойство, а не текст вывода грепом; только POSIX-инструменты (ripgrep на машине стенда нет); вердикт «недостижимо» перепроверять на достижимость, а не на факт.
Сначала прочитать
docs/architecture/orders/ingestion.md— «Поток данных», «Одно чтениеKafka», «Что проверено»;
docs/adr/0008-order-ingestion.md— настройки чтеца и «осталосьпроверить»;
docs/adr/0009-world-control.mdи даги пульта — как триггерят и ждут;docs/adr/0003-dag-run-state.md— конец графа не зеленеет при отказевыше;
docs/architecture/storage.md— служебные колонки, конвенции, раздел прозаказы.
Проверка
make lint,make config-test; на стендеmake smoke,make check-clickhouseзелёные. Ожидание доезда — опросом с таймаутом, неsleep(образцы —scripts/stand-services.sh); инструменты POSIX.world_next_dayдня D кладёт слепок дня D−1 в топик изабирает его в STG: счёт строк с
_load_idпрогона равен числуотправленных (печать из #92). Разовая приёмка допущения «один запуск —
одно чтение», не постоянный сторож.
_load_idи слепки двух соседнихдней.
stg.orders_rawзаполнены; результат опыта свиртуальными колонками — в ADR 0008.
же день.