Форма пробника Kafka: переписать по образцу test_clickhouse #22

Closed
opened 2026-07-31 19:21:01 +03:00 by ddmitry · 1 comment
Owner

Пробник Kafka — второй и последний пример DAG в стенде, и после #20 единственный,
который не показывает, как здесь пишут код. Он рабочий и проверяет ровно то, что
должен: менять надо текст, а не охват.

Тикет про то, как код читается, поэтому исполнитель — модель со вкусом
(AGENTS.md, «Кому что поручать»). По той же причине решения о форме приняты
здесь, а не оставлены на реализацию.

Что не так сейчас

  • Тело задачи — 60 строк одним куском внутри одного try: настройка продюсера,
    доставка, сверка доставки, настройка консьюмера, цикл чтения, сверка маркера.
    На вопрос «что проверяет пробник» файл отвечает только целиком.
  • Часовые producer = None / consumer = None и finally с двумя проверками
    на None — четыре строки на то, чтобы закрыть то, что могло не открыться.
  • Проверки не названы: if undelivered or delivery_errors or len(delivered_offsets) != 1 — три разных отказа в одном условии и с одним
    текстом ошибки.
  • Числа 10 000, 5 000, 10, 20 и 0.5 стоят без имён; что из них «сколько ждём
    доставки», а что «сколько ждём чтения», видно только из места.
  • Успех — return из середины цикла, отказ — после цикла.
  • Из уроков, которые файл мог бы дать про Kafka, в тексте видны два, и те по
    обрывкам.

Как переписываем

Форма — та же, что у test_clickhouse после #20: модульные константы, приватные
помощники, отдельные утверждения, короткое тело задачи. Менти открывает два
файла и видит один почерк — это и есть цель.

  1. Две задачи: write_marker пишет маркер и возвращает адрес записи,
    read_marker читает по этому адресу и сверяет. Половины круга ломаются
    по-разному, и в интерфейсе Airflow должно быть видно, какая из них отказала;
    заодно на живом примере виден XCom.
  2. Каждая задача заводит своего клиента и закрывает его сама: под
    LocalExecutor это два процесса, и одно подключение границу между ними не
    переживает. Продюсер и консьюмер — не два клиента одного вида, а два разных
    мира с разными настройками и разными способами сломаться; форма должна это
    показывать.
  3. Настройки клиентов — именованные модульные константы; у каждого срока
    ожидания имя, говорящее, чего именно мы ждём.
  4. Утверждения — функции с именами по способу отказа.
  5. Адрес записи (раздел и смещение) собирается NamedTuple уровня модуля:
    два соседних целых в сигнатуре переставляются молча, и проверка после этого
    краснеет загадочно. Границу задач тип не переходит — XCom идёт через JSON,
    и кортеж вернулся бы списком. Между задачами адрес едет именованными полями
    словаря вместе с маркером.
  6. Успех не возвращается из середины цикла.

Шапка файла

По образцу test_clickhouse (коммит c130319): «Тест проверяет:» списком того,
что утверждается, и абзац о том, чем тест не является — не образец боевого
потребителя: не подписывается на топик, не входит в группу, не коммитит
смещения, не создаёт топик. Шапка уходит в doc_md и видна в интерфейсе
Airflow. Правило комментариев из #20 строк описания DAG не касается.

Комментарии

Правило #20 остаётся в силе: комментарий объясняет решение, а не пересказывает
строку. Список мест расширяется явно — комментируем там, где незнакома модель
Kafka, и молчим там, где незнаком просто Python. Первое читатель из кода не
выведет никогда, второе выведет за минуту. Шесть мест:

  1. автосоздание топика (комментарий уже есть, у константы TOPIC);
  2. одноразовая группа и выключенный автокоммит: проверка не двигает ничьих
    смещений и не оставляет следов на брокере;
  3. три разных срока ожидания: сколько библиотека держит сообщение у себя,
    сколько ждём сеть, сколько ждём возврата;
  4. produce() не пишет, а ставит в очередь: подтверждение приходит колбэком,
    гарантию даёт flush(), он же возвращает число недоставленных. Здесь же —
    что ключ сообщения выбирает раздел;
  5. назначение на известный адрес вместо подписки: адрес известен, группа и
    перебалансировка не нужны (что на этом месте делает боевой потребитель,
    говорит шапка);
  6. None от poll() — норма, а не отказ; молчание брокера неотличимо от «ещё
    не приехало», поэтому у чтения обязан быть крайний срок.

