Files
clickstream-data-platform/tests/dag-probes-unit.py
T
ddadmin dd491ebe33 refactor(airflow): пробник Kafka разрезан на запись и чтение
Зачем: пробник Kafka — второй и последний пример DAG в стенде, и после #20
единственный, который не показывает, как здесь пишут код. Одним красным
квадратом он к тому же не отвечал на вопрос, какая половина круга отказала:
брокер не принял маркер или не отдал его обратно.

Что:
- вместо одной задачи check_round_trip две: write_marker пишет маркер и
  возвращает адрес записи, read_marker читает по этому адресу и сверяет;
- адрес и маркер едут между задачами XCom — именованными полями словаря:
  XCom проходит через JSON, и кортеж вернулся бы списком;
- внутри записи адрес собирается NamedTuple RecordAddress — два соседних
  целых в сигнатуре переставляются молча;
- каждая задача заводит своего клиента и закрывает его сама, поэтому часовые
  = None и finally с проверками на None ушли; у продюсера закрыть за собой —
  это flush(): своего close() у него нет, и он же возвращает число
  недоставленных;
- настройки клиентов и все сроки ожидания стали именованными модульными
  константами; безымянных чисел в телах задач не осталось;
- слитное условие доставки разобрано на шесть утверждений, каждое со своим
  именем и своим текстом ошибки;
- успех больше не возвращается из середины цикла: чтение выходит из цикла по
  сообщению или по крайнему сроку, а сверка идёт после;
- шапка файла приведена к форме «Тест проверяет: ...» с абзацем о том, чем
  тест не является; она же уходит в doc_md;
- комментарии стоят ровно в шести местах, где незнакома модель Kafka.

Из настроек консьюмера убран session.timeout.ms: он про членство в группе и
удары сердца координатору, а пробник назначает себе адрес и в группу не
входит — почему его нет, объясняет шапка файла. Поведение не меняется.
Сверено по документации confluent-kafka-python через Context7:
session.timeout.ms описан как срок сессии группы, flush() возвращает число
оставшихся в очереди сообщений.

KafkaProbeTests держался за flush_timeouts == [10, 1], то есть за устройство
finally, которого больше нет. На его место встали две проверки свойств:
продюсер закрыт даже тогда, когда отказала запись, и чтение назначается ровно
на тот адрес, который вернул брокер.

Тело issue #22 поправлено тем же изменением: там было записано «задача
остаётся одна» — это расхождение с тем, о чём договаривались в гриллинге.

Проверка: make config-test зелен; make smoke — оба пробника зелены, красной
осталась только проверка памяти стенда по причине из #21.

Closes #22
2026-07-31 20:10:00 +03:00

308 lines
10 KiB
Python

