Зачем. Демонстрационный DAG example_clickstream_hello ничего не проверял: он не обращался ни к ClickHouse, ни к Kafka, поэтому его зелёный результат ничего не говорил о стенде. Пробники проверяют связи по-настоящему — и тем же клиентом, каким будут ходить рабочие DAG. Что. - test_clickhouse: пишет строку в ReplicatedMergeTree на ноде 1 и читает её с ноды 2 через Distributed. Данные проходят путь «нода 2 → все шарды → шард ноды 1», то есть проверяется межшардовое чтение, а не одна нода. - test_kafka: пишет в постоянный топик сообщение с меткой прогона и вычитывает его обратно. - infra/airflow/Dockerfile: clickhouse-connect 1.6.0 и confluent-kafka 2.15.0 вшиты в образ, импорт проверяется на сборке — при запуске контейнера пакеты не доустанавливаются. - Проверки: scripts/stand-smoke.sh гоняет оба пробника через API Airflow, scripts/config-test.sh разбирает DAG без стенда, tests/stand-smoke-guards.sh проверяет красный путь, tests/dag-probes-unit.py — модульные проверки разбора. - README и ADR 0001 обновлены тем же изменением. - Удалён dags/example_clickstream_hello.py. Проверка. make config-test — пройдено 3, 3 и 6, ошибок 0. make clean; cp .env.example .env; make up — 116 с на чистых томах. make smoke — пройдено 25, ошибок 0; стенд занимает 2244,0 MiB. make smoke-cluster — все 8 проверок кластера. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
99 lines
3.4 KiB
Python
99 lines
3.4 KiB
Python
"""Проверка записи и чтения сообщения Airflow через Kafka."""
|
|
|
|
from __future__ import annotations
|
|
|
|
import datetime
|
|
import time
|
|
import uuid
|
|
|
|
from airflow.sdk import dag, get_current_context, task
|
|
|
|
BROKER = "kafka:9092"
|
|
TOPIC = "airflow_integration_probe"
|
|
|
|
|
|
@dag(
|
|
dag_id="test_kafka",
|
|
schedule=None,
|
|
start_date=datetime.datetime(2026, 1, 1, tzinfo=datetime.timezone.utc),
|
|
catchup=False,
|
|
tags=["проверка"],
|
|
doc_md=__doc__,
|
|
)
|
|
def test_kafka():
|
|
@task
|
|
def check_round_trip() -> None:
|
|
from confluent_kafka import Consumer, Producer, TopicPartition
|
|
|
|
marker = f"{get_current_context()['run_id']}:{uuid.uuid4()}"
|
|
marker_bytes = marker.encode()
|
|
group_id = f"airflow-probe-{uuid.uuid4()}"
|
|
producer = None
|
|
consumer = None
|
|
|
|
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}"
|
|
)
|
|
|
|
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:
|
|
if producer is not None:
|
|
producer.flush(1)
|
|
if consumer is not None:
|
|
consumer.close()
|
|
|
|
check_round_trip()
|
|
|
|
|
|
test_kafka()
|