Многие тестировщики при тестировании брокеров сообщений создают новое соединение перед каждым тестом. Если топиков в интеграционном тестировании несколько, то тесты выполняются очень долго, и каждое сообщение валидируется отдельно и последовательно в тесте.
Как можно улучшить, чтобы переиспользовать соединение и при этом проводить валидацию моментально при наступлении события? При этом делать это расширяемо и удобно.
Мы можем пойти немного в инженерию.
А что если у нас будет один Kafka listener, который будет вычитывать Kafka топики, и несколько подписчиков, которые при получении события (например, сообщения) будут выполнять свою логику?
Так мы с вами реализуем паттерн Observer.
Сначала реализуем Subject или наш KaflaListener, для простоты там будет один метод publish, который и будет уведомлять подписчиков Observer о наступлении событий.
Этот класс так же имеет метод регистрации подписчиков, которых можно сделать очень много.
class KafkaListener:
def __init__(self):
self._clients: list[Subscriber] = []
self._is_running: bool = False
def start(self):
if self._is_running:
raise RuntimeError("already running")
self._is_running = True
def publish(self, message: str) -> None:
for client in self._clients:
client.read_message(message)
def subscribe(self, client: Subscriber) -> None:
if self._is_running:
raise RuntimeError("already running")
self._clients.append(client)
Дальше опишем интерфейс наших подписчиков.
from typing import Protocol
class Subscriber(Protocol):
def read_message(self, message: str) -> None: ...
Теперь описываем конкретные подписчики, где можем сказать что делать при выполнении read_message, например валидацию или перекладывание в очередь.
class FirstClient:
def read_message(self, message: str) -> None:
print(f"Get message from {self.__class__.__name__}, {message}")
class SecondClient:
def read_message(self, message: str) -> None:
print(f"Get message from {self.__class__.__name__}, {message}")
class ThirdClient:
def read_message(self, message: str) -> None:
print(f"Get message from {self.__class__.__name__}, {message}")
Теперь мы можем увидеть, что при publish от нашего listener, наше сообщение попадет сразу во все подписчики.
kafka = KafkaListener()
kafka.subscribe(FirstClient())
kafka.subscribe(SecondClient())
kafka.subscribe(ThirdClient())
for _ in range(10):
kafka.publish(f"message {_}")
time.sleep(2)
Получим такой результат:
Get message from FirstClient, message 0
Get message from SecondClient, message 0
Get message from ThirdClient, message 0
Get message from FirstClient, message 1
Get message from SecondClient, message 1
Get message from ThirdClient, message 1
...
Использовать такую модель можно и в кейсе с обновлением токенов, где подписчиками будут наши клиенты, и по push-модели там будет обновляться авторизация.
Изучение данного паттерна уже включено в курс по брокерам сообщений, который выйдет после Нового года.
А пока вы можете изучить другие паттерны, такие как Proxy, Facade, Decorator, в курсе Advanced.






