Merge pull request 'refactor(airflow): пробник ClickHouse разбит на четыре задачи' (#23) from feat/20-dag-probe-form into main

Reviewed-on: #23
This commit was merged in pull request #23.
This commit is contained in:
2026-07-31 19:23:17 +03:00
5 changed files with 240 additions and 100 deletions
+105 -76
View File
@@ -1,4 +1,18 @@
"""Сквозная проверка подключения Airflow к кластеру ClickHouse.""" """Сквозная проверка связности Airflow с кластером ClickHouse.
Тест проверяет:
- до запуска служебных таблиц нет ни на одной ноде;
- после создания на обеих лежит ожидаемая пара движков;
- маркер запуска вставлен в локальную таблицу ноды 1 и найден там ровно одной
строкой;
- нода 2 читает этот маркер через Distributed и видит его на первом шарде;
- после уборки таблиц не осталось ни на одной ноде.
Этот тест — не образец ETL. Он ходит ON CLUSTER на каждом шаге, заводит и
сносит собственные служебные таблицы за один запуск и читает ноду 2 запросом
remote(): это приёмы проверки связности, а не загрузки данных.
"""
from __future__ import annotations from __future__ import annotations
@@ -10,14 +24,42 @@ from airflow.sdk import Connection, dag, get_current_context, task
CLUSTER = "clickstream_cluster" CLUSTER = "clickstream_cluster"
LOCAL_TABLE = "airflow_probe_local" LOCAL_TABLE = "airflow_probe_local"
DISTRIBUTED_TABLE = "airflow_probe_distributed" 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]]: def _clickhouse_client():
source = ( # clickhouse_connect стоит только в образе Airflow, а малые проверки грузят
"remote('clickhouse-02:9000', system.tables)" # этот модуль обычным интерпретатором, где пакета нет. Импорт верхнего
if query_node_2 # уровня красит make config-test, поэтому он живёт здесь.
else "system.tables" 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( result = client.query(
f""" f"""
SELECT name, engine SELECT name, engine
@@ -41,27 +83,33 @@ def _drop_tables(client) -> None:
def _assert_tables_absent(client) -> None: def _assert_tables_absent(client) -> None:
for query_node_2, node_name in ((False, "ноде 1"), (True, "ноде 2")): for node_name, source in NODES:
remaining = _table_engines(client, query_node_2=query_node_2) remaining = _table_engines(client, source)
if remaining: if remaining:
raise RuntimeError( raise RuntimeError(
f"служебные таблицы остались на {node_name}: {remaining}" 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( def _assert_marker_path(
*, *,
local_rows: list[tuple[str, str]],
distributed_rows: list[tuple[int, str, str]], distributed_rows: list[tuple[int, str, str]],
write_hostname: str,
node_2_hostname: str, node_2_hostname: str,
marker: str, marker: str,
) -> None: ) -> None:
if len(local_rows) != 1 or local_rows[0][1] != marker: if write_hostname == node_2_hostname:
raise RuntimeError(f"маркер не найден в локальной таблице: {marker}")
node_1_hostname = local_rows[0][0]
if node_1_hostname == node_2_hostname:
raise RuntimeError("запись и чтение маркера должны выполняться с разных нод") raise RuntimeError("запись и чтение маркера должны выполняться с разных нод")
expected_rows = [(1, node_1_hostname, marker)] expected_rows = [(1, write_hostname, marker)]
if distributed_rows != expected_rows: if distributed_rows != expected_rows:
raise RuntimeError( raise RuntimeError(
"нода 2 не прочитала маркер первого шарда через Distributed: " "нода 2 не прочитала маркер первого шарда через Distributed: "
@@ -69,30 +117,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(
dag_id="test_clickhouse", dag_id="test_clickhouse",
schedule=None, schedule=None,
@@ -103,26 +127,8 @@ def _cleanup_clickhouse_client(client, original_error: BaseException | None) ->
) )
def test_clickhouse(): def test_clickhouse():
@task @task
def check_cluster_path() -> None: def prepare_tables() -> None:
import clickhouse_connect client = _clickhouse_client()
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
try: try:
_drop_tables(client) _drop_tables(client)
_assert_tables_absent(client) _assert_tables_absent(client)
@@ -151,17 +157,15 @@ def test_clickhouse():
) )
""" """
) )
_assert_tables_created(client)
finally:
client.close()
for query_node_2, node_name in ((False, "ноде 1"), (True, "ноде 2")): @task
actual_tables = _table_engines( def write_marker() -> dict[str, str]:
client, client = _clickhouse_client()
query_node_2=query_node_2, try:
) marker = f"{get_current_context()['run_id']}:{uuid.uuid4()}"
if actual_tables != expected_tables:
raise RuntimeError(
f"неверный набор таблиц на {node_name}: {actual_tables}"
)
client.insert( client.insert(
f"default.{LOCAL_TABLE}", f"default.{LOCAL_TABLE}",
[[marker]], [[marker]],
@@ -175,6 +179,16 @@ def test_clickhouse():
""", """,
parameters={"marker": marker}, parameters={"marker": marker},
).result_rows ).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( node_2_rows = client.query(
""" """
SELECT hostName() SELECT hostName()
@@ -195,21 +209,36 @@ def test_clickhouse():
) )
WHERE marker = {{marker:String}} WHERE marker = {{marker:String}}
""", """,
parameters={"marker": marker}, parameters={"marker": written["marker"]},
).result_rows ).result_rows
_assert_marker_path( _assert_marker_path(
local_rows=local_rows,
distributed_rows=distributed_rows, distributed_rows=distributed_rows,
write_hostname=written["hostname"],
node_2_hostname=node_2_rows[0][0], node_2_hostname=node_2_rows[0][0],
marker=marker, marker=written["marker"],
) )
except BaseException as error:
probe_error = error
raise
finally: 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() test_clickhouse()
+2 -2
View File
@@ -9,6 +9,8 @@ import uuid
from airflow.sdk import dag, get_current_context, task from airflow.sdk import dag, get_current_context, task
BROKER = "kafka:9092" BROKER = "kafka:9092"
# Топик пробника постоянный, и стенд опирается на автосоздание: срок хранения
# чистит записи, а не топик. В бою автосоздание обычно выключают.
TOPIC = "airflow_integration_probe" TOPIC = "airflow_integration_probe"
@@ -32,8 +34,6 @@ def test_kafka():
consumer = None consumer = None
try: try:
# Стенд опирается на автосоздание: срок хранения чистит записи, а
# не топик; в бою автосоздание обычно выключают.
delivery_errors: list[str] = [] delivery_errors: list[str] = []
delivered_offsets: list[tuple[int, int]] = [] delivered_offsets: list[tuple[int, int]] = []
+55
View File
@@ -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`, и которая сама не такова.
+73 -20
View File
@@ -9,9 +9,17 @@ import unittest
from pathlib import Path 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") airflow_module = types.ModuleType("airflow")
sdk_module = types.ModuleType("airflow.sdk") sdk_module = types.ModuleType("airflow.sdk")
captured_tasks = {}
def dag(**_kwargs): def dag(**_kwargs):
def decorate(function): def decorate(function):
@@ -19,11 +27,16 @@ def load_clickhouse_dag_module():
return decorate return decorate
def task(function): def task(function=None, **_kwargs):
def declare_task(*_args, **_kwargs): def capture(target):
return None 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.Connection = object
sdk_module.dag = dag sdk_module.dag = dag
@@ -39,7 +52,7 @@ def load_clickhouse_dag_module():
raise RuntimeError("не удалось загрузить модуль пробника ClickHouse") raise RuntimeError("не удалось загрузить модуль пробника ClickHouse")
module = importlib.util.module_from_spec(spec) module = importlib.util.module_from_spec(spec)
spec.loader.exec_module(module) spec.loader.exec_module(module)
return module return module, captured_tasks
def load_kafka_dag_task(): def load_kafka_dag_task():
@@ -77,12 +90,26 @@ def load_kafka_dag_task():
return captured_tasks["check_round_trip"] return captured_tasks["check_round_trip"]
class FailingCleanupClient: class QueryResult:
def __init__(self) -> None: def __init__(self, result_rows: list[tuple]) -> None:
self.closed = False 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: def close(self) -> None:
self.closed = True self.closed = True
@@ -91,40 +118,66 @@ class FailingCleanupClient:
class ClickHouseProbeTests(unittest.TestCase): class ClickHouseProbeTests(unittest.TestCase):
@classmethod @classmethod
def setUpClass(cls) -> None: 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: def test_marker_path_requires_different_nodes_and_first_shard(self) -> None:
marker = "свой-маркер" marker = "свой-маркер"
self.module._assert_marker_path( self.module._assert_marker_path(
local_rows=[("clickhouse-01-host", marker)],
distributed_rows=[(1, "clickhouse-01-host", marker)], distributed_rows=[(1, "clickhouse-01-host", marker)],
write_hostname="clickhouse-01-host",
node_2_hostname="clickhouse-02-host", node_2_hostname="clickhouse-02-host",
marker=marker, marker=marker,
) )
with self.assertRaisesRegex(RuntimeError, "разных нод"): with self.assertRaisesRegex(RuntimeError, "разных нод"):
self.module._assert_marker_path( self.module._assert_marker_path(
local_rows=[("clickhouse-02-host", marker)],
distributed_rows=[(2, "clickhouse-02-host", marker)], distributed_rows=[(2, "clickhouse-02-host", marker)],
write_hostname="clickhouse-02-host",
node_2_hostname="clickhouse-02-host", node_2_hostname="clickhouse-02-host",
marker=marker, marker=marker,
) )
with self.assertRaisesRegex(RuntimeError, "первого шарда"): with self.assertRaisesRegex(RuntimeError, "первого шарда"):
self.module._assert_marker_path( self.module._assert_marker_path(
local_rows=[("clickhouse-01-host", marker)],
distributed_rows=[(2, "clickhouse-01-host", marker)], distributed_rows=[(2, "clickhouse-01-host", marker)],
write_hostname="clickhouse-01-host",
node_2_hostname="clickhouse-02-host", node_2_hostname="clickhouse-02-host",
marker=marker, marker=marker,
) )
def test_cleanup_keeps_original_error_as_primary(self) -> None: def test_cleanup_drops_tables_and_checks_both_nodes(self) -> None:
client = FailingCleanupClient() client = RecordingClient()
original_error = RuntimeError("маркер не найден")
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.assertTrue(client.closed)
self.assertIn("очистка недоступна", "\n".join(original_error.__notes__))
class KafkaProbeTests(unittest.TestCase): class KafkaProbeTests(unittest.TestCase):
+5 -2
View File
@@ -75,17 +75,20 @@ break_clickhouse_probe() {
return 1 return 1
} }
# Пробник разбит на четыре задачи, и упасть может любая из них: поломка на ноде
# 2 видна и проверке набора таблиц, и чтению маркера. Поэтому берём последний
# каталог запуска целиком и ищем образец по журналам всех его задач.
clickhouse_break_is_reported() { clickhouse_break_is_reported() {
compose exec -T airflow-scheduler bash -ceu ' compose exec -T airflow-scheduler bash -ceu '
latest="$( latest="$(
find /opt/airflow/logs/dag_id=test_clickhouse \ 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 | sort -nr |
head -n 1 head -n 1
)" )"
latest="${latest#* }" latest="${latest#* }"
test -n "$latest" test -n "$latest"
grep -Eq \ grep -Eqr --include="attempt=1.log" \
"неверный набор таблиц на ноде 2|Unknown table expression identifier '\''default.airflow_probe_local'\''" \ "неверный набор таблиц на ноде 2|Unknown table expression identifier '\''default.airflow_probe_local'\''" \
"$latest" "$latest"
' '