Форма пробных DAG: разбить на задачи, снять ручную машинерию, объяснить решения #20
Notifications
Due Date
No due date set.
Blocks
#22 Форма пробника Kafka: переписать по образцу test_clickhouse
ddmitry/clickstream-data-platform
Reference: ddmitry/clickstream-data-platform#20
Reference in New Issue
Block a user
Пробники из #18 работают и проверяют ровно то, что должны. Но они — первый и
пока единственный пример DAG в стенде: по ним будут учиться писать остальные,
а форма у них неудачная. Задача — привести форму в порядок, не трогая того,
что они проверяют.
Тикет про то, как код читается, поэтому исполнитель — модель со вкусом
(
AGENTS.md, раздел «Кому что поручать»). По той же причине решения о формеприняты здесь, а не оставлены на реализацию: тикет, который просит «убрать
ручную машинерию» и не говорит, что ею считать, передаёт исполнителю ровно ту
работу, ради которой заведён.
Что не так сейчас
test_clickhouse— одна задача на всё: снести таблицы, создать, проверитьна обеих нодах, записать, прочитать локально, прочитать с ноды 2, сверить,
убрать. В интерфейсе Airflow это один красный квадрат: где сломалось —
иди читай журнал. Ради этого DAG и заводят вместо скрипта.
probe_errorчерезexcept BaseException,сбор ошибок в список,
ExceptionGroup,add_note— двадцать с лишнимстрок машинерии, через которую читатель продирается раньше, чем понимает,
что вообще проверяет пробник. У Airflow для этого есть свой механизм.
неочевидное решение репозитория: почему нода 2 читается через
remote('clickhouse-02:9000', ...)с ноды 1, а не прямым подключением._table_engines(client, *, query_node_2: bool)переключается булевымфлагом между источниками, а вызывающий тащит рядом параллельный список
подписей
((False, "ноде 1"), (True, "ноде 2")).test_kafkaстоит внутриtryнадсписком ошибок доставки — то есть там, где закрывался критерий, а не там,
где он помогает читателю.
Как разбиваем
test_clickhouseЧетыре задачи, выстроенные в цепочку в этом порядке:
prepare_tables— снести остатки прошлого запуска и убедиться, что их нетна обеих нодах; создать локальную и распределённую таблицы
ON CLUSTER;проверить, что на обеих нодах лежит ожидаемый набор движков.
write_marker— вставить маркер в локальную таблицу ноды 1, прочитать егооттуда же, вернуть маркер и имя ноды, принявшей запись.
read_from_node_2— определить имя ноды 2, прочитатьDistributedс ноды 2через
remote(), сверить путь маркера.cleanup_tables— снести таблицы и убедиться, что их нет на обеих нодах.Как задачи делятся состоянием
Это самое сложное в разбиении, поэтому решено здесь.
помощник
_clickhouse_client()и закрывает его вfinally. Одноподключение границу процесса не переживает: под LocalExecutor четыре задачи
— это четыре процесса.
write_markerи едет дальше значением, которое задачавозвращает:
read_from_node_2принимает его аргументом. Это XCom в своёмобычном виде, и показать его на живом примере — часть смысла разбиения.
результата запроса не передавать: XCom проходит через JSON, кортежи
превращаются в списки, а сверка сравнивает кортежи. Сигнатуру
_assert_marker_pathпод это менять можно, малую проверку обновить тем жеPR.
DROP ... IF EXISTS ON CLUSTERотрабатывает и тогда, когда
prepare_tablesне дошла до создания.Комментарии
Комментарий объясняет решение, а не пересказывает строку. Он уместен там, где
читатель иначе спросит «почему так, а не очевидным способом», и неуместен
там, где код читается сам. В пробниках комментарий заслуживают ровно четыре
места:
remote()с ноды 1;trigger_rule="all_done"у уборки и отказ от setup/teardown;test_kafka.Больше комментариев
#в пробниках быть не должно. Если тянет пояснитьчто-то ещё — это признак, что неудачно выбрано имя задачи или функции.
Строки описания DAG это правило не касается.
Критерии приёмки
test_clickhouseсостоит ровно из четырёх задач с идентификаторамиprepare_tables,write_marker,read_from_node_2,cleanup_tablesи содержанием из раздела «Как разбиваем».
trigger_rule="all_done". Неsetup/teardown: у задач-teardown отказ по умолчанию не влияет на
состояние запуска (
on_failure_fail_dagrun=False), а стенду нужнообратное — не убравшийся за собой пробник обязан краснеть, на этом
стоит красный путь
make smoke-guards.шагами: ни
except BaseException, ни накопления ошибок в список, ниExceptionGroup, ниadd_note, ни склейки текстов ошибок. Каждаязадача падает своим исключением на своей строке; связывает шаги
Airflow.
try/finallyдопустим только для закрытия клиента внутриодной задачи.
источник для запроса», например:
python NODES = ( ("ноде 1", "system.tables"), ("ноде 2", "remote('clickhouse-02:9000', system.tables)"), )_table_enginesпринимает готовый источник строкой и ничего невыбирает; оба места, где раньше стоял параллельный список подписей,
ходят по
NODES. Ни булева, ни строкового переключателя междуисточниками не осталось.
подготовлено одно подключение — к
clickhouse-01, и второго радипробника не заводят;
remote('clickhouse-02:9000', ...)делаетинициатором распределённого запроса саму ноду 2 — именно это и
проверяется, а не доступность ноды 2 по сети; порт
9000—межсерверный, тогда как подключение Airflow ходит по HTTP на
8123.test_kafkaстоит рядом сконстантой
TOPIC.утверждает всё перечисленное:
1. до создания служебных таблиц нет ни на ноде 1, ни на ноде 2;
2. после создания на обеих нодах лежит ровно пара
airflow_probe_distributed/Distributedиairflow_probe_local/ReplicatedMergeTree;3. маркер уникален для запуска, и в локальной таблице находится ровно
одна строка с ним;
4. имя ноды 2 определено (ровно одна строка ответа) и отличается от
ноды, принявшей запись;
5. чтение
Distributedс ноды 2 возвращает ровно одну строку с_shard_num = 1;6. после уборки таблиц нет ни на одной ноде.
Сообщение об ошибке для пункта 2 сохраняет текст «неверный набор
таблиц на ноде 2»: по нему опознаёт поломку
tests/stand-smoke-guards.sh.tests/dag-probes-unit.pyпоправлен тем же PR. Заглушкаtaskдолжнапринимать обе формы —
@taskи@task(trigger_rule=...); сейчас онаобъявлена как
def task(function)и от аргументов декоратора падает,унося с собой
make config-test. Годится форма:python def task(function=None, **_kwargs): def declare_task(*_args, **_kwargs): return None return declare_task if function is not None else lambda _f: declare_taskПравится только заглушка в
load_clickhouse_dag_module: аргументыдекоратора получает она одна. Заглушку в
load_kafka_dag_taskнетрогать — она запоминает функцию по имени
(
captured_tasks[function.__name__]), и приведённая форма эту памятьтеряет:
captured_tasks["check_round_trip"]упадёт сKeyError.test_cleanup_keeps_original_error_as_primaryснята вместе с_cleanup_clickhouse_client, и на её место встала проверка того, чтоуборка сносит таблицы и закрывает клиент даже после падения
предыдущих шагов. Свойство «ошибка шага остаётся главной» теперь
обеспечивает Airflow, а не наш код, и проверяется красным путём
make smoke-guards._assert_marker_pathостаётся модульной функциейи по-прежнему покрыт малой проверкой.
tests/stand-smoke-guards.shпоправлен тем же PR. Сейчасclickhouse_break_is_reportedберёт один самый свежийattempt=1.log;после разбиения самым свежим станет журнал уборки, и образец в нём не
найдётся. Проверка должна искать по журналам всех задач запуска: брать
самый свежий каталог
run_id=*и грепать всеattempt=1.logвнутри.Состав проверок не меняется — меняется способ найти журнал.
trigger_rule="all_done"сверен по документации Airflow 3.3.0через Context7 и записан комментарием в коде: что проверили и почему
выбран этот вариант.
make config-testзелен: малые проверки грузят модуль DAG с@task(trigger_rule="all_done")и печатают «ошибок 0».make smokeзелен, 25 проверок.make smoke-guardsзелен: красный путь по-прежнему связывает отказtest_clickhouseс удалённой на ноде 2 таблицей — по журналу тойзадачи, которая упала.
Что сообщить в отчёте
Сколько заняла разбитая версия
test_clickhouse. Предел ожидания вrun_airflow_probe— 60 опросов по 2 секунды (scripts/stand-smoke.sh);четыре задачи вместо одной этот запас едят. Если запаса осталось меньше
двукратного — вернуться с числом; менять предел самостоятельно не надо.
Границы
infra/airflow/Dockerfile,compose.yamlи состав проверок.scripts/stand-smoke.shне трогать: он знает только идентификаторы DAG(
test_clickhouse,test_kafka), а они не меняются. Имена задач наружуутекают в одном месте —
tests/dag-probes-unit.pyберёт задачу Kafka поимени
check_round_trip; если оно меняется, правится и там.import clickhouse_connect(иconfluent_kafkaв пробнике Kafka)остаётся внутри функций: малые проверки грузят модуль обычным
интерпретатором, где этих пакетов нет, и импорт верхнего уровня красит
make config-test. Помощник_clickhouse_client()прячет импорт в себя,и каждой задаче достаётся одна строка. Причину записать комментарием —
иначе следующий читатель поднимет импорт наверх.
test_kafkaостаётся одной задачейcheck_round_trip: продюсер иконсьюмер там связаны смещением сообщения, разрывать их по задачам значило
бы гонять смещение через XCom ради шага, который и так виден по тексту
ошибки. В этом пробнике меняется только место комментария.
Сначала прочитать
AGENTS.md— разделы «Цель репозитория», «Кому что поручать», «Язык».dags/test_clickhouse.pyиdags/test_kafka.pyцеликом.tests/dag-probes-unit.py— малые проверки логики пробников: они грузятмодуль DAG с заглушенным
airflow.sdkи падают от любого нового импортаверхнего уровня и от аргументов у
@task.scripts/config-test.sh— как эти проверки запускаются вmake config-test.tests/stand-smoke-guards.sh— красный путь: как ломается нода 2 и покакому журналу проверяется, что пробник назвал поломку.
docs/adr/0001-stand-services.md— что о пробниках уже решено.Проверка