refactor(airflow): пробник Kafka разрезан на запись и чтение #24

Merged
ddmitry merged 1 commits from feat/22-kafka-probe-form into main 2026-07-31 20:25:54 +03:00
Owner

Пробник Kafka приведён к форме test_clickhouse из #20 и разрезан на две
задачи. Менти открывает два файла и видит один почерк — это и есть цель.

Что меняется

  • Вместо одной задачи check_round_trip две: write_marker пишет маркер и
    возвращает адрес записи, read_marker читает по этому адресу и сверяет.
    В интерфейсе Airflow теперь видно, какая половина круга отказала: брокер не
    принял маркер или не отдал его обратно.
  • Адрес и маркер едут между задачами XCom — именованными полями словаря.
    Внутри записи адрес собирается NamedTuple RecordAddress: два соседних
    целых в сигнатуре переставляются молча. Границу задач тип не переходит —
    XCom идёт через JSON, и кортеж вернулся бы списком.
  • Каждая задача заводит своего клиента и закрывает его сама: часовые = None
    и finally с проверками на None ушли. У продюсера закрыть за собой — это
    flush(): своего close() у него нет, и он же возвращает число
    недоставленных.
  • Настройки клиентов и все сроки ожидания — именованные модульные константы.
  • Слитное условие доставки разобрано на шесть утверждений, каждое со своим
    именем и своим текстом ошибки.
  • Успех не возвращается из середины цикла: чтение выходит по сообщению или по
    крайнему сроку, сверка идёт после.
  • Шапка файла — «Тест проверяет: ...» плюс абзац о том, чем тест не является;
    она же уходит в doc_md. Комментарии стоят ровно в шести местах, где
    незнакома модель Kafka.

Охват не менялся: пробник утверждает ровно то же, что и раньше.

Тело тикета поправлено тем же PR

В #22 было записано «задача остаётся одна» — это расхождение с тем, о чём
договаривались в гриллинге. Откуда оно взялось, записано комментарием в
тикете.

Отдельное решение

Из настроек консьюмера убран session.timeout.ms: он про членство в группе и
удары сердца координатору, а пробник назначает себе адрес и в группу не
входит. Поведение не меняется; почему настройки нет, объясняет шапка файла.
Сверено по документации confluent-kafka-python через Context7.

Проверка

  • make config-test — зелен.
  • make smoke — оба пробника зелены. Проверка памяти стенда колеблется вокруг
    порога 3,4 ГБ и бывает красной по причине из #21; к этому PR отношения не
    имеет.
  • make smoke-guards не гонялся: граф задач стал шире, но тексты ошибок, по
    которым красный путь опознаёт поломку, не менялись, а сам путь работает на
    уровне DAG.

Closes #22