Больше комментариев # в файле быть не должно.

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

  • test_kafka — две задачи, write_marker и read_marker; ни часовых
    = None, ни безымянных настроек клиентов в их телах не осталось.
  • Пишущая и читающая половины — отдельные задачи, каждая заводит и
    закрывает своего клиента.
  • Настройки клиентов и все сроки ожидания — именованные модульные
    константы; безымянных чисел в теле не осталось.
  • Проверки вынесены в функции, названные по способу отказа; ни одно условие
    не смешивает два разных отказа.
  • Адрес записи — NamedTuple уровня модуля с полями раздела и смещения;
    через XCom он едет именованными полями словаря, и читающая задача
    принимает этот словарь одним аргументом.
  • Шапка файла в форме «Тест проверяет: ...» плюс абзац, чем тест не
    является.
  • Комментарии стоят ровно в шести перечисленных местах.
  • Охват не сократился: пробник по-прежнему утверждает, что брокер подтвердил
    запись ровно одним сообщением и вернул её адрес; что чтение по этому
    адресу возвращает записанный маркер; что ошибка брокера при чтении и
    молчание дольше срока — отказы.
  • KafkaProbeTests в tests/dag-probes-unit.py переписан. Сейчас проверка
    держится за flush_timeouts == [10, 1], то есть за устройство finally,
    которого после правки не будет. На её место встают проверки свойств:
    продюсер закрыт даже тогда, когда запись отказала, и чтение назначается
    ровно на тот адрес, который вернул брокер.
  • Новых зависимостей нет.
  • make config-test зелен.
  • make smoke: оба пробника зелены. Проверка памяти стенда может быть
    красной по причине из #21 — это не отказ этого тикета.

Границы

  • Охват не расширять: тикет про форму, не про то, что пробник проверяет.
  • scripts/stand-smoke.sh, compose.yaml и состав проверок не трогать.
  • import confluent_kafka остаётся внутри функций: малые проверки грузят модуль
    обычным интерпретатором, где пакета нет, и импорт верхнего уровня красит
    make config-test.
  • Идентификатор check_round_trip уходит вместе с единственной задачей;
    малую проверку, которая достаёт задачу по имени, поправить тем же PR.
  • test_clickhouse не трогать: он приведён в порядок в #20.
  • ADR не заводить: решение о форме записано этим тикетом и показано кодом.

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

  1. AGENTS.md — разделы «Цель репозитория», «Кому что поручать», «Язык».
  2. dags/test_clickhouse.py целиком — образец формы.
  3. dags/test_kafka.py целиком.
  4. tests/dag-probes-unit.py — как малые проверки грузят модуль DAG с
    заглушенным airflow.sdk.
  5. issue #20 — правило комментариев и его оговорка про строки описания DAG.

Проверка

make config-test
make smoke

make smoke-guards не нужен: граф задач и тексты ошибок, по которым красный
путь опознаёт поломку, в этом тикете не меняются.

