Форма пробных DAG: разбить на задачи, снять ручную машинерию, объяснить решения #20

Closed
opened 2026-07-31 15:39:37 +03:00 by ddmitry · 0 comments
Owner

Пробники из #18 работают и проверяют ровно то, что должны. Но они — первый и
пока единственный пример DAG в стенде: по ним будут учиться писать остальные,
а форма у них неудачная. Задача — привести форму в порядок, не трогая того,
что они проверяют.

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

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

  • test_clickhouse — одна задача на всё: снести таблицы, создать, проверить
    на обеих нодах, записать, прочитать локально, прочитать с ноды 2, сверить,
    убрать. В интерфейсе Airflow это один красный квадрат: где сломалось —
    иди читай журнал. Ради этого DAG и заводят вместо скрипта.
  • Уборка написана руками: проброс probe_error через except BaseException,
    сбор ошибок в список, ExceptionGroup, add_note — двадцать с лишним
    строк машинерии, через которую читатель продирается раньше, чем понимает,
    что вообще проверяет пробник. У Airflow для этого есть свой механизм.
  • Ноль комментариев на 215 строк — при том, что здесь лежит самое
    неочевидное решение репозитория: почему нода 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 каждая задача заводит себе сама через общий модульный
    помощник _clickhouse_client() и закрывает его в finally. Одно
    подключение границу процесса не переживает: под LocalExecutor четыре задачи
    — это четыре процесса.
  • Маркер строится в write_marker и едет дальше значением, которое задача
    возвращает: read_from_node_2 принимает его аргументом. Это XCom в своём
    обычном виде, и показать его на живом примере — часть смысла разбиения.
  • Через XCom едут только строки: маркер и имя ноды, принявшей запись. Строки
    результата запроса не передавать: XCom проходит через JSON, кортежи
    превращаются в списки, а сверка сравнивает кортежи. Сигнатуру
    _assert_marker_path под это менять можно, малую проверку обновить тем же
    PR.
  • Уборка ничего не ждёт от предыдущих задач: DROP ... IF EXISTS ON CLUSTER
    отрабатывает и тогда, когда prepare_tables не дошла до создания.

Комментарии

Комментарий объясняет решение, а не пересказывает строку. Он уместен там, где
читатель иначе спросит «почему так, а не очевидным способом», и неуместен
там, где код читается сам. В пробниках комментарий заслуживают ровно четыре
места:

  1. чтение ноды 2 через remote() с ноды 1;
  2. trigger_rule="all_done" у уборки и отказ от setup/teardown;
  3. импорт клиента внутри функции;
  4. автосоздание топика в 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. Ни булева, ни строкового переключателя между
    источниками не осталось.
  • Комментарий рядом с чтением ноды 2 объясняет три вещи: у Airflow
    подготовлено одно подключение — к 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_connectconfluent_kafka в пробнике Kafka)
    остаётся внутри функций: малые проверки грузят модуль обычным
    интерпретатором, где этих пакетов нет, и импорт верхнего уровня красит
    make config-test. Помощник _clickhouse_client() прячет импорт в себя,
    и каждой задаче достаётся одна строка. Причину записать комментарием —
    иначе следующий читатель поднимет импорт наверх.
  • test_kafka остаётся одной задачей check_round_trip: продюсер и
    консьюмер там связаны смещением сообщения, разрывать их по задачам значило
    бы гонять смещение через XCom ради шага, который и так виден по тексту
    ошибки. В этом пробнике меняется только место комментария.
  • Не расширять то, что пробники проверяют: тикет про форму, не про охват.

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

  1. AGENTS.md — разделы «Цель репозитория», «Кому что поручать», «Язык».
  2. dags/test_clickhouse.py и dags/test_kafka.py целиком.
  3. tests/dag-probes-unit.py — малые проверки логики пробников: они грузят
    модуль DAG с заглушенным airflow.sdk и падают от любого нового импорта
    верхнего уровня и от аргументов у @task.
  4. scripts/config-test.sh — как эти проверки запускаются в
    make config-test.
  5. tests/stand-smoke-guards.sh — красный путь: как ломается нода 2 и по
    какому журналу проверяется, что пробник назвал поломку.
  6. docs/adr/0001-stand-services.md — что о пробниках уже решено.

