Files
clickstream-data-platform/dags/test_clickhouse.py
T
ddadmin c130319ac7 docs(airflow): шапка теста ClickHouse объясняет, что он проверяет
Зачем: пробники — единственный пример 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.
2026-07-31 19:15:03 +03:00

245 lines
9.0 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""Сквозная проверка связности Airflow с кластером ClickHouse.
Тест проверяет:
- до запуска служебных таблиц нет ни на одной ноде;
- после создания на обеих лежит ожидаемая пара движков;
- маркер запуска вставлен в локальную таблицу ноды 1 и найден там ровно одной
строкой;
- нода 2 читает этот маркер через Distributed и видит его на первом шарде;
- после уборки таблиц не осталось ни на одной ноде.
Этот тест — не образец ETL. Он ходит ON CLUSTER на каждом шаге, заводит и
сносит собственные служебные таблицы за один запуск и читает ноду 2 запросом
remote(): это приёмы проверки связности, а не загрузки данных.
"""
from __future__ import annotations
import datetime
import uuid
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 _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
FROM {source}
WHERE database = 'default'
AND name IN ('{LOCAL_TABLE}', '{DISTRIBUTED_TABLE}')
ORDER BY name
"""
)
return result.result_rows
def _drop_tables(client) -> None:
client.command(
f"DROP TABLE IF EXISTS default.{DISTRIBUTED_TABLE} "
f"ON CLUSTER {CLUSTER} SYNC"
)
client.command(
f"DROP TABLE IF EXISTS default.{LOCAL_TABLE} ON CLUSTER {CLUSTER} SYNC"
)
def _assert_tables_absent(client) -> None:
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(
*,
distributed_rows: list[tuple[int, str, str]],
write_hostname: str,
node_2_hostname: str,
marker: str,
) -> None:
if write_hostname == node_2_hostname:
raise RuntimeError("запись и чтение маркера должны выполняться с разных нод")
expected_rows = [(1, write_hostname, marker)]
if distributed_rows != expected_rows:
raise RuntimeError(
"нода 2 не прочитала маркер первого шарда через Distributed: "
f"{marker}, получено {distributed_rows}"
)
@dag(
dag_id="test_clickhouse",
schedule=None,
start_date=datetime.datetime(2026, 1, 1, tzinfo=datetime.timezone.utc),
catchup=False,
tags=["проверка"],
doc_md=__doc__,
)
def test_clickhouse():
@task
def prepare_tables() -> None:
client = _clickhouse_client()
try:
_drop_tables(client)
_assert_tables_absent(client)
client.command(
f"""
CREATE TABLE default.{LOCAL_TABLE} ON CLUSTER {CLUSTER}
(
marker String
)
ENGINE = ReplicatedMergeTree(
'/clickhouse/tables/{{shard}}/{LOCAL_TABLE}',
'{{replica}}'
)
ORDER BY marker
"""
)
client.command(
f"""
CREATE TABLE default.{DISTRIBUTED_TABLE} ON CLUSTER {CLUSTER}
AS default.{LOCAL_TABLE}
ENGINE = Distributed(
'{CLUSTER}',
'default',
'{LOCAL_TABLE}',
cityHash64(marker)
)
"""
)
_assert_tables_created(client)
finally:
client.close()
@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]],
column_names=["marker"],
)
local_rows = client.query(
f"""
SELECT hostName(), marker
FROM default.{LOCAL_TABLE}
WHERE marker = {{marker:String}}
""",
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()
FROM remote('clickhouse-02:9000', system.one)
"""
).result_rows
if len(node_2_rows) != 1:
raise RuntimeError(
f"не удалось определить имя ноды 2: {node_2_rows}"
)
distributed_rows = client.query(
f"""
SELECT _shard_num, hostName(), marker
FROM remote(
'clickhouse-02:9000',
'default',
'{DISTRIBUTED_TABLE}'
)
WHERE marker = {{marker:String}}
""",
parameters={"marker": written["marker"]},
).result_rows
_assert_marker_path(
distributed_rows=distributed_rows,
write_hostname=written["hostname"],
node_2_hostname=node_2_rows[0][0],
marker=written["marker"],
)
finally:
client.close()
# Уборка идёт только после успеха: упавший пробник оставляет кластер таким,
# каким сломался, а остатки сносит начало следующего запуска. Правило
# запуска решает здесь и то, что стенд увидит снаружи — с "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()