Сериализатор, приёмники и проигрыватель с CLI #41

Closed
opened 2026-08-01 21:19:41 +03:00 by ddmitry · 5 comments
Owner

Part of #4.

Цель

Довести события до стенда: единственный канонический сериализатор, глупые
приёмники (файл, Kafka пачкой, Kafka с темпом), проигрыватель и интерфейс
запуска. После этого тикета топик hits наполняется настоящими данными
генератора в обоих режимах.

Тикет идёт до типизированного ODS (#43): порядок развёрнут 6 августа
2026 года. Раньше #43 шёл первым и вынужден был выдумывать событие сам —
своим модулем, печатающим JSON, то есть вторым местом сериализации против
правила «сериализатор один». Теперь отправитель идёт первым, а хранилище
принимает настоящие данные. Разбор наверх (ods.event) появится в #43, и
утверждения про него живут там же.

Что войдёт

  • Канонический сериализатор на orjson: единственное место, где событие
    целиком
    превращается в JSON; канонический порядок ключей и строк; прямых
    json.dumps в коде нет. Граница правила: вложенный блок ecommerce
    собирает день (commerce.py, решение #40), и это не второй сериализатор —
    правило «сериализатор один» про событие, а не про блок внутри него.
  • Форма на проводе — ISO-8601, задана разделом 4 спеки генератора:
    EventDate как 2026-06-01, UTCEventTime как 2026-06-01T12:34:56Z,
    ecommerce — строка с экранированным JSON внутри. Этот тикет реализует
    форму первым; хранилище (#43) потом читает то, что он положил. Довод за
    ISO — читаемость сырья: весь смысл слоя STG в том, что менти открывает
    колонку raw в обычном клиенте и разбирает событие глазами.
  • Приёмники: файл (одно событие — одна строка, файл читается построчно),
    Kafka пачкой (пакетный режим), Kafka с темпом ×60 по умолчанию (живой день;
    ускорение — флаг проигрывателя, лаг логируется). Путь файла, адрес брокера
    и имя топика — параметры запуска, а не константы в коде: генератор задуман
    отдельной сущностью, которую можно унести в другой проект. Выделенного
    места файлу в репозитории не отводится — спека (раздел 5) обещает
    «потерял — пересчитал»; под мусорные прогоны заводится каталог tmp/
    целиком в .gitignore, чтобы файл не рвался в git.
  • Событие несёт все 47 ключей всегда. «Пусто» по контракту — пустое
    значение: пустой массив, пустая строка, ноль, — а не отсутствие ключа.
    Требование ADR 0005: строгий приём хранилища держится на сверке набора
    ключей, и пропавший ключ уводит событие в брак целиком. Раньше это
    обязательство пинил бы #43 своим модулем; после разворота его несёт
    этот тикет.
  • Переигровка дня целиком при обрыве: повтор даёт те же WatchID.
  • Тайминги: генерация и доставка печатаются отдельно.
  • Свой образ генератора. База python:3.14 (генератор пинит
    >=3.14,<3.15, лок — ==3.14.*), зависимости ставятся uv sync --frozen
    из uv.lock: побайтовая воспроизводимость привязана к версиям (раздел 4
    спеки), и второй список пакетов писал бы на стенде не те байты, что
    сторожит make test. Спека сходится сама — раздел 2 обещает детерминизм
    «внутри канонического контейнера», им и становится этот образ. Образ
    Airflow этот тикет не трогает вовсе: ни базу, ни контекст сборки, ни его
    список пакетов для пробников.
    • Каталог товаров должен оказаться в образе в той же взаимной
      раскладке
      , что и пакет: путь считается от модуля вверх по дереву
      (catalog.py:41, parents[3]), а каталог по решению #39 — часть мира,
      а не запуска, и снаружи не настраивается. Копировать в образ или
      монтировать томом — решить при исполнении, но решить явно: иначе
      служба падает на первом же импорте.
    • .dockerignore заводится тем же PR. Сейчас его нет, а в контекст
      сборки иначе уедут generator/.venv (125 МБ) и мусорный tmp/ — кеш
      слоёв будет сбрасываться на каждом прогоне.
  • Разовая служба compose на образе генератора, под профилем (либо через
    docker compose run --rm): в make up до #42 она не входит, обычная
    служба поднялась бы при up сама и стёрла границу с соседом. Оговорка
    для #42: up --wait считает успешную разовую службу упавшей, если у неё
    нет зависимого с service_completed_successfully — предупреждение об этом
    уже стоит в compose.yaml у clickhouse-init.

Решить и внести в спеку генератора (раздел 9) тем же PR

Часть пунктов решена грилингом 7 августа 2026 года — их не переоткрывать, а
внести в спеку как есть; остальные решаются при исполнении.

  • Интерфейс запуска: CLI-команды и/или цели make, как делятся режимы
    проигрывателя. Считать, что зовущий — контейнер: параметры приходят
    аргументами или окружением, итог виден кодом возврата, а логи идут в
    стандартный вывод. Даги world_init/next_day запустят его так на
    этапе 5 — интерфейс не должен им мешать.

  • Размещение решено грилингом 7 августа 2026 года; внести в спеку как
    есть.
    Генератор живёт в своём контейнере, и на этапе 5 даг будет
    запускать контейнер, а не импортировать пакет в процессе воркера. Довод
    учебный: так граница между оркестратором и задачей видна глазами — у
    задачи своё окружение, свой образ, свой процесс, а оркестратор передаёт
    параметры и читает результат. Это та же ступенька, с которой в бою вырастает
    KubernetesPodOperator. Побочно это единственный вариант, при котором
    «генератор — отдельная и заменяемая сущность» перестаёт быть словами: его
    зависимости не смешиваются с окружением Airflow и его провайдеров.
    Хостовый uv остаётся для make test, make lint и make typecheck.

    Отклонено: генератор в образе Airflow, даг импортирует пакет (так у
    стенда v1) — механика проще, но менти видит лишь «даг позвал функцию», а
    платить пришлось бы сменой базового образа всего стенда на python3.14,
    переездом контекста сборки и смешением зависимостей генератора с
    зависимостями Airflow. ExternalPythonOperator — задача в отдельном venv
    внутри образа Airflow
    : изоляцию даёт и сокета не требует, но границу видно
    в конфигурации, а не в docker ps; остаётся запасным путём, если проброс
    сокета однажды окажется неприемлемым.

    Цена названа и принимается: чтобы даг запускал контейнер, в Airflow
    пробрасывается сокет докера, а это доступ, равный root на хосте. На учебном
    стенде размен допустим — но записывается уроком, а не прячется: в бою так
    не делают. Механика этапа 5 (провайдер docker, сеть compose у
    запускаемого контейнера, доступ к группе сокета) проверяется там же, не
    здесь.

  • Клиент Kafka — confluent-kafka (решение того же грилинга). Он уже
    стоит в образе Airflow, на нём написан пробник dags/test_kafka.py, и его
    же использует официальный провайдер Kafka для Airflow. Один клиент на
    стенде вместо двух. Объявляется необязательной группой зависимостей
    генератора
    , а образ ставит её из лока: приёмник импортирует клиента, и
    необъявленный импорт отнял бы у пакета переносимость — довод, на котором
    стоит половина решений этого тикета. Список пакетов в образе Airflow
    (clickhouse-connect, тот же confluent-kafka) правилу «зависимости из
    лока» не противоречит и не трогается: он для пробников dags/, а к
    генератору отношения не имеет.

  • CLI без состояния. Зерно и номер дня приходят параметрами; позицию на
    оси проигрыватель не хранит и не выводит из данных — её ведёт тот, кто
    зовёт. Потребителей интерфейса трое, и третий обычно забывается: даги
    world_init/next_day этапа 5, проверки #43 — и заливка зернового мира
    из #42, восемь дней подряд
    . Восемь запусков службы внутри make up
    плохой ответ; способ проиграть подряд несколько дней одним запуском
    назначается здесь, вместе с остальным интерфейсом. Хвост этапу 5
    (внести в раздел 8 спеки): позиция живёт переменной
    Airflow, отсутствие переменной означает мир в зерновом состоянии, заводит
    и двигает её только даг next_day и только по успеху. Довод — генератор
    отдельная и заменяемая сущность, к хранилищу его привязывать не надо, а
    make clean сносит том метаданных Airflow вместе с данными ClickHouse,
    так что позиция и мир чистятся одной командой. Расхождение переменной с
    данными (менти почистил партицию руками) лечится документацией, а не
    сторожем: постоянная сверка была бы проверкой, красной в норме
    (ADR 0002).

  • Поправить абзац раздела 4 про порядок тикетов. Там записано: «Форму
    пинит хранилище (#43) как первый потребитель… Порядок тикетов обратный
    порядку зависимости, поэтому здесь она и записана». После разворота
    6 августа 2026 года это неверно: форму реализует #41, а хранилище её
    читает. Сама запись формы в спеке остаётся — она и есть источник истины,
    просто довод у неё теперь другой.

  • Ключ сообщения Kafka: назначить или явно отказаться. Довод, стоявший
    здесь раньше, — «у сообщений без ключа работает липкое распределение:
    подряд идущая пачка ложится в один раздел» — оказался свойством клиента, а
    не Kafka: у kafka-python сообщение без ключа уходит в случайный раздел на
    каждое сообщение, липкий партиционер (KIP-480) там опция (сверено через
    Context7 7 августа 2026 года). Для librdkafka, то есть для выбранного
    confluent-kafka, ответ пока неизвестен: в CONFIGURATION.md есть
    sticky.partitioning.linger.ms, но умолчание документация не называет, а
    умолчание топикового партиционера — consistent_random, то есть случайный
    раздел. Сверить до того, как довод ляжет в спеку, и опереть его на
    CONFIGURATION.md, а не на общую справку. Решение (WatchID ключом либо
    «ключа нет, и вот почему») интерфейсное, и место ему здесь.

  • Снять фразу «Сторожится тестом приёмника» в абзаце «Одно событие —
    одно сообщение Kafka» (раздел 4). Постоянного теста приёмника не будет —
    решение владельца 7 августа 2026 года: свойство уже сторожат постоянная
    проверка инварианта слоёв из #43 на настоящих данных и разовая приёмка
    ниже, а дешёвая половина достаётся построчностью файла в тесте
    детерминизма. Второй сторож того же свойства с подставным продюсером
    ничему не учит и живёт зря. Само утверждение в спеке остаётся: это
    контракт транспорта, а не деталь реализации.

  • Способ взять ограниченную пачку — сколько-то событий, а не день
    целиком. Довод от менти: тот, кто хочет посмотреть на конвейер, не должен
    ждать полсотни тысяч событий. На живой день пачка не распространяется —
    там ожидание снимает флаг ускорения, который и так обязан быть проверен.
    Три свойства обязательны, потому что на
    них стоят проверки #43: пачка сочетается с файловым приёмником, даёт
    взять срез, ещё не приезжавший в ODS, и печатает число отправленного
    счёт нужен обеим сторонам. Флаг, срез дня или отдельный режим — решить
    здесь, потому что здесь живёт интерфейс. Записывать в раздел 9 решением,
    а открытый пункт про интерфейс запуска — закрыть: дописывать в открытый
    список нельзя, за #42 остаются два последних пункта, и раздел он закрывает
    целиком.

Критерии приёмки

Приёмка идёт по слою сырья: ods.event и *_errors появятся в #43, и
утверждения про них — его критерии, не эти.

Критерии ниже, помеченные разовыми, требуют прогона целого дня
через стенд и постоянными целями make быть не могут — по тому же правилу,
по которому #43 не ставит день в повторяемую цель. Их след — запись в теле PR.

Убирать за этими прогонами не нужно, и вот почему. Правило
docs/architecture/testing.md требует, чтобы вкладывающая данные проверка
либо прибиралась за собой, либо не жила постоянной целью: вложенное оседает
в мире менти и врёт его лабам. Здесь оно не оседает — на момент этого тикета
типизированного ODS ещё нет (его приносит #43), день доезжает только до
stg.hits_raw_dist, а у сырья срок жизни трое суток, и оно уходит само.
Задним числом матвью не досыпают, так что в ods.event этот день не появится
и потом. Оговорка на будущее: если день гонять уже после того, как #43 слит,
он доедет до ODS и останется там навсегда — тогда действует общее правило,
и день берётся за границей оси мира со сносом партиции.

  • Побайтовый детерминизм: два прогона дня в файл — идентичные байты,
    сравнение хешей (тест). Тем же тестом — построчность: строк в файле
    ровно столько, сколько событий в дне. Это дешёвая половина контракта
    транспорта: склейка событий сломала бы разбор целиком, потому что
    хранилище читает топик байтами.
  • Разовая. Пакетный режим: день доезжает до stg.hits_raw_dist; счёт по
    Distributed сходится с числом сгенерированных событий. Счёт брать по
    своему прогону — по диапазону _load_ts или на чистом стенде: дедупа
    у сырья нет, и вторая заливка того же дня честно его удваивает.
  • Разовая. Топик прочитан обеими нодами — посмотреть consumer_host
    у доехавшего дня и записать увиденное в тело PR. Чекбоксом на каждый
    прогон это не делается: раскладку разделов решает ребаланс, а короткая
    пачка без ключа, если клиент раскладывает липко, вообще уезжает в один
    раздел — красное тут значило бы «сегодня так легло», а не «сломано».
  • Разовая. Живой день: поток идёт с темпом ×60, лаг виден в логе, флаг
    ускорения работает.
  • Форма на проводе соблюдена: в колонке raw даты читаются глазами
    (2026-06-01, 2026-06-01T12:34:56Z), ecommerce лежит строкой.
  • Образ генератора собирается и работает на стенде: разовая служба
    проигрывает день, каталог товаров в образе находится. make up,
    make smoke и make check-services остаются зелёными — образ Airflow
    этот тикет не трогает, и красное там значило бы, что тронул.
  • Решения тикета внесены в спеку: интерфейс запуска, размещение, клиент
    Kafka, ключ сообщения и способ взять ограниченную пачку — в раздел 9
    решениями, а открытый пункт про интерфейс запуска там закрыт; хвост про
    позицию на оси — в раздел 8; правки абзацев раздела 4 (порядок тикетов,
    снятая фраза про тест приёмника) сделаны.
  • Документация обновлена тем же PR: быстрый старт в README — как позвать
    генератор и чем режимы отличаются. Карту целей в
    docs/architecture/testing.md цель запуска не пополняет: там живут
    проверки, а это не проверка. Make-цели добавлены, каталог tmp/
    заведён в .gitignore.

Границы

  • Новых топиков нет; заказы — этап 3.
  • Снимок, мини-манифест и заливка при make up — тикет #42, идущий после
    #43 (порядок в чек-листе #4). Обвязка запуска заводится здесь: образ
    генератора, разовая служба compose под профилем и цели make. Иначе Сериализатор, приёмники и проигрыватель с CLI (#41)
    не проверит собственные разовые критерии. #42 берёт готовый запуск и
    включает его в make up.
  • Даги world_init/next_day и проброс сокета докера в Airflow — этап 5.
    Здесь только решение о том, что запуск идёт контейнером, и хвост в спеку.
  • Типизированный слой ODS не трогать: его приносит #43.

Сначала прочитать

  • docs/specs/2026-08-01-generator.md — разделы 4, 5, 9.
  • docs/research/2026-08-01-python-batch-generation-speed.md.
  • docs/architecture/storage.md — раздел «Приём»: что уже принимает сырьё и
    с какими свойствами.
  • Спорные API (confluent-kafka, настройки отправителя, партиционер) сверять
    через MCP Context7.
  • Образец обвязки — стенд v1, рабочая копия рядом
    (../clickstream-ch-kafka-superset-demo): служба generator в
    docker-compose.yml под профилем live-generator и её generator/Dockerfile
    (там uv и свой venv в образе). Переносится оттуда форма службы, не модель
    мира. А вот способ запуска у нас другой: v1 звал генератор импортом внутри
    воркера (airflow/dags/world_init_dag.py), мы запускаем контейнер.

Проверка

  • make test, make lint, make typecheck
  • make up с нуля, make smoke, make check-services — стенд не поехал
  • пакетная заливка дня, SELECT-счёт из stg.hits_raw_dist
  • SELECT колонки raw глазами: даты в ISO, ecommerce строкой
Part of #4. ## Цель Довести события до стенда: единственный канонический сериализатор, глупые приёмники (файл, Kafka пачкой, Kafka с темпом), проигрыватель и интерфейс запуска. После этого тикета топик `hits` наполняется настоящими данными генератора в обоих режимах. Тикет идёт **до** типизированного ODS (#43): порядок развёрнут 6 августа 2026 года. Раньше #43 шёл первым и вынужден был выдумывать событие сам — своим модулем, печатающим JSON, то есть вторым местом сериализации против правила «сериализатор один». Теперь отправитель идёт первым, а хранилище принимает настоящие данные. Разбор наверх (`ods.event`) появится в #43, и утверждения про него живут там же. ## Что войдёт - Канонический сериализатор на orjson: единственное место, где **событие целиком** превращается в JSON; канонический порядок ключей и строк; прямых `json.dumps` в коде нет. Граница правила: вложенный блок `ecommerce` собирает день (`commerce.py`, решение #40), и это не второй сериализатор — правило «сериализатор один» про событие, а не про блок внутри него. - **Форма на проводе — ISO-8601**, задана разделом 4 спеки генератора: `EventDate` как `2026-06-01`, `UTCEventTime` как `2026-06-01T12:34:56Z`, `ecommerce` — строка с экранированным JSON внутри. Этот тикет реализует форму первым; хранилище (#43) потом читает то, что он положил. Довод за ISO — читаемость сырья: весь смысл слоя STG в том, что менти открывает колонку `raw` в обычном клиенте и разбирает событие глазами. - Приёмники: файл (**одно событие — одна строка**, файл читается построчно), Kafka пачкой (пакетный режим), Kafka с темпом ×60 по умолчанию (живой день; ускорение — флаг проигрывателя, лаг логируется). Путь файла, адрес брокера и имя топика — параметры запуска, а не константы в коде: генератор задуман отдельной сущностью, которую можно унести в другой проект. Выделенного места файлу в репозитории не отводится — спека (раздел 5) обещает «потерял — пересчитал»; под мусорные прогоны заводится каталог `tmp/` целиком в `.gitignore`, чтобы файл не рвался в git. - **Событие несёт все 47 ключей всегда.** «Пусто» по контракту — пустое значение: пустой массив, пустая строка, ноль, — а не отсутствие ключа. Требование ADR 0005: строгий приём хранилища держится на сверке набора ключей, и пропавший ключ уводит событие в брак целиком. Раньше это обязательство пинил бы #43 своим модулем; после разворота его несёт этот тикет. - Переигровка дня целиком при обрыве: повтор даёт те же `WatchID`. - Тайминги: генерация и доставка печатаются отдельно. - **Свой образ генератора.** База `python:3.14` (генератор пинит `>=3.14,<3.15`, лок — `==3.14.*`), зависимости ставятся `uv sync --frozen` из `uv.lock`: побайтовая воспроизводимость привязана к версиям (раздел 4 спеки), и второй список пакетов писал бы на стенде не те байты, что сторожит `make test`. Спека сходится сама — раздел 2 обещает детерминизм «внутри канонического контейнера», им и становится этот образ. Образ Airflow этот тикет не трогает вовсе: ни базу, ни контекст сборки, ни его список пакетов для пробников. - **Каталог товаров должен оказаться в образе в той же взаимной раскладке**, что и пакет: путь считается от модуля вверх по дереву (`catalog.py:41`, `parents[3]`), а каталог по решению #39 — часть мира, а не запуска, и снаружи не настраивается. Копировать в образ или монтировать томом — решить при исполнении, но решить явно: иначе служба падает на первом же импорте. - **`.dockerignore` заводится тем же PR.** Сейчас его нет, а в контекст сборки иначе уедут `generator/.venv` (125 МБ) и мусорный `tmp/` — кеш слоёв будет сбрасываться на каждом прогоне. - **Разовая служба compose на образе генератора**, под профилем (либо через `docker compose run --rm`): в `make up` до #42 она не входит, обычная служба поднялась бы при `up` сама и стёрла границу с соседом. Оговорка для #42: `up --wait` считает успешную разовую службу упавшей, если у неё нет зависимого с `service_completed_successfully` — предупреждение об этом уже стоит в `compose.yaml` у `clickhouse-init`. ## Решить и внести в спеку генератора (раздел 9) тем же PR Часть пунктов решена грилингом 7 августа 2026 года — их не переоткрывать, а внести в спеку как есть; остальные решаются при исполнении. - Интерфейс запуска: CLI-команды и/или цели make, как делятся режимы проигрывателя. Считать, что зовущий — контейнер: параметры приходят аргументами или окружением, итог виден кодом возврата, а логи идут в стандартный вывод. Даги `world_init`/`next_day` запустят его так на этапе 5 — интерфейс не должен им мешать. - **Размещение решено грилингом 7 августа 2026 года; внести в спеку как есть.** Генератор живёт **в своём контейнере**, и на этапе 5 даг будет запускать контейнер, а не импортировать пакет в процессе воркера. Довод учебный: так граница между оркестратором и задачей видна глазами — у задачи своё окружение, свой образ, свой процесс, а оркестратор передаёт параметры и читает результат. Это та же ступенька, с которой в бою вырастает `KubernetesPodOperator`. Побочно это единственный вариант, при котором «генератор — отдельная и заменяемая сущность» перестаёт быть словами: его зависимости не смешиваются с окружением Airflow и его провайдеров. Хостовый `uv` остаётся для `make test`, `make lint` и `make typecheck`. Отклонено: *генератор в образе Airflow, даг импортирует пакет* (так у стенда v1) — механика проще, но менти видит лишь «даг позвал функцию», а платить пришлось бы сменой базового образа всего стенда на `python3.14`, переездом контекста сборки и смешением зависимостей генератора с зависимостями Airflow. *`ExternalPythonOperator` — задача в отдельном venv внутри образа Airflow*: изоляцию даёт и сокета не требует, но границу видно в конфигурации, а не в `docker ps`; остаётся запасным путём, если проброс сокета однажды окажется неприемлемым. Цена названа и принимается: чтобы даг запускал контейнер, в Airflow пробрасывается сокет докера, а это доступ, равный root на хосте. На учебном стенде размен допустим — но записывается уроком, а не прячется: в бою так не делают. Механика этапа 5 (провайдер `docker`, сеть compose у запускаемого контейнера, доступ к группе сокета) проверяется там же, не здесь. - **Клиент Kafka — `confluent-kafka`** (решение того же грилинга). Он уже стоит в образе Airflow, на нём написан пробник `dags/test_kafka.py`, и его же использует официальный провайдер Kafka для Airflow. Один клиент на стенде вместо двух. **Объявляется необязательной группой зависимостей генератора**, а образ ставит её из лока: приёмник импортирует клиента, и необъявленный импорт отнял бы у пакета переносимость — довод, на котором стоит половина решений этого тикета. Список пакетов в образе Airflow (`clickhouse-connect`, тот же `confluent-kafka`) правилу «зависимости из лока» не противоречит и не трогается: он для пробников `dags/`, а к генератору отношения не имеет. - **CLI без состояния.** Зерно и номер дня приходят параметрами; позицию на оси проигрыватель не хранит и не выводит из данных — её ведёт тот, кто зовёт. Потребителей интерфейса трое, и третий обычно забывается: даги `world_init`/`next_day` этапа 5, проверки #43 — и **заливка зернового мира из #42, восемь дней подряд**. Восемь запусков службы внутри `make up` — плохой ответ; способ проиграть подряд несколько дней одним запуском назначается здесь, вместе с остальным интерфейсом. Хвост этапу 5 (внести в раздел 8 спеки): позиция живёт переменной Airflow, отсутствие переменной означает мир в зерновом состоянии, заводит и двигает её только даг `next_day` и только по успеху. Довод — генератор отдельная и заменяемая сущность, к хранилищу его привязывать не надо, а `make clean` сносит том метаданных Airflow вместе с данными ClickHouse, так что позиция и мир чистятся одной командой. Расхождение переменной с данными (менти почистил партицию руками) лечится документацией, а не сторожем: постоянная сверка была бы проверкой, красной в норме (ADR 0002). - **Поправить абзац раздела 4 про порядок тикетов.** Там записано: «Форму пинит хранилище (#43) как первый потребитель… Порядок тикетов обратный порядку зависимости, поэтому здесь она и записана». После разворота 6 августа 2026 года это неверно: форму реализует #41, а хранилище её читает. Сама запись формы в спеке остаётся — она и есть источник истины, просто довод у неё теперь другой. - **Ключ сообщения Kafka: назначить или явно отказаться.** Довод, стоявший здесь раньше, — «у сообщений без ключа работает липкое распределение: подряд идущая пачка ложится в один раздел» — оказался свойством клиента, а не Kafka: у `kafka-python` сообщение без ключа уходит в случайный раздел на каждое сообщение, липкий партиционер (KIP-480) там опция (сверено через Context7 7 августа 2026 года). Для librdkafka, то есть для выбранного `confluent-kafka`, ответ пока неизвестен: в `CONFIGURATION.md` есть `sticky.partitioning.linger.ms`, но умолчание документация не называет, а умолчание топикового партиционера — `consistent_random`, то есть случайный раздел. **Сверить до того**, как довод ляжет в спеку, и опереть его на `CONFIGURATION.md`, а не на общую справку. Решение (`WatchID` ключом либо «ключа нет, и вот почему») интерфейсное, и место ему здесь. - **Снять фразу «Сторожится тестом приёмника»** в абзаце «Одно событие — одно сообщение Kafka» (раздел 4). Постоянного теста приёмника не будет — решение владельца 7 августа 2026 года: свойство уже сторожат постоянная проверка инварианта слоёв из #43 на настоящих данных и разовая приёмка ниже, а дешёвая половина достаётся построчностью файла в тесте детерминизма. Второй сторож того же свойства с подставным продюсером ничему не учит и живёт зря. Само утверждение в спеке остаётся: это контракт транспорта, а не деталь реализации. - **Способ взять ограниченную пачку** — сколько-то событий, а не день целиком. Довод от менти: тот, кто хочет посмотреть на конвейер, не должен ждать полсотни тысяч событий. На живой день пачка не распространяется — там ожидание снимает флаг ускорения, который и так обязан быть проверен. Три свойства обязательны, потому что на них стоят проверки #43: пачка сочетается с **файловым** приёмником, даёт взять срез, ещё не приезжавший в ODS, и **печатает число отправленного** — счёт нужен обеим сторонам. Флаг, срез дня или отдельный режим — решить здесь, потому что здесь живёт интерфейс. Записывать в раздел 9 **решением**, а открытый пункт про интерфейс запуска — закрыть: дописывать в открытый список нельзя, за #42 остаются два последних пункта, и раздел он закрывает целиком. ## Критерии приёмки Приёмка идёт **по слою сырья**: `ods.event` и `*_errors` появятся в #43, и утверждения про них — его критерии, не эти. Критерии ниже, помеченные **разовыми**, требуют прогона целого дня через стенд и постоянными целями `make` быть не могут — по тому же правилу, по которому #43 не ставит день в повторяемую цель. Их след — запись в теле PR. **Убирать за этими прогонами не нужно, и вот почему.** Правило `docs/architecture/testing.md` требует, чтобы вкладывающая данные проверка либо прибиралась за собой, либо не жила постоянной целью: вложенное оседает в мире менти и врёт его лабам. Здесь оно не оседает — на момент этого тикета типизированного ODS ещё нет (его приносит #43), день доезжает только до `stg.hits_raw_dist`, а у сырья срок жизни трое суток, и оно уходит само. Задним числом матвью не досыпают, так что в `ods.event` этот день не появится и потом. Оговорка на будущее: если день гонять уже после того, как #43 слит, он доедет до ODS и останется там навсегда — тогда действует общее правило, и день берётся за границей оси мира со сносом партиции. - [x] Побайтовый детерминизм: два прогона дня в файл — идентичные байты, сравнение хешей (тест). Тем же тестом — построчность: строк в файле ровно столько, сколько событий в дне. Это дешёвая половина контракта транспорта: склейка событий сломала бы разбор целиком, потому что хранилище читает топик байтами. - [x] **Разовая.** Пакетный режим: день доезжает до `stg.hits_raw_dist`; счёт по Distributed сходится с числом сгенерированных событий. Счёт брать по своему прогону — по диапазону `_load_ts` или на чистом стенде: дедупа у сырья нет, и вторая заливка того же дня честно его удваивает. - [x] **Разовая.** Топик прочитан обеими нодами — посмотреть `consumer_host` у доехавшего дня и записать увиденное в тело PR. Чекбоксом на каждый прогон это не делается: раскладку разделов решает ребаланс, а короткая пачка без ключа, если клиент раскладывает липко, вообще уезжает в один раздел — красное тут значило бы «сегодня так легло», а не «сломано». - [x] **Разовая.** Живой день: поток идёт с темпом ×60, лаг виден в логе, флаг ускорения работает. - [x] Форма на проводе соблюдена: в колонке `raw` даты читаются глазами (`2026-06-01`, `2026-06-01T12:34:56Z`), `ecommerce` лежит строкой. - [x] Образ генератора собирается и работает на стенде: разовая служба проигрывает день, каталог товаров в образе находится. `make up`, `make smoke` и `make check-services` остаются зелёными — образ Airflow этот тикет не трогает, и красное там значило бы, что тронул. - [x] Решения тикета внесены в спеку: интерфейс запуска, размещение, клиент Kafka, ключ сообщения и способ взять ограниченную пачку — в раздел 9 решениями, а открытый пункт про интерфейс запуска там закрыт; хвост про позицию на оси — в раздел 8; правки абзацев раздела 4 (порядок тикетов, снятая фраза про тест приёмника) сделаны. - [x] Документация обновлена тем же PR: быстрый старт в README — как позвать генератор и чем режимы отличаются. Карту целей в `docs/architecture/testing.md` цель запуска не пополняет: там живут проверки, а это не проверка. Make-цели добавлены, каталог `tmp/` заведён в `.gitignore`. ## Границы - Новых топиков нет; заказы — этап 3. - Снимок, мини-манифест и заливка при `make up` — тикет #42, идущий после #43 (порядок в чек-листе #4). Обвязка запуска заводится здесь: образ генератора, разовая служба compose под профилем и цели `make`. Иначе #41 не проверит собственные разовые критерии. #42 берёт готовый запуск и включает его в `make up`. - Даги `world_init`/`next_day` и проброс сокета докера в Airflow — этап 5. Здесь только решение о том, что запуск идёт контейнером, и хвост в спеку. - Типизированный слой ODS не трогать: его приносит #43. ## Сначала прочитать - docs/specs/2026-08-01-generator.md — разделы 4, 5, 9. - docs/research/2026-08-01-python-batch-generation-speed.md. - docs/architecture/storage.md — раздел «Приём»: что уже принимает сырьё и с какими свойствами. - Спорные API (`confluent-kafka`, настройки отправителя, партиционер) сверять через MCP Context7. - Образец обвязки — стенд v1, рабочая копия рядом (`../clickstream-ch-kafka-superset-demo`): служба `generator` в `docker-compose.yml` под профилем `live-generator` и её `generator/Dockerfile` (там `uv` и свой venv в образе). Переносится оттуда форма службы, не модель мира. А вот способ запуска у нас другой: v1 звал генератор импортом внутри воркера (`airflow/dags/world_init_dag.py`), мы запускаем контейнер. ## Проверка - `make test`, `make lint`, `make typecheck` - `make up` с нуля, `make smoke`, `make check-services` — стенд не поехал - пакетная заливка дня, `SELECT`-счёт из `stg.hits_raw_dist` - `SELECT` колонки `raw` глазами: даты в ISO, `ecommerce` строкой
ddmitry added the ready-for-agent label 2026-08-01 21:20:03 +03:00
ddmitry added a new dependency 2026-08-01 21:20:07 +03:00
Author
Owner

Тикет переставлен вперёд #43: теперь producer идёт первым, а хранилище
принимает настоящие данные. Зависимость в Gitea развёрнута — #41 больше
не ждёт #43, а блокирует его.

Почему. Прежний порядок стоял на доводе «форму на проводе пинит первый
потребитель». Довод отработал: форма записана словами в разделе 4 спеки
5 августа. Осталась только цена — #43 пришлось бы выдумывать событие
собственным модулем, печатающим JSON, то есть заводить второе место
сериализации против правила «сериализатор один» из того же раздела 4.
Вместе с ним тянулись правило «тип ClickHouse → значение на проводе»,
вопрос о дате фикстуры и новая зависимость проверки от uv. Ни одного из
этих вопросов не существует, если #41 идёт первым.

Три критерия уехали в #43. «День доезжает до ods.event», «*_errors
пуст на честном прогоне» и «повторная заливка не меняет счёт после дедупа» —
это утверждения про хранилище, а не про сериализатор, и проверить их можно
только там, где есть ODS. Из-за них и стояла обратная зависимость. Теперь
приёмка #41 идёт по слою сырья, и круг разомкнут.

Что добавилось. Форма на проводе названа явно — этот тикет реализует её
первым. Проверка «топик читают обе ноды» через consumer_host. Проверка,
что даты в raw читаются глазами: на этом стоит весь довод за ISO. И два
пункта в «решить и внести в спеку» — способ взять ограниченную пачку (он
нужен проверкам #43, а день в полсотни тысяч событий в повторяемую цель не
поставишь) и правка абзаца раздела 4, который после разворота стал неверен.

Тикет переставлен вперёд #43: теперь producer идёт первым, а хранилище принимает настоящие данные. Зависимость в Gitea развёрнута — #41 больше не ждёт #43, а блокирует его. **Почему.** Прежний порядок стоял на доводе «форму на проводе пинит первый потребитель». Довод отработал: форма записана словами в разделе 4 спеки 5 августа. Осталась только цена — #43 пришлось бы выдумывать событие собственным модулем, печатающим JSON, то есть заводить второе место сериализации против правила «сериализатор один» из того же раздела 4. Вместе с ним тянулись правило «тип ClickHouse → значение на проводе», вопрос о дате фикстуры и новая зависимость проверки от `uv`. Ни одного из этих вопросов не существует, если #41 идёт первым. **Три критерия уехали в #43.** «День доезжает до `ods.event`», «`*_errors` пуст на честном прогоне» и «повторная заливка не меняет счёт после дедупа» — это утверждения про хранилище, а не про сериализатор, и проверить их можно только там, где есть ODS. Из-за них и стояла обратная зависимость. Теперь приёмка #41 идёт по слою сырья, и круг разомкнут. **Что добавилось.** Форма на проводе названа явно — этот тикет реализует её первым. Проверка «топик читают обе ноды» через `consumer_host`. Проверка, что даты в `raw` читаются глазами: на этом стоит весь довод за ISO. И два пункта в «решить и внести в спеку» — способ взять ограниченную пачку (он нужен проверкам #43, а день в полсотни тысяч событий в повторяемую цель не поставишь) и правка абзаца раздела 4, который после разворота стал неверен.
Author
Owner

Два критерия сняты решением владельца.

«Топик читают обе ноды» перестал быть чекбоксом на каждый прогон и стал разовым
наблюдением с записью в тело PR. Довод владельца: прогнали день, увидели обе
ноды в consumer_host, поняли, что группа потребителей настроена верно, — и
довольно; настройка сама не портится. Довод против чекбокса: раскладку разделов
решает ребаланс, а у сообщений без ключа короткая пачка вообще уезжает в один
раздел, так что красное значило бы «сегодня так легло», а не «сломано».

«Повторная заливка даёт те же WatchID» снят целиком. Это побайтовый
детерминизм, пересказанный через полный стенд, Kafka и SELECT. Свойство чистой
функции уже доказано сравнением хешей файла — первым критерием, дёшево и без
стенда. Транспорт WatchID не трогает, а дедуп, ради которого свойство нужно,
живёт в #43 и проверяется там.

Ранее в тот же заход по холодному ревью: обещан формат файлового приёмника —
одно событие на строку, на этом стоит оснастка проверок #43; выписано
обязательство ADR 0005 слать все 47 ключей всегда, которое после разворота
порядка стало обязательством этого тикета; заведён тест приёмника «одно событие
— одно сообщение», который спека прямо называет сторожем контракта транспорта;
добавлено решение о ключе сообщения Kafka; у «ограниченной пачки» назван довод
от менти и три свойства, на которых стоят проверки соседа; критерии, требующие
прогона целого дня, помечены разовыми.

**Два критерия сняты решением владельца.** «Топик читают обе ноды» перестал быть чекбоксом на каждый прогон и стал разовым наблюдением с записью в тело PR. Довод владельца: прогнали день, увидели обе ноды в `consumer_host`, поняли, что группа потребителей настроена верно, — и довольно; настройка сама не портится. Довод против чекбокса: раскладку разделов решает ребаланс, а у сообщений без ключа короткая пачка вообще уезжает в один раздел, так что красное значило бы «сегодня так легло», а не «сломано». «Повторная заливка даёт те же `WatchID`» снят целиком. Это побайтовый детерминизм, пересказанный через полный стенд, Kafka и SELECT. Свойство чистой функции уже доказано сравнением хешей файла — первым критерием, дёшево и без стенда. Транспорт `WatchID` не трогает, а дедуп, ради которого свойство нужно, живёт в #43 и проверяется там. Ранее в тот же заход по холодному ревью: обещан формат файлового приёмника — одно событие на строку, на этом стоит оснастка проверок #43; выписано обязательство ADR 0005 слать все 47 ключей всегда, которое после разворота порядка стало обязательством этого тикета; заведён тест приёмника «одно событие — одно сообщение», который спека прямо называет сторожем контракта транспорта; добавлено решение о ключе сообщения Kafka; у «ограниченной пачки» назван довод от менти и три свойства, на которых стоят проверки соседа; критерии, требующие прогона целого дня, помечены разовыми.
Author
Owner

Грилинг постановки 7 августа 2026 года: шесть решений, тело тикета правлено.

Размещение — как у стенда v1, велосипед не изобретаем. У предшественника
генератор живёт в образе Airflow: исходник примонтирован в контейнеры, даг
world_init импортирует пакет и зовёт его в процессе воркера, а отдельная
служба генератора существует только под живой непрерывный поток. Берём эту
схему. Прогон на стенде — разовая служба compose на том же образе, по образцу
kafka-init и clickhouse-init. Отдельный образ отклонён: вторая сборка ради
того же кода, и этап 5 всё равно пришлось бы решать заново.

Одно следствие пришлось выписать отдельно: зависимости в образ ставятся из
uv.lock
, а не вторым списком в Dockerfile. Побайтовая воспроизводимость
привязана к версии orjson (раздел 4 спеки), и разъехавшийся образ писал бы на
стенде не те байты, что сторожит make test — молча, до первого красного
манифеста в #42.

Клиент Kafka — confluent-kafka. Он уже в образе Airflow, на нём написан
пробник, его же использует официальный провайдер. Заодно поправлен довод у
вопроса про ключ сообщения: «липкое распределение у сообщений без ключа» —
свойство клиента, а не Kafka. У kafka-python сообщение без ключа уходит в
случайный раздел на каждое сообщение, липкость там опция (сверено через
Context7). У librdkafka липкость есть, но её умолчание надо сверить до того,
как довод ляжет в спеку.

Граница правила «сериализатор один» названа. commerce.py уже вызывает
orjson.dumps, собирая вложенный блок ecommerce (решение #40). Это не второй
сериализатор: правило про событие целиком, а не про блок внутри него. Без этой
фразы ревьюер честно предъявил бы commerce.py, а исполнитель мог бы потащить
сборку ecommerce в новый модуль.

CLI без состояния, позиция на оси — переменной Airflow. Генератор — отдельная
и потенциально переиспользуемая сущность, привязывать его к ClickHouse не надо.
Поэтому и позицию не выводим из данных (max(EventDate) был рассмотрен и
отклонён): мир ведёт пульт, а не хранилище. Переменная Airflow отвечает и на
вопрос про чистку — make clean сносит том метаданных Airflow вместе с данными
ClickHouse, так что позиция и мир исчезают одной командой. Отсутствие переменной
означает мир в зерновом состоянии: тогда заливке при make up не нужно ничего
знать про Airflow. Компакт-топик, как в v1, закрыт спекой заранее — Kafka у нас
труба, не хранилище.

Постоянный тест «одно событие — одно сообщение Kafka» снят. Свойство уже
сторожат постоянная проверка инварианта слоёв из #43 на настоящих данных и
разовая приёмка этого тикета, а дешёвая половина достаётся построчностью файла
в тесте детерминизма. Второй сторож того же свойства с подставным продюсером
менти ничему не учит. Фразу «Сторожится тестом приёмника» из раздела 4 спеки
снять тем же PR; само утверждение остаётся — это контракт транспорта.

Файловому приёмнику места в репозитории не отводится. Разово нужен файл, а
не приёмник: путь становится параметром запуска, а под мусорные прогоны
заводится каталог tmp/ целиком в .gitignore, чтобы файл не мозолил глаза и
не рвался в git. Сам приёмник остаётся — на нём стоят фикстуры порчи для #43,
побайтовый детерминизм и возможность посмотреть событие глазами без стенда;
убрать его значило бы вернуть #43 собственный дамп, то есть второе место
сериализации.

Граница с #42 уточнена. Обвязка запуска — зависимости в образе,
монтирование, разовая служба, цели make — заводится здесь: без неё #41 не
проверит собственные разовые критерии. #42 берёт готовый запуск и включает его
в make up.

**Грилинг постановки 7 августа 2026 года: шесть решений, тело тикета правлено.** **Размещение — как у стенда v1, велосипед не изобретаем.** У предшественника генератор живёт в образе Airflow: исходник примонтирован в контейнеры, даг `world_init` импортирует пакет и зовёт его в процессе воркера, а отдельная служба генератора существует только под живой непрерывный поток. Берём эту схему. Прогон на стенде — разовая служба compose на том же образе, по образцу `kafka-init` и `clickhouse-init`. Отдельный образ отклонён: вторая сборка ради того же кода, и этап 5 всё равно пришлось бы решать заново. Одно следствие пришлось выписать отдельно: **зависимости в образ ставятся из `uv.lock`**, а не вторым списком в `Dockerfile`. Побайтовая воспроизводимость привязана к версии orjson (раздел 4 спеки), и разъехавшийся образ писал бы на стенде не те байты, что сторожит `make test` — молча, до первого красного манифеста в #42. **Клиент Kafka — `confluent-kafka`.** Он уже в образе Airflow, на нём написан пробник, его же использует официальный провайдер. Заодно поправлен довод у вопроса про ключ сообщения: «липкое распределение у сообщений без ключа» — свойство клиента, а не Kafka. У `kafka-python` сообщение без ключа уходит в случайный раздел на каждое сообщение, липкость там опция (сверено через Context7). У librdkafka липкость есть, но её умолчание надо сверить до того, как довод ляжет в спеку. **Граница правила «сериализатор один» названа.** `commerce.py` уже вызывает `orjson.dumps`, собирая вложенный блок `ecommerce` (решение #40). Это не второй сериализатор: правило про событие целиком, а не про блок внутри него. Без этой фразы ревьюер честно предъявил бы `commerce.py`, а исполнитель мог бы потащить сборку `ecommerce` в новый модуль. **CLI без состояния, позиция на оси — переменной Airflow.** Генератор — отдельная и потенциально переиспользуемая сущность, привязывать его к ClickHouse не надо. Поэтому и позицию не выводим из данных (`max(EventDate)` был рассмотрен и отклонён): мир ведёт пульт, а не хранилище. Переменная Airflow отвечает и на вопрос про чистку — `make clean` сносит том метаданных Airflow вместе с данными ClickHouse, так что позиция и мир исчезают одной командой. Отсутствие переменной означает мир в зерновом состоянии: тогда заливке при `make up` не нужно ничего знать про Airflow. Компакт-топик, как в v1, закрыт спекой заранее — Kafka у нас труба, не хранилище. **Постоянный тест «одно событие — одно сообщение Kafka» снят.** Свойство уже сторожат постоянная проверка инварианта слоёв из #43 на настоящих данных и разовая приёмка этого тикета, а дешёвая половина достаётся построчностью файла в тесте детерминизма. Второй сторож того же свойства с подставным продюсером менти ничему не учит. Фразу «Сторожится тестом приёмника» из раздела 4 спеки снять тем же PR; само утверждение остаётся — это контракт транспорта. **Файловому приёмнику места в репозитории не отводится.** Разово нужен файл, а не приёмник: путь становится параметром запуска, а под мусорные прогоны заводится каталог `tmp/` целиком в `.gitignore`, чтобы файл не мозолил глаза и не рвался в git. Сам приёмник остаётся — на нём стоят фикстуры порчи для #43, побайтовый детерминизм и возможность посмотреть событие глазами без стенда; убрать его значило бы вернуть #43 собственный дамп, то есть второе место сериализации. **Граница с #42 уточнена.** Обвязка запуска — зависимости в образе, монтирование, разовая служба, цели `make` — заводится здесь: без неё #41 не проверит собственные разовые критерии. #42 берёт готовый запуск и включает его в `make up`.
Author
Owner

Тёплое ревью постановки: восемь находок, тело правлено второй раз.

Ревью получило бриф с решениями и доводами, тикет читало само. Три самые
весомые находки перепроверены на месте, не приняты на слово.

Размещение в образе Airflow упиралось в интерпретатор. Штатный
apache/airflow:3.3.0 несёт Python 3.13.14, а генератор пинит >=3.14,<3.15,
и uv.lock==3.14.*. Из этого лока в тот образ не встало бы ничего, а
примонтированный исходник импортировался бы интерпретатором, который пакет от
себя отрезал. Решение уцелело: у Airflow 3.3.0 есть вариант образа под 3.14,
тег в реестре проверен (с контролем: заведомо несуществующий -python9.9
командой отвергается), база задаётся аргументом AIRFLOW_BASE_IMAGE. Спека
сходится сама — раздел 2 обещает детерминизм «внутри канонического
контейнера», им и становится образ Airflow. Запасной путь на случай, если
Airflow на 3.14 окажется битым, назван в тикете и объявлен вопросом владельцу.

Контекст сборки не видел лока. build.context: ./infra/airflow, а
uv.lock лежит в generator/ — требование «ставить из лока» было физически
невыполнимо. Контекст переезжает в корень.

up --wait не дожидается разовых служб сам. Предупреждение об этом уже
стоит в compose.yaml у clickhouse-init: без зависимого с
service_completed_successfully успешная разовая служба считается упавшей.
Мой «образец kafka-init» работал только за счёт зависимых. Вдобавок обычная
служба поднимается при up по умолчанию — то есть, заведя её здесь, мы бы
включили заливку в make up и стёрли границу с #42. Служба заводится под
профилем; зависимого придумывает #42, когда будет включать.

confluent-kafka получил дом. Он объявляется необязательной группой
зависимостей генератора, а Dockerfile берёт версию оттуда. Сейчас клиент
прописан в образе руками — ровно тем вторым списком, который этот же тикет
запретил; необъявленный импорт в приёмнике отнял бы у пакета переносимость.

Мелкие правки. Способ взять ограниченную пачку записывается в раздел 9
решением, а открытый пункт про интерфейс запуска там закрывается: дописывать в
открытый список нельзя, #42 стоит на том, что закрывает последний вопрос.
Утверждение «у librdkafka липкость есть» снято — в CONFIGURATION.md есть
sticky.partitioning.linger.ms, но умолчание не документировано, а умолчание
топикового партиционера consistent_random кладёт сообщение без ключа в
случайный раздел; остаётся требование сверить. «Следующий тикет (#42
исправлен: по чек-листу #4 порядок #41#43#42. Заведены критерии на
прогон стенда после смены базы образа и на быстрый старт в README; карту целей
testing.md цель запуска не пополняет — там живут проверки, а это не проверка.

Довод про ограниченную пачку и живой день снят (решение владельца): ждать
двадцать четыре минуты не приходится и без пачки — ожидание снимает флаг
ускорения, который и так обязан быть проверен. На живой день пачка не
распространяется.

Ревью сверило и нашло верным: commerce.py:537 — единственный вызов orjson
вне будущего сериализатора; make clean действительно сносит том метаданных
Airflow вместе с данными ClickHouse; в контракте ровно 47 колонок; цитаты из
спеки существуют дословно; дыр между #41 и #43 по приёмке нет.

**Тёплое ревью постановки: восемь находок, тело правлено второй раз.** Ревью получило бриф с решениями и доводами, тикет читало само. Три самые весомые находки перепроверены на месте, не приняты на слово. **Размещение в образе Airflow упиралось в интерпретатор.** Штатный `apache/airflow:3.3.0` несёт Python 3.13.14, а генератор пинит `>=3.14,<3.15`, и `uv.lock` — `==3.14.*`. Из этого лока в тот образ не встало бы ничего, а примонтированный исходник импортировался бы интерпретатором, который пакет от себя отрезал. Решение уцелело: у Airflow 3.3.0 есть вариант образа под 3.14, тег в реестре проверен (с контролем: заведомо несуществующий `-python9.9` командой отвергается), база задаётся аргументом `AIRFLOW_BASE_IMAGE`. Спека сходится сама — раздел 2 обещает детерминизм «внутри канонического контейнера», им и становится образ Airflow. Запасной путь на случай, если Airflow на 3.14 окажется битым, назван в тикете и объявлен вопросом владельцу. **Контекст сборки не видел лока.** `build.context: ./infra/airflow`, а `uv.lock` лежит в `generator/` — требование «ставить из лока» было физически невыполнимо. Контекст переезжает в корень. **`up --wait` не дожидается разовых служб сам.** Предупреждение об этом уже стоит в `compose.yaml` у `clickhouse-init`: без зависимого с `service_completed_successfully` успешная разовая служба считается упавшей. Мой «образец `kafka-init`» работал только за счёт зависимых. Вдобавок обычная служба поднимается при `up` по умолчанию — то есть, заведя её здесь, мы бы включили заливку в `make up` и стёрли границу с #42. Служба заводится под профилем; зависимого придумывает #42, когда будет включать. **`confluent-kafka` получил дом.** Он объявляется необязательной группой зависимостей генератора, а `Dockerfile` берёт версию оттуда. Сейчас клиент прописан в образе руками — ровно тем вторым списком, который этот же тикет запретил; необъявленный импорт в приёмнике отнял бы у пакета переносимость. **Мелкие правки.** Способ взять ограниченную пачку записывается в раздел 9 решением, а открытый пункт про интерфейс запуска там закрывается: дописывать в открытый список нельзя, #42 стоит на том, что закрывает последний вопрос. Утверждение «у librdkafka липкость есть» снято — в `CONFIGURATION.md` есть `sticky.partitioning.linger.ms`, но умолчание не документировано, а умолчание топикового партиционера `consistent_random` кладёт сообщение без ключа в случайный раздел; остаётся требование сверить. «Следующий тикет (#42)» исправлен: по чек-листу #4 порядок #41 → #43 → #42. Заведены критерии на прогон стенда после смены базы образа и на быстрый старт в README; карту целей `testing.md` цель запуска не пополняет — там живут проверки, а это не проверка. **Довод про ограниченную пачку и живой день снят** (решение владельца): ждать двадцать четыре минуты не приходится и без пачки — ожидание снимает флаг ускорения, который и так обязан быть проверен. На живой день пачка не распространяется. Ревью сверило и нашло верным: `commerce.py:537` — единственный вызов `orjson` вне будущего сериализатора; `make clean` действительно сносит том метаданных Airflow вместе с данными ClickHouse; в контракте ровно 47 колонок; цитаты из спеки существуют дословно; дыр между #41 и #43 по приёмке нет.
Author
Owner

Развилка размещения переиграна: генератор живёт в своём контейнере.

Вчерашнее решение — генератор в образе Airflow, даг импортирует пакет —
отменено. Довод владельца учебный: при импорте менти видит только «даг позвал
функцию», а границы между оркестратором и задачей не наблюдает. При запуске
контейнером она видна глазами: своё окружение, свой образ, свой процесс,
параметры внутрь, код возврата наружу. Это та же ступенька, с которой в бою
вырастает KubernetesPodOperator, и единственный вариант, при котором
«генератор — отдельная и заменяемая сущность» перестаёт быть словами в тикете.

Счёт при этом сошёлся в ту же сторону. Холодное ревью показало, чего стоил
образ Airflow: смена базы всего стенда на python3.14, переезд контекста
сборки в корень, .dockerignore, зависимости генератора вперемешку с
провайдерами Airflow — пять правок инфраструктуры ради одного import. Со
своим образом три из них отпадают: образ Airflow не трогается вовсе, а правило
«зависимости из лока» перестаёт спорить с clickhouse-connect, который в
образе Airflow нужен пробникам и к генератору отношения не имеет.

Цена названа и принята: чтобы даг запускал контейнер, на этапе 5 в Airflow
пробрасывается сокет докера — доступ, равный root на хосте. На учебном стенде
размен допустим, но записывается уроком, а не прячется: в бою так не делают.
Запасной путь, если сокет однажды окажется неприемлемым, — ExternalPythonOperator
с отдельным venv внутри образа Airflow: изоляция та же, границу видно в
конфигурации, а не в docker ps.

Из холодного ревью связки #41/#42/#43 сюда легли три находки.

Каталог товаров ищется от исходника. catalog.py:41 считает путь от модуля
вверх по дереву (parents[3]), и каталог по решению #39 снаружи не
настраивается. В образе раскладку надо сохранить осознанно — копированием или
монтированием, — иначе служба падает на первом же импорте. Ни постановка, ни
тёплое ревью этого не видели.

.dockerignore в репозитории нет, а в контекст сборки иначе уедут
generator/.venv (125 МБ) и мусорный tmp/, сбрасывая кеш слоёв на каждом
прогоне.

У интерфейса запуска забыт третий потребитель. Кроме дагов этапа 5 и
проверок #43 его зовёт заливка зернового мира из #42 — восемь дней подряд.
Восемь запусков службы внутри make up были бы плохим ответом, а решение к
тому времени объявлено закрытым, так что способ проиграть несколько дней одним
запуском назначается здесь.

Заодно поправлено: «#42 закрывает последний вопрос раздела 9» — на самом деле
за ним остаются два из трёх.

**Развилка размещения переиграна: генератор живёт в своём контейнере.** Вчерашнее решение — генератор в образе Airflow, даг импортирует пакет — отменено. Довод владельца учебный: при импорте менти видит только «даг позвал функцию», а границы между оркестратором и задачей не наблюдает. При запуске контейнером она видна глазами: своё окружение, свой образ, свой процесс, параметры внутрь, код возврата наружу. Это та же ступенька, с которой в бою вырастает `KubernetesPodOperator`, и единственный вариант, при котором «генератор — отдельная и заменяемая сущность» перестаёт быть словами в тикете. **Счёт при этом сошёлся в ту же сторону.** Холодное ревью показало, чего стоил образ Airflow: смена базы всего стенда на `python3.14`, переезд контекста сборки в корень, `.dockerignore`, зависимости генератора вперемешку с провайдерами Airflow — пять правок инфраструктуры ради одного `import`. Со своим образом три из них отпадают: образ Airflow не трогается вовсе, а правило «зависимости из лока» перестаёт спорить с `clickhouse-connect`, который в образе Airflow нужен пробникам и к генератору отношения не имеет. **Цена названа и принята:** чтобы даг запускал контейнер, на этапе 5 в Airflow пробрасывается сокет докера — доступ, равный root на хосте. На учебном стенде размен допустим, но записывается уроком, а не прячется: в бою так не делают. Запасной путь, если сокет однажды окажется неприемлемым, — `ExternalPythonOperator` с отдельным venv внутри образа Airflow: изоляция та же, границу видно в конфигурации, а не в `docker ps`. **Из холодного ревью связки #41/#42/#43 сюда легли три находки.** *Каталог товаров ищется от исходника.* `catalog.py:41` считает путь от модуля вверх по дереву (`parents[3]`), и каталог по решению #39 снаружи не настраивается. В образе раскладку надо сохранить осознанно — копированием или монтированием, — иначе служба падает на первом же импорте. Ни постановка, ни тёплое ревью этого не видели. *`.dockerignore` в репозитории нет*, а в контекст сборки иначе уедут `generator/.venv` (125 МБ) и мусорный `tmp/`, сбрасывая кеш слоёв на каждом прогоне. *У интерфейса запуска забыт третий потребитель.* Кроме дагов этапа 5 и проверок #43 его зовёт заливка зернового мира из #42 — восемь дней подряд. Восемь запусков службы внутри `make up` были бы плохим ответом, а решение к тому времени объявлено закрытым, так что способ проиграть несколько дней одним запуском назначается здесь. Заодно поправлено: «#42 закрывает последний вопрос раздела 9» — на самом деле за ним остаются два из трёх.
Sign in to join this conversation.
1 Participants
Notifications
Due Date
No due date set.
Reference: ddmitry/clickstream-data-platform#41