Топик orders, байтовый чтец и даг приёма: забор в STG #93

Closed
opened 2026-08-17 22:54:56 +03:00 by ddmitry · 0 comments
Owner

Part of #5.

Цель

Транспорт второго источника: топик orders, байтовый чтец без матвью и
отдельный даг приёма заказов, который одной порцией кладёт слепок в
stg.orders_raw. Работники пульта играют день и триггерят приём; тот же даг
разово дёргает заливка стартового мира (#96). Учебный результат: пакетный
забор, управляемый Airflow, — контраст потоковому приёму hits.

Что войдёт

  • kafka-init: топик orders, одна партиция (ADR 0008).
  • DDL: Kafka-чтец RawBLOB на ноде 1, без ON CLUSTER и без матвью;
    stg.orders_raw — сырой текст, метаданные доставки из виртуальных
    колонок, _load_id, _load_ts. Пары таблиц, сроки хранения и типы
    служебных колонок — по конвенциям docs/architecture/storage.md;
    расхождение с конвенцией — вопрос владельцу.
  • Отдельный даг приёма заказов (решение владельца при нарезке, 17 августа
    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).
  • Разгон «дней = N» отправляет N слепков; забор остаётся одной порцией —
    допущение «один запуск — одно чтение» покрывает разгоны порядка недели
    (~1,5 тыс. строк на слепок при пределе порции в десятки тысяч). Короткая
    порция — задержка, не потеря.
  • Из «осталось проверить» ADR 0008: доступны ли виртуальные колонки
    доставки при прямом чтении — опыт здесь, результат — в ADR тем же PR.

Границы

  • Переход STG → ODS, классы брака, ods.order_v#94.
  • Сенсоров и triggerer нет (ADR 0009); маркер конца слепка, опись ожидаемых
    строк, чтение до пустоты, фиксация офсетов — отклонены (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.
  • Обрыв работника не двигает позицию мира; повторный прогон играет тот
    же день.
Part of #5. ## Цель Транспорт второго источника: топик `orders`, байтовый чтец без матвью и отдельный даг приёма заказов, который одной порцией кладёт слепок в `stg.orders_raw`. Работники пульта играют день и триггерят приём; тот же даг разово дёргает заливка стартового мира (#96). Учебный результат: пакетный забор, управляемый Airflow, — контраст потоковому приёму `hits`. ## Что войдёт - `kafka-init`: топик `orders`, одна партиция (ADR 0008). - DDL: Kafka-чтец `RawBLOB` на ноде 1, без `ON CLUSTER` и без матвью; `stg.orders_raw` — сырой текст, метаданные доставки из виртуальных колонок, `_load_id`, `_load_ts`. Пары таблиц, сроки хранения и типы служебных колонок — по конвенциям `docs/architecture/storage.md`; расхождение с конвенцией — вопрос владельцу. - Отдельный даг приёма заказов (решение владельца при нарезке, 17 августа 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). - Разгон «дней = N» отправляет N слепков; забор остаётся одной порцией — допущение «один запуск — одно чтение» покрывает разгоны порядка недели (~1,5 тыс. строк на слепок при пределе порции в десятки тысяч). Короткая порция — задержка, не потеря. - Из «осталось проверить» ADR 0008: доступны ли виртуальные колонки доставки при прямом чтении — опыт здесь, результат — в ADR тем же PR. ## Границы - Переход STG → ODS, классы брака, `ods.order_v` — #94. - Сенсоров и triggerer нет (ADR 0009); маркер конца слепка, опись ожидаемых строк, чтение до пустоты, фиксация офсетов — отклонены (`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. - [x] Прогон `world_next_day` дня D кладёт слепок дня D−1 в топик и забирает его в STG: счёт строк с `_load_id` прогона равен числу отправленных (печать из #92). Разовая приёмка допущения «один запуск — одно чтение», не постоянный сторож. - [x] Два прогона подряд дают два разных `_load_id` и слепки двух соседних дней. - [x] Метаданные доставки в `stg.orders_raw` заполнены; результат опыта с виртуальными колонками — в ADR 0008. - [x] Обрыв работника не двигает позицию мира; повторный прогон играет тот же день.
ddmitry added the ready-for-agent label 2026-08-17 22:56:11 +03:00
Sign in to join this conversation.
1 Participants
Notifications
Due Date
No due date set.
Reference: ddmitry/clickstream-data-platform#93