diff --git a/dags/test_kafka.py b/dags/test_kafka.py index f6ed3df..812c2cc 100644 --- a/dags/test_kafka.py +++ b/dags/test_kafka.py @@ -1,10 +1,29 @@ -"""Проверка записи и чтения сообщения Airflow через Kafka.""" +"""Сквозная проверка записи и чтения сообщения Airflow через Kafka. + +Тест проверяет: + +- брокер принял маркер запуска и подтвердил его ровно одним сообщением; +- вместе с подтверждением брокер вернул адрес записи — раздел и смещение; +- чтение по этому адресу возвращает записанный маркер; +- ошибка брокера при чтении — отказ; +- молчание брокера дольше срока — отказ. + +Запись и чтение — две задачи, потому что ломаются они по-разному: в интерфейсе +Airflow сразу видно, брокер не принял маркер или не отдал его обратно. Адрес +записи едет от первой задачи ко второй обычным XCom. + +Этот тест — не образец боевого потребителя. Он не подписывается на топик, не +входит в группу, не коммитит смещения и не создаёт топик: сообщение читается +по адресу, который брокер назвал при записи. Это приёмы проверки связности, а +не работы с потоком. +""" from __future__ import annotations import datetime import time import uuid +from typing import NamedTuple from airflow.sdk import dag, get_current_context, task @@ -13,6 +32,89 @@ BROKER = "kafka:9092" # чистит записи, а не топик. В бою автосоздание обычно выключают. TOPIC = "airflow_integration_probe" +# Сроки ожидания у записи меряют разное, и путать их дорого: DELIVERY — +# сколько библиотека держит сообщение у себя, пока пробует его доставить; +# SOCKET — сколько ждём ответа сети на один запрос; FLUSH — сколько ждём, пока +# библиотека вернёт управление и скажет, чем доставка кончилась. +DELIVERY_TIMEOUT_MS = 10_000 +SOCKET_TIMEOUT_MS = 5_000 +FLUSH_TIMEOUT_SEC = 10 + +PRODUCER_CONFIG = { + "bootstrap.servers": BROKER, + "client.id": "airflow-test-kafka", + "message.timeout.ms": DELIVERY_TIMEOUT_MS, + "socket.timeout.ms": SOCKET_TIMEOUT_MS, +} + +# Группа у чтения одноразовая, автокоммит выключен: пробник не двигает ничьих +# смещений и не оставляет следов на брокере. Обойтись совсем без группы нельзя +# — без group.id библиотека консьюмера не создаст. +CONSUMER_GROUP_PREFIX = "airflow-probe-" +CONSUMER_CONFIG = { + "bootstrap.servers": BROKER, + "enable.auto.commit": False, + "socket.timeout.ms": SOCKET_TIMEOUT_MS, +} + +READ_DEADLINE_SEC = 20 +POLL_TIMEOUT_SEC = 0.5 + + +class RecordAddress(NamedTuple): + """Адрес записи в топике, каким его называет брокер в подтверждении. + + Тип живёт внутри задачи записи: два соседних целых в сигнатуре + переставляются молча, а имена полей этого не дают. Через границу задач + адрес едет отдельными полями словаря — XCom проходит через JSON, и кортеж + вернулся бы на той стороне списком. + """ + + partition: int + offset: int + + +def _assert_all_delivered(undelivered: int, marker: str) -> None: + if undelivered: + raise RuntimeError( + f"Kafka не приняла маркер за {FLUSH_TIMEOUT_SEC} с, " + f"не доставлено сообщений {undelivered}: {marker}" + ) + + +def _assert_no_delivery_errors(delivery_errors: list[str], marker: str) -> None: + if delivery_errors: + raise RuntimeError( + f"Kafka отказалась принять маркер {marker}: {delivery_errors}" + ) + + +def _assert_confirmed_once(addresses: list[RecordAddress], marker: str) -> None: + if len(addresses) != 1: + raise RuntimeError( + f"Kafka подтвердила запись маркера {marker} " + f"не одним сообщением: {addresses}" + ) + + +def _assert_message_arrived(message, marker: str) -> None: + if message is None: + raise RuntimeError( + f"Kafka молчала {READ_DEADLINE_SEC} с и не вернула маркер: {marker}" + ) + + +def _assert_no_read_error(message) -> None: + if message.error(): + raise RuntimeError(f"Kafka вернула ошибку чтения: {message.error()}") + + +def _assert_marker_matches(message, marker: str) -> None: + if message.value() != marker.encode(): + raise RuntimeError( + f"по адресу записи лежит не маркер запуска {marker}: {message.value()!r}" + ) + @dag( dag_id="test_kafka", @@ -24,75 +126,78 @@ TOPIC = "airflow_integration_probe" ) def test_kafka(): @task - def check_round_trip() -> None: - from confluent_kafka import Consumer, Producer, TopicPartition + def write_marker() -> dict[str, str | int]: + from confluent_kafka import Producer marker = f"{get_current_context()['run_id']}:{uuid.uuid4()}" - marker_bytes = marker.encode() - group_id = f"airflow-probe-{uuid.uuid4()}" - producer = None - consumer = None + payload = marker.encode() + delivery_errors: list[str] = [] + addresses: list[RecordAddress] = [] - try: - delivery_errors: list[str] = [] - delivered_offsets: list[tuple[int, int]] = [] - - def on_delivery(error, message) -> None: - if error is not None: - delivery_errors.append(str(error)) - else: - delivered_offsets.append((message.partition(), message.offset())) - - producer = Producer( - { - "bootstrap.servers": BROKER, - "client.id": "airflow-test-kafka", - "message.timeout.ms": 10_000, - "socket.timeout.ms": 5_000, - } - ) - producer.produce( - TOPIC, - key=marker_bytes, - value=marker_bytes, - on_delivery=on_delivery, - ) - undelivered = producer.flush(10) - if undelivered or delivery_errors or len(delivered_offsets) != 1: - raise RuntimeError( - "Kafka не подтвердила запись маркера: " - f"не доставлено {undelivered}, ошибки {delivery_errors}, " - f"смещения {delivered_offsets}" + def remember_delivery(error, message) -> None: + if error is not None: + delivery_errors.append(str(error)) + else: + addresses.append( + RecordAddress(message.partition(), message.offset()) ) - consumer = Consumer( - { - "bootstrap.servers": BROKER, - "group.id": group_id, - "enable.auto.commit": False, - "session.timeout.ms": 6_000, - "socket.timeout.ms": 5_000, - } + producer = Producer(PRODUCER_CONFIG) + # produce() не пишет, а ставит сообщение в очередь: о судьбе записи + # сообщает колбэк, а гарантию даёт flush() — он же возвращает число + # недоставленных и он же единственный способ закрыть за собой продюсера. + # Раздел выбирается по ключу сообщения, а маркер уникален для запуска, + # поэтому адрес записи мы не выбираем, а узнаём от брокера. + try: + producer.produce( + TOPIC, + key=payload, + value=payload, + on_delivery=remember_delivery, ) - partition, offset = delivered_offsets[0] - consumer.assign([TopicPartition(TOPIC, partition, offset)]) - read_deadline = time.monotonic() + 20 - while time.monotonic() < read_deadline: - message = consumer.poll(0.5) - if message is None: - continue - if message.error(): - raise RuntimeError(f"Kafka вернула ошибку чтения: {message.error()}") - if message.value() == marker_bytes: - return - raise RuntimeError(f"Kafka не вернула свой маркер за 20 секунд: {marker}") finally: - if producer is not None: - producer.flush(1) - if consumer is not None: - consumer.close() + undelivered = producer.flush(FLUSH_TIMEOUT_SEC) - check_round_trip() + _assert_all_delivered(undelivered, marker) + _assert_no_delivery_errors(delivery_errors, marker) + _assert_confirmed_once(addresses, marker) + return { + "marker": marker, + "partition": addresses[0].partition, + "offset": addresses[0].offset, + } + + @task + def read_marker(written: dict[str, str | int]) -> None: + from confluent_kafka import Consumer, TopicPartition + + marker = str(written["marker"]) + consumer = Consumer( + {**CONSUMER_CONFIG, "group.id": f"{CONSUMER_GROUP_PREFIX}{uuid.uuid4()}"} + ) + try: + # Пробник назначает себе известный адрес, а не подписывается на + # топик: адрес брокер назвал при записи, поэтому ни группа, которая + # делит разделы между читателями, ни перебалансировка чтению не + # нужны. + consumer.assign( + [TopicPartition(TOPIC, written["partition"], written["offset"])] + ) + deadline = time.monotonic() + READ_DEADLINE_SEC + message = None + # None от poll() — норма, а не отказ: брокер отвечает молчанием и + # когда сообщения нет, и когда оно ещё не доехало. Отличить одно от + # другого нечем, поэтому у чтения обязан быть крайний срок. + while message is None and time.monotonic() < deadline: + message = consumer.poll(POLL_TIMEOUT_SEC) + _assert_message_arrived(message, marker) + _assert_no_read_error(message) + _assert_marker_matches(message, marker) + finally: + consumer.close() + + written = write_marker() + read_marker(written) test_kafka() diff --git a/tests/dag-probes-unit.py b/tests/dag-probes-unit.py index 914b2a0..56b3031 100644 --- a/tests/dag-probes-unit.py +++ b/tests/dag-probes-unit.py @@ -1,4 +1,4 @@ -"""Малые проверки логики пробника ClickHouse без запуска Airflow.""" +"""Малые проверки логики пробников ClickHouse и Kafka без запуска Airflow.""" from __future__ import annotations @@ -55,7 +55,7 @@ def load_clickhouse_dag_tasks(): return module, captured_tasks -def load_kafka_dag_task(): +def load_kafka_dag_tasks(): airflow_module = types.ModuleType("airflow") sdk_module = types.ModuleType("airflow.sdk") captured_tasks = {} @@ -87,7 +87,7 @@ def load_kafka_dag_task(): raise RuntimeError("не удалось загрузить модуль пробника Kafka") module = importlib.util.module_from_spec(spec) spec.loader.exec_module(module) - return captured_tasks["check_round_trip"] + return captured_tasks class QueryResult: @@ -180,43 +180,113 @@ class ClickHouseProbeTests(unittest.TestCase): 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_flushes_when_consumer_creation_fails(self) -> None: - kafka_module = types.ModuleType("confluent_kafka") - flush_timeouts = [] + def test_producer_is_closed_when_write_fails(self) -> None: + producers, _ = install_kafka_stub(produce_error="брокер недоступен") + tasks = load_kafka_dag_tasks() - class Message: - def partition(self) -> int: - return 0 + with self.assertRaisesRegex(RuntimeError, "брокер недоступен"): + tasks["write_marker"]() - def offset(self) -> int: - return 1 + self.assertEqual(len(producers), 1) + self.assertTrue(producers[0].flushed) - class Producer: - def __init__(self, _config) -> None: - pass + def test_read_goes_to_the_address_broker_returned(self) -> None: + _, consumers = install_kafka_stub() + tasks = load_kafka_dag_tasks() - def produce(self, _topic, **kwargs) -> None: - kwargs["on_delivery"](None, Message()) + written = tasks["write_marker"]() + tasks["read_marker"](written) - def flush(self, timeout: int) -> int: - flush_timeouts.append(timeout) - return 0 - - class Consumer: - def __init__(self, _config) -> None: - raise RuntimeError("чтение недоступно") - - kafka_module.Consumer = Consumer - kafka_module.Producer = Producer - kafka_module.TopicPartition = object - sys.modules["confluent_kafka"] = kafka_module - check_round_trip = load_kafka_dag_task() - - with self.assertRaisesRegex(RuntimeError, "чтение недоступно"): - check_round_trip() - - self.assertEqual(flush_timeouts, [10, 1]) + 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: