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
This commit is contained in:
2026-07-31 20:10:00 +03:00
parent f7de99f10e
commit dd491ebe33
2 changed files with 271 additions and 96 deletions
+167 -62
View File
@@ -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()