Merge pull request 'refactor(airflow): пробник Kafka разрезан на запись и чтение' (#24) from feat/22-kafka-probe-form into main
Reviewed-on: #24
This commit was merged in pull request #24.
This commit is contained in:
+161
-56
@@ -1,10 +1,29 @@
|
|||||||
"""Проверка записи и чтения сообщения Airflow через Kafka."""
|
"""Сквозная проверка записи и чтения сообщения Airflow через Kafka.
|
||||||
|
|
||||||
|
Тест проверяет:
|
||||||
|
|
||||||
|
- брокер принял маркер запуска и подтвердил его ровно одним сообщением;
|
||||||
|
- вместе с подтверждением брокер вернул адрес записи — раздел и смещение;
|
||||||
|
- чтение по этому адресу возвращает записанный маркер;
|
||||||
|
- ошибка брокера при чтении — отказ;
|
||||||
|
- молчание брокера дольше срока — отказ.
|
||||||
|
|
||||||
|
Запись и чтение — две задачи, потому что ломаются они по-разному: в интерфейсе
|
||||||
|
Airflow сразу видно, брокер не принял маркер или не отдал его обратно. Адрес
|
||||||
|
записи едет от первой задачи ко второй обычным XCom.
|
||||||
|
|
||||||
|
Этот тест — не образец боевого потребителя. Он не подписывается на топик, не
|
||||||
|
входит в группу, не коммитит смещения и не создаёт топик: сообщение читается
|
||||||
|
по адресу, который брокер назвал при записи. Это приёмы проверки связности, а
|
||||||
|
не работы с потоком.
|
||||||
|
"""
|
||||||
|
|
||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
import datetime
|
import datetime
|
||||||
import time
|
import time
|
||||||
import uuid
|
import uuid
|
||||||
|
from typing import NamedTuple
|
||||||
|
|
||||||
from airflow.sdk import dag, get_current_context, task
|
from airflow.sdk import dag, get_current_context, task
|
||||||
|
|
||||||
@@ -13,6 +32,89 @@ BROKER = "kafka:9092"
|
|||||||
# чистит записи, а не топик. В бою автосоздание обычно выключают.
|
# чистит записи, а не топик. В бою автосоздание обычно выключают.
|
||||||
TOPIC = "airflow_integration_probe"
|
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(
|
||||||
dag_id="test_kafka",
|
dag_id="test_kafka",
|
||||||
@@ -24,75 +126,78 @@ TOPIC = "airflow_integration_probe"
|
|||||||
)
|
)
|
||||||
def test_kafka():
|
def test_kafka():
|
||||||
@task
|
@task
|
||||||
def check_round_trip() -> None:
|
def write_marker() -> dict[str, str | int]:
|
||||||
from confluent_kafka import Consumer, Producer, TopicPartition
|
from confluent_kafka import Producer
|
||||||
|
|
||||||
marker = f"{get_current_context()['run_id']}:{uuid.uuid4()}"
|
marker = f"{get_current_context()['run_id']}:{uuid.uuid4()}"
|
||||||
marker_bytes = marker.encode()
|
payload = marker.encode()
|
||||||
group_id = f"airflow-probe-{uuid.uuid4()}"
|
|
||||||
producer = None
|
|
||||||
consumer = None
|
|
||||||
|
|
||||||
try:
|
|
||||||
delivery_errors: list[str] = []
|
delivery_errors: list[str] = []
|
||||||
delivered_offsets: list[tuple[int, int]] = []
|
addresses: list[RecordAddress] = []
|
||||||
|
|
||||||
def on_delivery(error, message) -> None:
|
def remember_delivery(error, message) -> None:
|
||||||
if error is not None:
|
if error is not None:
|
||||||
delivery_errors.append(str(error))
|
delivery_errors.append(str(error))
|
||||||
else:
|
else:
|
||||||
delivered_offsets.append((message.partition(), message.offset()))
|
addresses.append(
|
||||||
|
RecordAddress(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 = Producer(PRODUCER_CONFIG)
|
||||||
|
# produce() не пишет, а ставит сообщение в очередь: о судьбе записи
|
||||||
|
# сообщает колбэк, а гарантию даёт flush() — он же возвращает число
|
||||||
|
# недоставленных и он же единственный способ закрыть за собой продюсера.
|
||||||
|
# Раздел выбирается по ключу сообщения, а маркер уникален для запуска,
|
||||||
|
# поэтому адрес записи мы не выбираем, а узнаём от брокера.
|
||||||
|
try:
|
||||||
producer.produce(
|
producer.produce(
|
||||||
TOPIC,
|
TOPIC,
|
||||||
key=marker_bytes,
|
key=payload,
|
||||||
value=marker_bytes,
|
value=payload,
|
||||||
on_delivery=on_delivery,
|
on_delivery=remember_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}"
|
|
||||||
)
|
|
||||||
|
|
||||||
consumer = Consumer(
|
|
||||||
{
|
|
||||||
"bootstrap.servers": BROKER,
|
|
||||||
"group.id": group_id,
|
|
||||||
"enable.auto.commit": False,
|
|
||||||
"session.timeout.ms": 6_000,
|
|
||||||
"socket.timeout.ms": 5_000,
|
|
||||||
}
|
|
||||||
)
|
|
||||||
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:
|
finally:
|
||||||
if producer is not None:
|
undelivered = producer.flush(FLUSH_TIMEOUT_SEC)
|
||||||
producer.flush(1)
|
|
||||||
if consumer is not None:
|
_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()
|
consumer.close()
|
||||||
|
|
||||||
check_round_trip()
|
written = write_marker()
|
||||||
|
read_marker(written)
|
||||||
|
|
||||||
|
|
||||||
test_kafka()
|
test_kafka()
|
||||||
|
|||||||
+91
-21
@@ -1,4 +1,4 @@
|
|||||||
"""Малые проверки логики пробника ClickHouse без запуска Airflow."""
|
"""Малые проверки логики пробников ClickHouse и Kafka без запуска Airflow."""
|
||||||
|
|
||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
@@ -55,7 +55,7 @@ def load_clickhouse_dag_tasks():
|
|||||||
return module, captured_tasks
|
return module, captured_tasks
|
||||||
|
|
||||||
|
|
||||||
def load_kafka_dag_task():
|
def load_kafka_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 = {}
|
captured_tasks = {}
|
||||||
@@ -87,7 +87,7 @@ def load_kafka_dag_task():
|
|||||||
raise RuntimeError("не удалось загрузить модуль пробника Kafka")
|
raise RuntimeError("не удалось загрузить модуль пробника Kafka")
|
||||||
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 captured_tasks["check_round_trip"]
|
return captured_tasks
|
||||||
|
|
||||||
|
|
||||||
class QueryResult:
|
class QueryResult:
|
||||||
@@ -180,43 +180,113 @@ class ClickHouseProbeTests(unittest.TestCase):
|
|||||||
self.assertTrue(client.closed)
|
self.assertTrue(client.closed)
|
||||||
|
|
||||||
|
|
||||||
class KafkaProbeTests(unittest.TestCase):
|
DELIVERED_PARTITION = 3
|
||||||
def test_producer_flushes_when_consumer_creation_fails(self) -> None:
|
DELIVERED_OFFSET = 42
|
||||||
kafka_module = types.ModuleType("confluent_kafka")
|
|
||||||
flush_timeouts = []
|
|
||||||
|
class StubMessage:
|
||||||
|
"""Сообщение Kafka в том объёме, в каком его читает пробник."""
|
||||||
|
|
||||||
|
def __init__(self, value: bytes) -> None:
|
||||||
|
self._value = value
|
||||||
|
|
||||||
class Message:
|
|
||||||
def partition(self) -> int:
|
def partition(self) -> int:
|
||||||
return 0
|
return DELIVERED_PARTITION
|
||||||
|
|
||||||
def offset(self) -> int:
|
def offset(self) -> int:
|
||||||
return 1
|
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:
|
class Producer:
|
||||||
def __init__(self, _config) -> None:
|
def __init__(self, _config) -> None:
|
||||||
pass
|
self.written = b""
|
||||||
|
self.flushed = False
|
||||||
|
producers.append(self)
|
||||||
|
|
||||||
def produce(self, _topic, **kwargs) -> None:
|
def produce(self, _topic, key=None, value=None, on_delivery=None) -> None:
|
||||||
kwargs["on_delivery"](None, Message())
|
if produce_error is not None:
|
||||||
|
raise RuntimeError(produce_error)
|
||||||
|
self.written = value
|
||||||
|
on_delivery(None, StubMessage(value))
|
||||||
|
|
||||||
def flush(self, timeout: int) -> int:
|
def flush(self, _timeout) -> int:
|
||||||
flush_timeouts.append(timeout)
|
self.flushed = True
|
||||||
return 0
|
return 0
|
||||||
|
|
||||||
class Consumer:
|
class Consumer:
|
||||||
def __init__(self, _config) -> None:
|
def __init__(self, _config) -> None:
|
||||||
raise RuntimeError("чтение недоступно")
|
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.Consumer = Consumer
|
||||||
kafka_module.Producer = Producer
|
kafka_module.Producer = Producer
|
||||||
kafka_module.TopicPartition = object
|
kafka_module.TopicPartition = StubTopicPartition
|
||||||
sys.modules["confluent_kafka"] = kafka_module
|
sys.modules["confluent_kafka"] = kafka_module
|
||||||
check_round_trip = load_kafka_dag_task()
|
return producers, consumers
|
||||||
|
|
||||||
with self.assertRaisesRegex(RuntimeError, "чтение недоступно"):
|
|
||||||
check_round_trip()
|
|
||||||
|
|
||||||
self.assertEqual(flush_timeouts, [10, 1])
|
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:
|
def run_tests() -> int:
|
||||||
|
|||||||
Reference in New Issue
Block a user