"""Малые проверки логики пробников ClickHouse и Kafka без запуска Airflow."""
from __future__ import annotations
import importlib.util
import sys
import types
import unittest
from pathlib import Path
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):
return function
return decorate
def task(function=None, **_kwargs):
def capture(target):
captured_tasks[target.__name__] = target
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
sdk_module.get_current_context = lambda: {}
sdk_module.task = task
airflow_module.sdk = sdk_module
sys.modules["airflow"] = airflow_module
sys.modules["airflow.sdk"] = sdk_module
dag_path = Path(__file__).resolve().parents[1] / "dags" / "test_clickhouse.py"
spec = importlib.util.spec_from_file_location("test_clickhouse_dag", dag_path)
if spec is None or spec.loader is None:
raise RuntimeError("не удалось загрузить модуль пробника ClickHouse")
module = importlib.util.module_from_spec(spec)
spec.loader.exec_module(module)
return module, captured_tasks
def load_kafka_dag_tasks():
airflow_module = types.ModuleType("airflow")
sdk_module = types.ModuleType("airflow.sdk")
captured_tasks = {}
def dag(**_kwargs):
def decorate(function):
return function
return decorate
def task(function):
captured_tasks[function.__name__] = function
def declare_task(*_args, **_kwargs):
return None
return declare_task
sdk_module.dag = dag
sdk_module.get_current_context = lambda: {"run_id": "unit-test"}
sdk_module.task = task
airflow_module.sdk = sdk_module
sys.modules["airflow"] = airflow_module
sys.modules["airflow.sdk"] = sdk_module
dag_path = Path(__file__).resolve().parents[1] / "dags" / "test_kafka.py"
spec = importlib.util.spec_from_file_location("test_kafka_dag", dag_path)
if spec is None or spec.loader is None:
raise RuntimeError("не удалось загрузить модуль пробника Kafka")
module = importlib.util.module_from_spec(spec)
spec.loader.exec_module(module)
return captured_tasks
class QueryResult:
def __init__(self, result_rows: list[tuple]) -> None:
self.result_rows = result_rows
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
class ClickHouseProbeTests(unittest.TestCase):
@classmethod
def setUpClass(cls) -> None:
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(
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(
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(
distributed_rows=[(2, "clickhouse-01-host", marker)],
write_hostname="clickhouse-01-host",
node_2_hostname="clickhouse-02-host",
marker=marker,
)
def test_cleanup_drops_tables_and_checks_both_nodes(self) -> None:
client = RecordingClient()
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)
DELIVERED_PARTITION = 3
DELIVERED_OFFSET = 42
class StubMessage:
"""Сообщение Kafka в том объёме, в каком его читает пробник."""
def __init__(self, value: bytes) -> None:
self._value = value
def partition(self) -> int:
return DELIVERED_PARTITION
def offset(self) -> int:
return DELIVERED_OFFSET
def error(self):
return None
def value(self) -> bytes:
return self._value
class StubTopicPartition:
def __init__(self, topic: str, partition: int, offset: int) -> None:
self.topic = topic
self.partition = partition
self.offset = offset
def install_kafka_stub(produce_error: str | None = None):
"""Ставит заглушку `confluent_kafka` и возвращает журналы её клиентов.
Продюсер подтверждает доставку сразу и по известному адресу, консьюмер
отдаёт записанное с первого опроса. Если задан `produce_error`, запись
падает — так проверяется, что продюсер закрывается и на пути отказа.
"""
producers = []
consumers = []
class Producer:
def __init__(self, _config) -> None:
self.written = b""
self.flushed = False
producers.append(self)
def produce(self, _topic, key=None, value=None, on_delivery=None) -> None:
if produce_error is not None:
raise RuntimeError(produce_error)
self.written = value
on_delivery(None, StubMessage(value))
def flush(self, _timeout) -> int:
self.flushed = True
return 0
class Consumer:
def __init__(self, _config) -> None:
self.assigned = []
self.closed = False
consumers.append(self)
def assign(self, partitions) -> None:
self.assigned = partitions
def poll(self, _timeout):
return StubMessage(producers[-1].written)
def close(self) -> None:
self.closed = True
kafka_module = types.ModuleType("confluent_kafka")
kafka_module.Consumer = Consumer
kafka_module.Producer = Producer
kafka_module.TopicPartition = StubTopicPartition
sys.modules["confluent_kafka"] = kafka_module
return producers, consumers
class KafkaProbeTests(unittest.TestCase):
def test_producer_is_closed_when_write_fails(self) -> None:
producers, _ = install_kafka_stub(produce_error="брокер недоступен")
tasks = load_kafka_dag_tasks()
with self.assertRaisesRegex(RuntimeError, "брокер недоступен"):
tasks["write_marker"]()
self.assertEqual(len(producers), 1)
self.assertTrue(producers[0].flushed)
def test_read_goes_to_the_address_broker_returned(self) -> None:
_, consumers = install_kafka_stub()
tasks = load_kafka_dag_tasks()
written = tasks["write_marker"]()
tasks["read_marker"](written)
self.assertEqual(
(written["partition"], written["offset"]),
(DELIVERED_PARTITION, DELIVERED_OFFSET),
)
self.assertEqual(len(consumers), 1)
self.assertEqual(len(consumers[0].assigned), 1)
assigned = consumers[0].assigned[0]
self.assertEqual(assigned.partition, DELIVERED_PARTITION)
self.assertEqual(assigned.offset, DELIVERED_OFFSET)
self.assertTrue(consumers[0].closed)
def run_tests() -> int:
suite = unittest.defaultTestLoader.loadTestsFromModule(sys.modules[__name__])
result = unittest.TestResult()
suite.run(result)
problems = result.failures + result.errors
for test, details in problems:
print(f"ОШИБКА: {test.id()}", file=sys.stderr)
print(details, file=sys.stderr)
passed = result.testsRun - len(problems) - len(result.skipped)
print(f"ИТОГ: пройдено {passed}, ошибок {len(problems)}")
return 0 if result.wasSuccessful() else 1
if __name__ == "__main__":
raise SystemExit(run_tests())