Форма пробника Kafka: переписать по образцу test_clickhouse #22
Notifications
Due Date
No due date set.
Depends on
#20 Форма пробных DAG: разбить на задачи, снять ручную машинерию, объяснить решения
ddmitry/clickstream-data-platform
Reference: ddmitry/clickstream-data-platform#22
Reference in New Issue
Block a user
Пробник Kafka — второй и последний пример DAG в стенде, и после #20 единственный,
который не показывает, как здесь пишут код. Он рабочий и проверяет ровно то, что
должен: менять надо текст, а не охват.
Тикет про то, как код читается, поэтому исполнитель — модель со вкусом
(AGENTS.md, «Кому что поручать»). По той же причине решения о форме приняты
здесь, а не оставлены на реализацию.
Что не так сейчас
try: настройка продюсера,доставка, сверка доставки, настройка консьюмера, цикл чтения, сверка маркера.
На вопрос «что проверяет пробник» файл отвечает только целиком.
producer = None/consumer = Noneиfinallyс двумя проверкамина
None— четыре строки на то, чтобы закрыть то, что могло не открыться.if undelivered or delivery_errors or len(delivered_offsets) != 1— три разных отказа в одном условии и с однимтекстом ошибки.
доставки», а что «сколько ждём чтения», видно только из места.
returnиз середины цикла, отказ — после цикла.обрывкам.
Как переписываем
Форма — та же, что у
test_clickhouseпосле #20: модульные константы, приватныепомощники, отдельные утверждения, короткое тело задачи. Менти открывает два
файла и видит один почерк — это и есть цель.
write_markerпишет маркер и возвращает адрес записи,read_markerчитает по этому адресу и сверяет. Половины круга ломаютсяпо-разному, и в интерфейсе Airflow должно быть видно, какая из них отказала;
заодно на живом примере виден XCom.
LocalExecutor это два процесса, и одно подключение границу между ними не
переживает. Продюсер и консьюмер — не два клиента одного вида, а два разных
мира с разными настройками и разными способами сломаться; форма должна это
показывать.
ожидания имя, говорящее, чего именно мы ждём.
NamedTupleуровня модуля:два соседних целых в сигнатуре переставляются молча, и проверка после этого
краснеет загадочно. Границу задач тип не переходит — XCom идёт через JSON,
и кортеж вернулся бы списком. Между задачами адрес едет именованными полями
словаря вместе с маркером.
Шапка файла
По образцу
test_clickhouse(коммитc130319): «Тест проверяет:» списком того,что утверждается, и абзац о том, чем тест не является — не образец боевого
потребителя: не подписывается на топик, не входит в группу, не коммитит
смещения, не создаёт топик. Шапка уходит в
doc_mdи видна в интерфейсеAirflow. Правило комментариев из #20 строк описания DAG не касается.
Комментарии
Правило #20 остаётся в силе: комментарий объясняет решение, а не пересказывает
строку. Список мест расширяется явно — комментируем там, где незнакома модель
Kafka, и молчим там, где незнаком просто Python. Первое читатель из кода не
выведет никогда, второе выведет за минуту. Шесть мест:
TOPIC);смещений и не оставляет следов на брокере;
сколько ждём сеть, сколько ждём возврата;
produce()не пишет, а ставит в очередь: подтверждение приходит колбэком,гарантию даёт
flush(), он же возвращает число недоставленных. Здесь же —что ключ сообщения выбирает раздел;
перебалансировка не нужны (что на этом месте делает боевой потребитель,
говорит шапка);
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.Сначала прочитать
AGENTS.md— разделы «Цель репозитория», «Кому что поручать», «Язык».dags/test_clickhouse.pyцеликом — образец формы.dags/test_kafka.pyцеликом.tests/dag-probes-unit.py— как малые проверки грузят модуль DAG сзаглушенным
airflow.sdk.Проверка
make smoke-guardsне нужен: граф задач и тексты ошибок, по которым красныйпуть опознаёт поломку, в этом тикете не меняются.
Тело тикета поправлено при реализации: пробник режется на две задачи,
write_markerиread_marker, а не остаётся одной задачей.Откуда взялась «одна задача». В гриллинге разбиение на две задачи предлагалось
и было названо лучшим вариантом. Дальше прозвучало «всякие xcom и прочее —
мишура», и это прочиталось как отказ от разбиения, хотя сказано было про XCom
как самоцель. Закрепила расхождение формула «режем по ролям: пишущая половина
и читающая»: в тексте «половина» означала модульный помощник, а читается как
задача — согласие получено на одно, записано другое.
Что поменялось в требованиях:
check_round_tripуходит;NamedTupleграницу задач не переходит: XCom идёт через JSON, кортежвернулся бы списком, поэтому между задачами адрес едет именованными полями
словаря вместе с маркером;
консьюмер живёт в другом процессе. На её месте два свойства: продюсер закрыт
и на пути отказа записи, и чтение назначается ровно на тот адрес, который
вернул брокер.
Всё остальное — форма, комментарии, шапка, охват утверждений — без изменений.