Пробник Kafka — второй и последний пример DAG в стенде, и после #20 единственный, который не показывает, как здесь пишут код. Он рабочий и проверяет ровно то, что должен: менять надо текст, а не охват. Тикет про то, как код читается, поэтому исполнитель — модель со вкусом (AGENTS.md, «Кому что поручать»). По той же причине решения о форме приняты здесь, а не оставлены на реализацию. ## Что не так сейчас - Тело задачи — 60 строк одним куском внутри одного `try`: настройка продюсера, доставка, сверка доставки, настройка консьюмера, цикл чтения, сверка маркера. На вопрос «что проверяет пробник» файл отвечает только целиком. - Часовые `producer = None` / `consumer = None` и `finally` с двумя проверками на `None` — четыре строки на то, чтобы закрыть то, что могло не открыться. - Проверки не названы: `if undelivered or delivery_errors or len(delivered_offsets) != 1` — три разных отказа в одном условии и с одним текстом ошибки. - Числа 10 000, 5 000, 10, 20 и 0.5 стоят без имён; что из них «сколько ждём доставки», а что «сколько ждём чтения», видно только из места. - Успех — `return` из середины цикла, отказ — после цикла. - Из уроков, которые файл мог бы дать про Kafka, в тексте видны два, и те по обрывкам. ## Как переписываем Форма — та же, что у `test_clickhouse` после #20: модульные константы, приватные помощники, отдельные утверждения, короткое тело задачи. Менти открывает два файла и видит один почерк — это и есть цель. 1. Две задачи: `write_marker` пишет маркер и возвращает адрес записи, `read_marker` читает по этому адресу и сверяет. Половины круга ломаются по-разному, и в интерфейсе Airflow должно быть видно, какая из них отказала; заодно на живом примере виден XCom. 2. Каждая задача заводит своего клиента и закрывает его сама: под LocalExecutor это два процесса, и одно подключение границу между ними не переживает. Продюсер и консьюмер — не два клиента одного вида, а два разных мира с разными настройками и разными способами сломаться; форма должна это показывать. 3. Настройки клиентов — именованные модульные константы; у каждого срока ожидания имя, говорящее, чего именно мы ждём. 4. Утверждения — функции с именами по способу отказа. 5. Адрес записи (раздел и смещение) собирается `NamedTuple` уровня модуля: два соседних целых в сигнатуре переставляются молча, и проверка после этого краснеет загадочно. Границу задач тип не переходит — XCom идёт через JSON, и кортеж вернулся бы списком. Между задачами адрес едет именованными полями словаря вместе с маркером. 6. Успех не возвращается из середины цикла. ## Шапка файла По образцу `test_clickhouse` (коммит c130319): «Тест проверяет:» списком того, что утверждается, и абзац о том, чем тест не является — не образец боевого потребителя: не подписывается на топик, не входит в группу, не коммитит смещения, не создаёт топик. Шапка уходит в `doc_md` и видна в интерфейсе Airflow. Правило комментариев из #20 строк описания DAG не касается. ## Комментарии Правило #20 остаётся в силе: комментарий объясняет решение, а не пересказывает строку. Список мест расширяется явно — комментируем там, где незнакома модель Kafka, и молчим там, где незнаком просто Python. Первое читатель из кода не выведет никогда, второе выведет за минуту. Шесть мест: 1. автосоздание топика (комментарий уже есть, у константы `TOPIC`); 2. одноразовая группа и выключенный автокоммит: проверка не двигает ничьих смещений и не оставляет следов на брокере; 3. три разных срока ожидания: сколько библиотека держит сообщение у себя, сколько ждём сеть, сколько ждём возврата; 4. `produce()` не пишет, а ставит в очередь: подтверждение приходит колбэком, гарантию даёт `flush()`, он же возвращает число недоставленных. Здесь же — что ключ сообщения выбирает раздел; 5. назначение на известный адрес вместо подписки: адрес известен, группа и перебалансировка не нужны (что на этом месте делает боевой потребитель, говорит шапка); 6. `None` от `poll()` — норма, а не отказ; молчание брокера неотличимо от «ещё не приехало», поэтому у чтения обязан быть крайний срок. Больше комментариев `#` в файле быть не должно. ## Критерии приёмки - [ ] `test_kafka` — две задачи, `write_marker` и `read_marker`; ни часовых `= None`, ни безымянных настроек клиентов в их телах не осталось. - [ ] Пишущая и читающая половины — отдельные задачи, каждая заводит и закрывает своего клиента. - [ ] Настройки клиентов и все сроки ожидания — именованные модульные константы; безымянных чисел в теле не осталось. - [ ] Проверки вынесены в функции, названные по способу отказа; ни одно условие не смешивает два разных отказа. - [ ] Адрес записи — `NamedTuple` уровня модуля с полями раздела и смещения; через XCom он едет именованными полями словаря, и читающая задача принимает этот словарь одним аргументом. - [ ] Шапка файла в форме «Тест проверяет: ...» плюс абзац, чем тест не является. - [ ] Комментарии стоят ровно в шести перечисленных местах. - [ ] Охват не сократился: пробник по-прежнему утверждает, что брокер подтвердил запись ровно одним сообщением и вернул её адрес; что чтение по этому адресу возвращает записанный маркер; что ошибка брокера при чтении и молчание дольше срока — отказы. - [ ] `KafkaProbeTests` в `tests/dag-probes-unit.py` переписан. Сейчас проверка держится за `flush_timeouts == [10, 1]`, то есть за устройство `finally`, которого после правки не будет. На её место встают проверки свойств: продюсер закрыт даже тогда, когда запись отказала, и чтение назначается ровно на тот адрес, который вернул брокер. - [ ] Новых зависимостей нет. - [ ] `make config-test` зелен. - [ ] `make smoke`: оба пробника зелены. Проверка памяти стенда может быть красной по причине из #21 — это не отказ этого тикета. ## Границы - Охват не расширять: тикет про форму, не про то, что пробник проверяет. - `scripts/stand-smoke.sh`, `compose.yaml` и состав проверок не трогать. - `import confluent_kafka` остаётся внутри функций: малые проверки грузят модуль обычным интерпретатором, где пакета нет, и импорт верхнего уровня красит `make config-test`. - Идентификатор `check_round_trip` уходит вместе с единственной задачей; малую проверку, которая достаёт задачу по имени, поправить тем же PR. - `test_clickhouse` не трогать: он приведён в порядок в #20. - ADR не заводить: решение о форме записано этим тикетом и показано кодом. ## Сначала прочитать 1. `AGENTS.md` — разделы «Цель репозитория», «Кому что поручать», «Язык». 2. `dags/test_clickhouse.py` целиком — образец формы. 3. `dags/test_kafka.py` целиком. 4. `tests/dag-probes-unit.py` — как малые проверки грузят модуль DAG с заглушенным `airflow.sdk`. 5. issue #20 — правило комментариев и его оговорка про строки описания DAG. ## Проверка make config-test make smoke `make smoke-guards` не нужен: граф задач и тексты ошибок, по которым красный путь опознаёт поломку, в этом тикете не меняются.
ddmitry added the ready-for-agent label 2026-07-31 19:21:01 +03:00
Author
Owner