Проверка

make config-test
make smoke
make smoke-guards
Пробники из #18 работают и проверяют ровно то, что должны. Но они — первый и пока единственный пример DAG в стенде: по ним будут учиться писать остальные, а форма у них неудачная. Задача — привести форму в порядок, не трогая того, что они проверяют. Тикет про то, как код читается, поэтому исполнитель — модель со вкусом (`AGENTS.md`, раздел «Кому что поручать»). По той же причине решения о форме приняты здесь, а не оставлены на реализацию: тикет, который просит «убрать ручную машинерию» и не говорит, что ею считать, передаёт исполнителю ровно ту работу, ради которой заведён. ## Что не так сейчас - `test_clickhouse` — одна задача на всё: снести таблицы, создать, проверить на обеих нодах, записать, прочитать локально, прочитать с ноды 2, сверить, убрать. В интерфейсе Airflow это один красный квадрат: где сломалось — иди читай журнал. Ради этого DAG и заводят вместо скрипта. - Уборка написана руками: проброс `probe_error` через `except BaseException`, сбор ошибок в список, `ExceptionGroup`, `add_note` — двадцать с лишним строк машинерии, через которую читатель продирается раньше, чем понимает, что вообще проверяет пробник. У Airflow для этого есть свой механизм. - Ноль комментариев на 215 строк — при том, что здесь лежит самое неочевидное решение репозитория: почему нода 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 каждая задача заводит себе сама через общий модульный помощник `_clickhouse_client()` и закрывает его в `finally`. Одно подключение границу процесса не переживает: под LocalExecutor четыре задачи — это четыре процесса. - Маркер строится в `write_marker` и едет дальше значением, которое задача возвращает: `read_from_node_2` принимает его аргументом. Это XCom в своём обычном виде, и показать его на живом примере — часть смысла разбиения. - Через XCom едут только строки: маркер и имя ноды, принявшей запись. Строки результата запроса не передавать: XCom проходит через JSON, кортежи превращаются в списки, а сверка сравнивает кортежи. Сигнатуру `_assert_marker_path` под это менять можно, малую проверку обновить тем же PR. - Уборка ничего не ждёт от предыдущих задач: `DROP ... IF EXISTS ON CLUSTER` отрабатывает и тогда, когда `prepare_tables` не дошла до создания. ## Комментарии Комментарий объясняет решение, а не пересказывает строку. Он уместен там, где читатель иначе спросит «почему так, а не очевидным способом», и неуместен там, где код читается сам. В пробниках комментарий заслуживают ровно четыре места: 1. чтение ноды 2 через `remote()` с ноды 1; 2. `trigger_rule="all_done"` у уборки и отказ от setup/teardown; 3. импорт клиента внутри функции; 4. автосоздание топика в `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`. Ни булева, ни строкового переключателя между источниками не осталось. - [ ] Комментарий рядом с чтением ноды 2 объясняет три вещи: у Airflow подготовлено одно подключение — к `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 ради шага, который и так виден по тексту ошибки. В этом пробнике меняется только место комментария. - Не расширять то, что пробники проверяют: тикет про форму, не про охват. ## Сначала прочитать 1. `AGENTS.md` — разделы «Цель репозитория», «Кому что поручать», «Язык». 2. `dags/test_clickhouse.py` и `dags/test_kafka.py` целиком. 3. `tests/dag-probes-unit.py` — малые проверки логики пробников: они грузят модуль DAG с заглушенным `airflow.sdk` и падают от любого нового импорта верхнего уровня и от аргументов у `@task`. 4. `scripts/config-test.sh` — как эти проверки запускаются в `make config-test`. 5. `tests/stand-smoke-guards.sh` — красный путь: как ломается нода 2 и по какому журналу проверяется, что пробник назвал поломку. 6. `docs/adr/0001-stand-services.md` — что о пробниках уже решено. ## Проверка ``` make config-test make smoke make smoke-guards ```
ddmitry added the ready-for-agent label 2026-07-31 17:03:03 +03:00
Sign in to join this conversation.
1 Participants
Notifications
Due Date
No due date set.
Reference: ddmitry/clickstream-data-platform#20