Пробник Kafka приведён к форме `test_clickhouse` из #20 и разрезан на две задачи. Менти открывает два файла и видит один почерк — это и есть цель. ## Что меняется - Вместо одной задачи `check_round_trip` две: `write_marker` пишет маркер и возвращает адрес записи, `read_marker` читает по этому адресу и сверяет. В интерфейсе Airflow теперь видно, какая половина круга отказала: брокер не принял маркер или не отдал его обратно. - Адрес и маркер едут между задачами XCom — именованными полями словаря. Внутри записи адрес собирается `NamedTuple` `RecordAddress`: два соседних целых в сигнатуре переставляются молча. Границу задач тип не переходит — XCom идёт через JSON, и кортеж вернулся бы списком. - Каждая задача заводит своего клиента и закрывает его сама: часовые `= None` и `finally` с проверками на `None` ушли. У продюсера закрыть за собой — это `flush()`: своего `close()` у него нет, и он же возвращает число недоставленных. - Настройки клиентов и все сроки ожидания — именованные модульные константы. - Слитное условие доставки разобрано на шесть утверждений, каждое со своим именем и своим текстом ошибки. - Успех не возвращается из середины цикла: чтение выходит по сообщению или по крайнему сроку, сверка идёт после. - Шапка файла — «Тест проверяет: ...» плюс абзац о том, чем тест не является; она же уходит в `doc_md`. Комментарии стоят ровно в шести местах, где незнакома модель Kafka. Охват не менялся: пробник утверждает ровно то же, что и раньше. ## Тело тикета поправлено тем же PR В #22 было записано «задача остаётся одна» — это расхождение с тем, о чём договаривались в гриллинге. Откуда оно взялось, записано комментарием в тикете. ## Отдельное решение Из настроек консьюмера убран `session.timeout.ms`: он про членство в группе и удары сердца координатору, а пробник назначает себе адрес и в группу не входит. Поведение не меняется; почему настройки нет, объясняет шапка файла. Сверено по документации `confluent-kafka-python` через Context7. ## Проверка - `make config-test` — зелен. - `make smoke` — оба пробника зелены. Проверка памяти стенда колеблется вокруг порога 3,4 ГБ и бывает красной по причине из #21; к этому PR отношения не имеет. - `make smoke-guards` не гонялся: граф задач стал шире, но тексты ошибок, по которым красный путь опознаёт поломку, не менялись, а сам путь работает на уровне DAG. Closes #22
ddmitry added 1 commit 2026-07-31 20:21:18 +03:00
Зачем: пробник Kafka — второй и последний пример DAG в стенде, и после #20
единственный, который не показывает, как здесь пишут код. Одним красным
квадратом он к тому же не отвечал на вопрос, какая половина круга отказала:
брокер не принял маркер или не отдал его обратно.

Что:
- вместо одной задачи check_round_trip две: write_marker пишет маркер и
  возвращает адрес записи, read_marker читает по этому адресу и сверяет;
- адрес и маркер едут между задачами XCom — именованными полями словаря:
  XCom проходит через JSON, и кортеж вернулся бы списком;
- внутри записи адрес собирается NamedTuple RecordAddress — два соседних
  целых в сигнатуре переставляются молча;
- каждая задача заводит своего клиента и закрывает его сама, поэтому часовые
  = None и finally с проверками на None ушли; у продюсера закрыть за собой —
  это flush(): своего close() у него нет, и он же возвращает число
  недоставленных;
- настройки клиентов и все сроки ожидания стали именованными модульными
  константами; безымянных чисел в телах задач не осталось;
- слитное условие доставки разобрано на шесть утверждений, каждое со своим
  именем и своим текстом ошибки;
- успех больше не возвращается из середины цикла: чтение выходит из цикла по
  сообщению или по крайнему сроку, а сверка идёт после;
- шапка файла приведена к форме «Тест проверяет: ...» с абзацем о том, чем
  тест не является; она же уходит в doc_md;
- комментарии стоят ровно в шести местах, где незнакома модель Kafka.

Из настроек консьюмера убран session.timeout.ms: он про членство в группе и
удары сердца координатору, а пробник назначает себе адрес и в группу не
входит — почему его нет, объясняет шапка файла. Поведение не меняется.
Сверено по документации confluent-kafka-python через Context7:
session.timeout.ms описан как срок сессии группы, flush() возвращает число
оставшихся в очереди сообщений.

KafkaProbeTests держался за flush_timeouts == [10, 1], то есть за устройство
finally, которого больше нет. На его место встали две проверки свойств:
продюсер закрыт даже тогда, когда отказала запись, и чтение назначается ровно
на тот адрес, который вернул брокер.

Тело issue #22 поправлено тем же изменением: там было записано «задача
остаётся одна» — это расхождение с тем, о чём договаривались в гриллинге.

Проверка: make config-test зелен; make smoke — оба пробника зелены, красной
осталась только проверка памяти стенда по причине из #21.

Closes #22
ddmitry merged commit d0c9e8ff43 into main 2026-07-31 20:25:54 +03:00
ddmitry deleted branch feat/22-kafka-probe-form 2026-07-31 20:25:54 +03:00
Sign in to join this conversation.
No Reviewers
1 Participants
Notifications
Due Date
No due date set.
Dependencies

No dependencies set.

Reference: ddmitry/clickstream-data-platform#24