Тело тикета поправлено при реализации: пробник режется на две задачи,
write_marker и read_marker, а не остаётся одной задачей.

Откуда взялась «одна задача». В гриллинге разбиение на две задачи предлагалось
и было названо лучшим вариантом. Дальше прозвучало «всякие xcom и прочее —
мишура», и это прочиталось как отказ от разбиения, хотя сказано было про XCom
как самоцель. Закрепила расхождение формула «режем по ролям: пишущая половина
и читающая»: в тексте «половина» означала модульный помощник, а читается как
задача — согласие получено на одно, записано другое.

Что поменялось в требованиях:

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

Всё остальное — форма, комментарии, шапка, охват утверждений — без изменений.

Тело тикета поправлено при реализации: пробник режется на две задачи, `write_marker` и `read_marker`, а не остаётся одной задачей. Откуда взялась «одна задача». В гриллинге разбиение на две задачи предлагалось и было названо лучшим вариантом. Дальше прозвучало «всякие xcom и прочее — мишура», и это прочиталось как отказ от разбиения, хотя сказано было про XCom как самоцель. Закрепила расхождение формула «режем по ролям: пишущая половина и читающая»: в тексте «половина» означала модульный помощник, а читается как задача — согласие получено на одно, записано другое. Что поменялось в требованиях: - две задачи вместо одной, идентификатор `check_round_trip` уходит; - `NamedTuple` границу задач не переходит: XCom идёт через JSON, кортеж вернулся бы списком, поэтому между задачами адрес едет именованными полями словаря вместе с маркером; - малая проверка «продюсер закрыт, когда консьюмер не поднялся» теряет смысл — консьюмер живёт в другом процессе. На её месте два свойства: продюсер закрыт и на пути отказа записи, и чтение назначается ровно на тот адрес, который вернул брокер. Всё остальное — форма, комментарии, шапка, охват утверждений — без изменений.
Sign in to join this conversation.