From dc12953cac61165e18c80f8c9ae64d5905a078bd Mon Sep 17 00:00:00 2001 From: Dmitry Dementiev Date: Fri, 31 Jul 2026 18:26:39 +0300 Subject: [PATCH 1/2] =?UTF-8?q?refactor(airflow):=20=D0=BF=D1=80=D0=BE?= =?UTF-8?q?=D0=B1=D0=BD=D0=B8=D0=BA=20ClickHouse=20=D1=80=D0=B0=D0=B7?= =?UTF-8?q?=D0=B1=D0=B8=D1=82=20=D0=BD=D0=B0=20=D1=87=D0=B5=D1=82=D1=8B?= =?UTF-8?q?=D1=80=D0=B5=20=D0=B7=D0=B0=D0=B4=D0=B0=D1=87=D0=B8?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Зачем: пробник был одной задачей — в интерфейсе Airflow один красный квадрат, а место отказа приходилось искать по журналу. Ручная машинерия проброса и сведения ошибок занимала больше места, чем сама проверка, и читатель продирался через неё раньше, чем понимал, что пробник проверяет. Пробники — единственный образец DAG в стенде, по ним будут писать остальные. Что: test_clickhouse разбит на prepare_tables, write_marker, read_from_node_2 и cleanup_tables; маркер и имя принявшей запись ноды едут между задачами через XCom строками. Снято сведение ошибок: except BaseException, ExceptionGroup, add_note и накопление ошибок в список; клиент каждая задача заводит общим помощником и закрывает в finally. Ноды описаны константой NODES парами «имя для человека — источник для запроса», булев переключатель и параллельные списки подписей ушли. Уборка идёт обычным правилом запуска, а не all_done: состояние запуска Airflow считает по концам графа, и уборка, отработавшая после отказа, покрасила бы в зелёный запуск с упавшей проверкой — решение записано в ADR 0003. Комментарии остались в четырёх местах: чтение ноды 2 через remote(), импорт клиента внутри функции, правило запуска уборки и автосоздание топика в test_kafka. Малые проверки: заглушка task принимает обе формы декоратора, проверка сведения ошибок заменена проверками уборки. Красный путь ищет образец по журналам всех задач последнего запуска, а не в одном самом свежем. Проверка: make config-test, make smoke (25 проверок) и make smoke-guards зелены. Разбитый пробник укладывается в 5 секунд из 120, отведённых run_airflow_probe, — предел не трогаем. Co-Authored-By: Claude Opus 5 (1M context) --- dags/test_clickhouse.py | 165 ++++++++++++++++++--------------- dags/test_kafka.py | 4 +- docs/adr/0003-dag-run-state.md | 55 +++++++++++ tests/dag-probes-unit.py | 93 +++++++++++++++---- tests/stand-smoke-guards.sh | 7 +- 5 files changed, 225 insertions(+), 99 deletions(-) create mode 100644 docs/adr/0003-dag-run-state.md diff --git a/dags/test_clickhouse.py b/dags/test_clickhouse.py index fa34817..7a0c630 100644 --- a/dags/test_clickhouse.py +++ b/dags/test_clickhouse.py @@ -10,14 +10,42 @@ from airflow.sdk import Connection, dag, get_current_context, task CLUSTER = "clickstream_cluster" LOCAL_TABLE = "airflow_probe_local" DISTRIBUTED_TABLE = "airflow_probe_distributed" +EXPECTED_TABLES = [ + (DISTRIBUTED_TABLE, "Distributed"), + (LOCAL_TABLE, "ReplicatedMergeTree"), +] + +# Ноду 2 пробник читает не своим подключением, а запросом remote() с ноды 1: у +# Airflow подготовлено одно подключение — к clickhouse-01, и второго ради +# пробника не заводят. При этом remote('clickhouse-02:9000', ...) делает +# инициатором распределённого запроса саму ноду 2 — проверяется именно это, а +# не доступность ноды 2 по сети. Порт 9000 — межсерверный, тогда как +# подключение Airflow ходит по HTTP на 8123. +NODES = ( + ("ноде 1", "system.tables"), + ("ноде 2", "remote('clickhouse-02:9000', system.tables)"), +) -def _table_engines(client, *, query_node_2: bool) -> list[tuple[str, str]]: - source = ( - "remote('clickhouse-02:9000', system.tables)" - if query_node_2 - else "system.tables" +def _clickhouse_client(): + # clickhouse_connect стоит только в образе Airflow, а малые проверки грузят + # этот модуль обычным интерпретатором, где пакета нет. Импорт верхнего + # уровня красит make config-test, поэтому он живёт здесь. + import clickhouse_connect + + connection = Connection.get("clickhouse_default") + return clickhouse_connect.get_client( + host=connection.host, + port=connection.port, + username=connection.login or "default", + password=connection.password or "", + database=connection.schema or "default", + connect_timeout=5, + send_receive_timeout=30, ) + + +def _table_engines(client, source: str) -> list[tuple[str, str]]: result = client.query( f""" SELECT name, engine @@ -41,27 +69,33 @@ def _drop_tables(client) -> None: def _assert_tables_absent(client) -> None: - for query_node_2, node_name in ((False, "ноде 1"), (True, "ноде 2")): - remaining = _table_engines(client, query_node_2=query_node_2) + for node_name, source in NODES: + remaining = _table_engines(client, source) if remaining: raise RuntimeError( f"служебные таблицы остались на {node_name}: {remaining}" ) +def _assert_tables_created(client) -> None: + for node_name, source in NODES: + actual_tables = _table_engines(client, source) + if actual_tables != EXPECTED_TABLES: + raise RuntimeError( + f"неверный набор таблиц на {node_name}: {actual_tables}" + ) + + def _assert_marker_path( *, - local_rows: list[tuple[str, str]], distributed_rows: list[tuple[int, str, str]], + write_hostname: str, node_2_hostname: str, marker: str, ) -> None: - if len(local_rows) != 1 or local_rows[0][1] != marker: - raise RuntimeError(f"маркер не найден в локальной таблице: {marker}") - node_1_hostname = local_rows[0][0] - if node_1_hostname == node_2_hostname: + if write_hostname == node_2_hostname: raise RuntimeError("запись и чтение маркера должны выполняться с разных нод") - expected_rows = [(1, node_1_hostname, marker)] + expected_rows = [(1, write_hostname, marker)] if distributed_rows != expected_rows: raise RuntimeError( "нода 2 не прочитала маркер первого шарда через Distributed: " @@ -69,30 +103,6 @@ def _assert_marker_path( ) -def _cleanup_clickhouse_client(client, original_error: BaseException | None) -> None: - cleanup_errors: list[Exception] = [] - try: - _drop_tables(client) - _assert_tables_absent(client) - except Exception as error: - cleanup_errors.append(error) - try: - client.close() - except Exception as error: - cleanup_errors.append(error) - - if original_error is not None: - for error in cleanup_errors: - original_error.add_note( - f"Дополнительная ошибка очистки ClickHouse: {error}" - ) - return - if len(cleanup_errors) == 1: - raise cleanup_errors[0] - if cleanup_errors: - raise ExceptionGroup("ошибки очистки ClickHouse", cleanup_errors) - - @dag( dag_id="test_clickhouse", schedule=None, @@ -103,26 +113,8 @@ def _cleanup_clickhouse_client(client, original_error: BaseException | None) -> ) def test_clickhouse(): @task - def check_cluster_path() -> None: - import clickhouse_connect - - connection = Connection.get("clickhouse_default") - client = clickhouse_connect.get_client( - host=connection.host, - port=connection.port, - username=connection.login or "default", - password=connection.password or "", - database=connection.schema or "default", - connect_timeout=5, - send_receive_timeout=30, - ) - marker = f"{get_current_context()['run_id']}:{uuid.uuid4()}" - expected_tables = [ - (DISTRIBUTED_TABLE, "Distributed"), - (LOCAL_TABLE, "ReplicatedMergeTree"), - ] - probe_error = None - + def prepare_tables() -> None: + client = _clickhouse_client() try: _drop_tables(client) _assert_tables_absent(client) @@ -151,17 +143,15 @@ def test_clickhouse(): ) """ ) + _assert_tables_created(client) + finally: + client.close() - for query_node_2, node_name in ((False, "ноде 1"), (True, "ноде 2")): - actual_tables = _table_engines( - client, - query_node_2=query_node_2, - ) - if actual_tables != expected_tables: - raise RuntimeError( - f"неверный набор таблиц на {node_name}: {actual_tables}" - ) - + @task + def write_marker() -> dict[str, str]: + client = _clickhouse_client() + try: + marker = f"{get_current_context()['run_id']}:{uuid.uuid4()}" client.insert( f"default.{LOCAL_TABLE}", [[marker]], @@ -175,6 +165,16 @@ def test_clickhouse(): """, parameters={"marker": marker}, ).result_rows + if len(local_rows) != 1 or local_rows[0][1] != marker: + raise RuntimeError(f"маркер не найден в локальной таблице: {marker}") + return {"marker": marker, "hostname": local_rows[0][0]} + finally: + client.close() + + @task + def read_from_node_2(written: dict[str, str]) -> None: + client = _clickhouse_client() + try: node_2_rows = client.query( """ SELECT hostName() @@ -195,21 +195,36 @@ def test_clickhouse(): ) WHERE marker = {{marker:String}} """, - parameters={"marker": marker}, + parameters={"marker": written["marker"]}, ).result_rows _assert_marker_path( - local_rows=local_rows, distributed_rows=distributed_rows, + write_hostname=written["hostname"], node_2_hostname=node_2_rows[0][0], - marker=marker, + marker=written["marker"], ) - except BaseException as error: - probe_error = error - raise finally: - _cleanup_clickhouse_client(client, probe_error) + client.close() - check_cluster_path() + # Уборка идёт только после успеха: упавший пробник оставляет кластер таким, + # каким сломался, а остатки сносит начало следующего запуска. Правило + # запуска решает здесь и то, что стенд увидит снаружи — с "all_done" + # уборка стала бы зелёным концом графа и покрасила бы в зелёный запуск + # с упавшей проверкой (ADR 0003). + @task + def cleanup_tables() -> None: + client = _clickhouse_client() + try: + _drop_tables(client) + _assert_tables_absent(client) + finally: + client.close() + + prepared = prepare_tables() + written = write_marker() + checked = read_from_node_2(written) + + prepared >> written >> checked >> cleanup_tables() test_clickhouse() diff --git a/dags/test_kafka.py b/dags/test_kafka.py index dec32c9..f6ed3df 100644 --- a/dags/test_kafka.py +++ b/dags/test_kafka.py @@ -9,6 +9,8 @@ import uuid from airflow.sdk import dag, get_current_context, task BROKER = "kafka:9092" +# Топик пробника постоянный, и стенд опирается на автосоздание: срок хранения +# чистит записи, а не топик. В бою автосоздание обычно выключают. TOPIC = "airflow_integration_probe" @@ -32,8 +34,6 @@ def test_kafka(): consumer = None try: - # Стенд опирается на автосоздание: срок хранения чистит записи, а - # не топик; в бою автосоздание обычно выключают. delivery_errors: list[str] = [] delivered_offsets: list[tuple[int, int]] = [] diff --git a/docs/adr/0003-dag-run-state.md b/docs/adr/0003-dag-run-state.md new file mode 100644 index 0000000..24dfec0 --- /dev/null +++ b/docs/adr/0003-dag-run-state.md @@ -0,0 +1,55 @@ +# ADR 0003. Уборка в DAG и состояние запуска + +Дата: 31 июля 2026 года. Статус: принято. + +## Решение + +Задача уборки в DAG стенда идёт обычным правилом запуска: только после успеха +всех предыдущих шагов. `trigger_rule="all_done"` и задачи-teardown в хвосте +DAG не используются. Упавший DAG оставляет кластер таким, каким сломался; +чистое состояние обеспечивает первая задача следующего запуска — она сносит +остатки через `DROP ... IF EXISTS` и убеждается, что их нет. + +Правило шире уборки: в конце DAG не должно стоять задачи, которая зеленеет при +отказе предыдущих. + +## Почему + +Состояние запуска Airflow считает по концам графа: зелёные концы — зелёный +запуск, а упавшее выше для итога немо. Уборка с `all_done` написана так, чтобы +отработать при любом исходе, и, встав в конец, красит запуск в зелёный именно +тогда, когда проверка упала. Для стенда это худший исход: `make smoke` +спрашивает у Airflow одно поле — состояние запуска, — и на сломанном кластере +докладывает, что пробник прошёл. + +Обычное правило запуска переворачивает счёт по концам в нашу пользу. После +отказа выше уборка уходит в `upstream_failed`, это состояние отказа, и запуск +краснеет по ней. Видны обе половины: и отказ рабочей задачи, и отказ самой +уборки, если кластер жив, а снести таблицы не вышло. Документированный приём +Airflow для тех, кому `all_done` в хвосте всё же нужен, — отдельная +задача-сторож с правилом `one_failed`; здесь она не понадобилась. + +Цена решения — между падением и следующим запуском в `default` лежат служебные +таблицы пробника. Размен принят: тот же выбор уже сделан для постоянного топика +Kafka, который пробник не удаляет. + +## Что проверено + +31 июля 2026 года на живом стенде с Airflow 3.3.0. Три прогона `test_clickhouse` +через API: зелёный (четыре задачи `success`, таблиц после прогона нет), красный +с удалённой на ноде 2 локальной таблицей (`prepare_tables` — `failed`, +остальные — `upstream_failed`, состояние запуска `failed`, следы поломки на +кластере), затем снова зелёный (остатки снесены первой задачей, таблиц 0). + +Оба отвергнутых варианта проверены так же и оба дали `success` при упавшей +проверке: уборка с `trigger_rule="all_done"` и уборка +`as_teardown(on_failure_fail_dagrun=True)`. + +Семантика правил запуска и задач-teardown сверена через MCP Context7 по +документации Airflow: раздел про setup/teardown (умолчание +`on_failure_fail_dagrun=False`, правило `ALL_DONE_SETUP_SUCCESS`) и раздел +лучших практик про задачу-сторож. Оттуда же взята формулировка ловушки. +Правило прочитано в исходниках +установленного Airflow — `airflow/models/dagrun.py`, `is_effective_leaf`: +концом графа считается задача, ниже которой только задачи-teardown с +`on_failure_fail_dagrun=False`, и которая сама не такова. diff --git a/tests/dag-probes-unit.py b/tests/dag-probes-unit.py index 78995ed..914b2a0 100644 --- a/tests/dag-probes-unit.py +++ b/tests/dag-probes-unit.py @@ -9,9 +9,17 @@ import unittest from pathlib import Path -def load_clickhouse_dag_module(): +class DeclaredTask: + """Заглушка объявленной задачи: держит только цепочку через `>>`.""" + + def __rshift__(self, other): + return other + + +def load_clickhouse_dag_tasks(): airflow_module = types.ModuleType("airflow") sdk_module = types.ModuleType("airflow.sdk") + captured_tasks = {} def dag(**_kwargs): def decorate(function): @@ -19,11 +27,16 @@ def load_clickhouse_dag_module(): return decorate - def task(function): - def declare_task(*_args, **_kwargs): - return None + def task(function=None, **_kwargs): + def capture(target): + captured_tasks[target.__name__] = target - return declare_task + def declare_task(*_args, **_kwargs): + return DeclaredTask() + + return declare_task + + return capture(function) if function is not None else capture sdk_module.Connection = object sdk_module.dag = dag @@ -39,7 +52,7 @@ def load_clickhouse_dag_module(): raise RuntimeError("не удалось загрузить модуль пробника ClickHouse") module = importlib.util.module_from_spec(spec) spec.loader.exec_module(module) - return module + return module, captured_tasks def load_kafka_dag_task(): @@ -77,12 +90,26 @@ def load_kafka_dag_task(): return captured_tasks["check_round_trip"] -class FailingCleanupClient: - def __init__(self) -> None: - self.closed = False +class QueryResult: + def __init__(self, result_rows: list[tuple]) -> None: + self.result_rows = result_rows - def command(self, _sql: str) -> None: - raise RuntimeError("очистка недоступна") + +class RecordingClient: + """Клиент ClickHouse, который запоминает запросы и отвечает заготовкой.""" + + def __init__(self, remaining_tables: list[tuple[str, str]] | None = None) -> None: + self.commands: list[str] = [] + self.queries: list[str] = [] + self.closed = False + self._remaining_tables = remaining_tables or [] + + def command(self, sql: str) -> None: + self.commands.append(" ".join(sql.split())) + + def query(self, sql: str, parameters=None) -> QueryResult: + self.queries.append(" ".join(sql.split())) + return QueryResult(list(self._remaining_tables)) def close(self) -> None: self.closed = True @@ -91,40 +118,66 @@ class FailingCleanupClient: class ClickHouseProbeTests(unittest.TestCase): @classmethod def setUpClass(cls) -> None: - cls.module = load_clickhouse_dag_module() + cls.module, cls.tasks = load_clickhouse_dag_tasks() + + def run_cleanup_task(self, client: RecordingClient) -> None: + original_client_factory = self.module._clickhouse_client + self.module._clickhouse_client = lambda: client + try: + self.tasks["cleanup_tables"]() + finally: + self.module._clickhouse_client = original_client_factory def test_marker_path_requires_different_nodes_and_first_shard(self) -> None: marker = "свой-маркер" self.module._assert_marker_path( - local_rows=[("clickhouse-01-host", marker)], distributed_rows=[(1, "clickhouse-01-host", marker)], + write_hostname="clickhouse-01-host", node_2_hostname="clickhouse-02-host", marker=marker, ) with self.assertRaisesRegex(RuntimeError, "разных нод"): self.module._assert_marker_path( - local_rows=[("clickhouse-02-host", marker)], distributed_rows=[(2, "clickhouse-02-host", marker)], + write_hostname="clickhouse-02-host", node_2_hostname="clickhouse-02-host", marker=marker, ) with self.assertRaisesRegex(RuntimeError, "первого шарда"): self.module._assert_marker_path( - local_rows=[("clickhouse-01-host", marker)], distributed_rows=[(2, "clickhouse-01-host", marker)], + write_hostname="clickhouse-01-host", node_2_hostname="clickhouse-02-host", marker=marker, ) - def test_cleanup_keeps_original_error_as_primary(self) -> None: - client = FailingCleanupClient() - original_error = RuntimeError("маркер не найден") + def test_cleanup_drops_tables_and_checks_both_nodes(self) -> None: + client = RecordingClient() - self.module._cleanup_clickhouse_client(client, original_error) + self.run_cleanup_task(client) + + dropped = { + table + for table in (self.module.LOCAL_TABLE, self.module.DISTRIBUTED_TABLE) + if any( + command.startswith(f"DROP TABLE IF EXISTS default.{table} ") + for command in client.commands + ) + } + self.assertEqual( + dropped, {self.module.LOCAL_TABLE, self.module.DISTRIBUTED_TABLE} + ) + self.assertEqual(len(client.queries), len(self.module.NODES)) + self.assertTrue(client.closed) + + def test_cleanup_closes_client_when_tables_survive(self) -> None: + client = RecordingClient(remaining_tables=[("airflow_probe_local", "Log")]) + + with self.assertRaisesRegex(RuntimeError, "служебные таблицы остались"): + self.run_cleanup_task(client) self.assertTrue(client.closed) - self.assertIn("очистка недоступна", "\n".join(original_error.__notes__)) class KafkaProbeTests(unittest.TestCase): diff --git a/tests/stand-smoke-guards.sh b/tests/stand-smoke-guards.sh index b91b3a2..a33b52f 100755 --- a/tests/stand-smoke-guards.sh +++ b/tests/stand-smoke-guards.sh @@ -75,17 +75,20 @@ break_clickhouse_probe() { return 1 } +# Пробник разбит на четыре задачи, и упасть может любая из них: поломка на ноде +# 2 видна и проверке набора таблиц, и чтению маркера. Поэтому берём последний +# каталог запуска целиком и ищем образец по журналам всех его задач. clickhouse_break_is_reported() { compose exec -T airflow-scheduler bash -ceu ' latest="$( find /opt/airflow/logs/dag_id=test_clickhouse \ - -type f -name "attempt=1.log" -printf "%T@ %p\n" | + -mindepth 1 -maxdepth 1 -type d -name "run_id=*" -printf "%T@ %p\n" | sort -nr | head -n 1 )" latest="${latest#* }" test -n "$latest" - grep -Eq \ + grep -Eqr --include="attempt=1.log" \ "неверный набор таблиц на ноде 2|Unknown table expression identifier '\''default.airflow_probe_local'\''" \ "$latest" ' From c130319ac7b3a019b7baf2840e8f4c9f314a9588 Mon Sep 17 00:00:00 2001 From: Dmitry Dementiev Date: Fri, 31 Jul 2026 19:15:03 +0300 Subject: [PATCH 2/2] =?UTF-8?q?docs(airflow):=20=D1=88=D0=B0=D0=BF=D0=BA?= =?UTF-8?q?=D0=B0=20=D1=82=D0=B5=D1=81=D1=82=D0=B0=20ClickHouse=20=D0=BE?= =?UTF-8?q?=D0=B1=D1=8A=D1=8F=D1=81=D0=BD=D1=8F=D0=B5=D1=82,=20=D1=87?= =?UTF-8?q?=D1=82=D0=BE=20=D0=BE=D0=BD=20=D0=BF=D1=80=D0=BE=D0=B2=D0=B5?= =?UTF-8?q?=D1=80=D1=8F=D0=B5=D1=82?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Зачем: пробники — единственный пример DAG в стенде, по ним будут учиться. Читатель начинал с кода, не зная, что тест утверждает и чем он отличается от боевого кода, — и рисковал скопировать приёмы проверки связности в ETL. Что: строка описания test_clickhouse развёрнута в список проверяемых утверждений и оговорку, что образцом ETL этот тест не является. Описание уходит в doc_md и видно в интерфейсе Airflow. Правило комментариев из #20 строк описания DAG не касается, границы тикета не задеты. Проверка: make config-test — 4/3/6, ошибок 0. make smoke — 24 проверки зелены, включая оба пробника; красной осталась только память стенда (3250 MiB против порога 3242,5), причина известна и разбирается в #21. --- dags/test_clickhouse.py | 16 +++++++++++++++- 1 file changed, 15 insertions(+), 1 deletion(-) diff --git a/dags/test_clickhouse.py b/dags/test_clickhouse.py index 7a0c630..17d9808 100644 --- a/dags/test_clickhouse.py +++ b/dags/test_clickhouse.py @@ -1,4 +1,18 @@ -"""Сквозная проверка подключения Airflow к кластеру ClickHouse.""" +"""Сквозная проверка связности Airflow с кластером ClickHouse. + +Тест проверяет: + +- до запуска служебных таблиц нет ни на одной ноде; +- после создания на обеих лежит ожидаемая пара движков; +- маркер запуска вставлен в локальную таблицу ноды 1 и найден там ровно одной + строкой; +- нода 2 читает этот маркер через Distributed и видит его на первом шарде; +- после уборки таблиц не осталось ни на одной ноде. + +Этот тест — не образец ETL. Он ходит ON CLUSTER на каждом шаге, заводит и +сносит собственные служебные таблицы за один запуск и читает ноду 2 запросом +remote(): это приёмы проверки связности, а не загрузки данных. +""" from __future__ import annotations