RPC через брокер
Запрос-ответ поверх очередей: как это устроено, зачем нужны reply_to и correlation_id — и почему в девяти случаях из десяти лучше взять обычный HTTP.
RPC через брокер — приём, при котором клиент кладёт запрос в очередь и ждёт ответ в другой очереди, сопоставляя его с запросом по идентификатору
correlation_id.
Зачем это нужно (и когда не нужно)
Брокер по своей природе асинхронен: отправил и забыл. RPC ломает эту логику — клиент отправляет и ждёт. Возникает законный вопрос: зачем городить огород, если у нас есть HTTP, где запрос-ответ встроен в протокол?
Честный ответ: чаще всего — незачем. Но пара сценариев есть.
- Воркеры без адреса. Рендер-ферма из тридцати машин за NAT, которые сами подключаются к брокеру. Их не поставишь за балансировщик — у них нет входящих портов.
- Естественная очередь и всплески. Генерация PDF занимает 20 секунд. Когда приходит сотня запросов разом, HTTP-балансировщик начнёт сыпать 503 или таймаутить. Брокер спокойно копит их в очереди, а
prefetch_count=1раздаёт по мере освобождения воркеров — бесплатный backpressure. - Уже есть брокер и воркеры. Если инфраструктура задач построена, добавить «а верни ещё и результат» дешевле, чем поднимать отдельный HTTP-сервис.
Во всех остальных случаях синхронный вызов поверх асинхронного транспорта — лишняя сложность. Об этом подробно в конце урока.
Как устроен обмен
Схема простая. Клиент публикует запрос в очередь rpc.pdf и добавляет в свойства сообщения два поля:
| Свойство | Смысл |
reply_to | Имя очереди, куда сервер должен положить ответ. Обратный адрес на конверте. |
correlation_id | Уникальный идентификатор запроса. Возвращается в ответе, чтобы клиент понял, на какой из своих вопросов ему ответили. |
Сервер их не придумывает и не хранит — просто читает из входящего сообщения и вежливо копирует в ответ. Вся память о том, кто чего ждёт, живёт на клиенте.
Клиент
import json
import time
import uuid
import pika
class PdfRpcClient:
def __init__(self, host="localhost"):
self.connection = pika.BlockingConnection(
pika.ConnectionParameters(host=host, heartbeat=60)
)
self.channel = self.connection.channel()
self.responses = {}
# amq.rabbitmq.reply-to — псевдоочередь брокера, её не надо создавать
self.channel.basic_consume(
queue="amq.rabbitmq.reply-to",
on_message_callback=self.on_response,
auto_ack=True, # для direct reply-to подтверждения запрещены
)
def on_response(self, ch, method, props, body):
self.responses[props.correlation_id] = body
def call(self, payload, timeout=30):
corr_id = str(uuid.uuid4())
self.channel.basic_publish(
exchange="",
routing_key="rpc.pdf",
properties=pika.BasicProperties(
reply_to="amq.rabbitmq.reply-to",
correlation_id=corr_id,
content_type="application/json",
expiration="30000", # запрос протухнет в очереди через 30 c
),
body=json.dumps(payload).encode(),
)
deadline = time.monotonic() + timeout
while corr_id not in self.responses:
self.connection.process_data_events(time_limit=1)
if time.monotonic() > deadline:
raise TimeoutError("сервер не ответил за {} c".format(timeout))
return json.loads(self.responses.pop(corr_id))
client = PdfRpcClient()
result = client.call({"invoice_id": 3391, "template": "ru-default"})
print("готов PDF:", result["url"])
Что здесь важно.
uuid4()на каждый запрос. Клиент может выстрелить пятью запросами подряд и получить ответы вразнобой —correlation_idединственное, что позволяет их разложить по местам. Словарьresponses— это и есть та самая «память ожидающих».amq.rabbitmq.reply-to— встроенная псевдоочередь RabbitMQ (режим direct reply-to). Она не существует физически: брокер держит ответный канал в памяти, пока клиент подключён. Раньше приходилось создавать эксклюзивную очередь на каждого клиента; direct reply-to избавляет от этой возни и не мусорит очередями. Единственное ограничение —auto_ack=Trueобязателен, подтверждать ответы нельзя.process_data_events(time_limit=1)— способ «покрутить» соединение, пока мы ждём. Без него pika просто спит и не читает входящие кадры, и ответ никогда не появится.- Таймаут. Ждать вечно нельзя: воркер мог упасть, очередь
rpc.pdfмогла вообще не существовать. Плюсexpirationна сообщении — чтобы запрос не пролежал в очереди сутки и не был обработан, когда его уже никто не ждёт.
Сервер
import json
import pika
def render_pdf(payload):
# ... долгая работа ...
return {"url": "s3://invoices/{}.pdf".format(payload["invoice_id"]), "pages": 2}
def on_request(ch, method, props, body):
payload = json.loads(body)
print("рендерю", payload["invoice_id"], flush=True)
try:
result = render_pdf(payload)
except Exception as exc:
result = {"error": str(exc)} # ошибку тоже надо ВЕРНУТЬ, иначе клиент повиснет
ch.basic_publish(
exchange="",
routing_key=props.reply_to, # обратный адрес из запроса
properties=pika.BasicProperties(
correlation_id=props.correlation_id, # тот же id, что прислал клиент
content_type="application/json",
),
body=json.dumps(result).encode(),
)
ch.basic_ack(delivery_tag=method.delivery_tag)
connection = pika.BlockingConnection(pika.ConnectionParameters(host="localhost"))
channel = connection.channel()
channel.queue_declare(queue="rpc.pdf", durable=True)
channel.basic_qos(prefetch_count=1)
channel.basic_consume(queue="rpc.pdf", on_message_callback=on_request)
print("RPC-сервер ждёт запросы")
channel.start_consuming()
Сервер — обычный воркер из первого урока плюс три строки: публикация ответа в props.reply_to, копирование correlation_id и обязательный ответ даже при ошибке. Последнее упускают чаще всего: если исключение просто улетит в лог, клиент будет висеть до таймаута, хотя ответ «не смог» можно вернуть мгновенно.
Как это работает
Ответная очередь у direct reply-to не настоящая: брокер запоминает канал клиента и, увидев публикацию с ключом вида amq.rabbitmq.reply-to.HASH, доставляет сообщение прямо в этот канал. Если клиент отключился — ответ выбрасывается. Это правильное поведение: ответ, которого никто не ждёт, никому не нужен.
Балансировка достаётся бесплатно. Поднимите три RPC-сервера на одной очереди rpc.pdf — запросы разойдутся между ними, как обычные задачи. Клиенту всё равно, кто ответил: он ждёт свой correlation_id, а не конкретный хост.
Почему обычно лучше HTTP
Сравним честно.
| Аспект | HTTP / gRPC | RPC через брокер |
| Сетевых хопов | Один: клиент — сервер | Два: клиент — брокер — сервер, и столько же обратно |
| Точка отказа | Сервер | Сервер и брокер: упал брокер — не работают даже живые сервисы |
| Таймауты, ретраи | Встроены, есть в любой библиотеке | Пишете руками, включая дедупликацию |
| Коды ошибок | Стандартные статусы | Свой протокол ошибок в теле ответа |
| Трассировка | Работает из коробки | Нужно самому таскать trace-id через свойства сообщения |
| Всплески нагрузки | 503 и переполненный пул соединений | Очередь копит, воркеры разгребают |
Брокер придуман, чтобы развязать отправителя и получателя во времени. RPC связывает их обратно — и оставляет всю сложность брокера, забирая его главное преимущество. Если вам нужен ответ здесь и сейчас, а на другом конце обычный сервис с доступным адресом — берите HTTP. Брокерный RPC оправдан там, где вы вынуждены ждать ответа от очереди задач, а не от сервиса.
И есть третий путь, который часто оказывается лучшим: не ждать вовсе. Клиент отправляет задачу, получает 202 Accepted и task_id, а результат забирает поллингом или получает через WebSocket, когда тот будет готов. Никаких висящих соединений — и никакого RPC.
Частые ошибки
- Одна ответная очередь без проверки
correlation_id. Два параллельных запроса — и клиент забирает чужой ответ. Ошибка живёт незаметно на тестах (там запрос всегда один) и вылезает на проде под нагрузкой. - Новая эксклюзивная очередь на каждый вызов. Кажется логичным, но объявление очереди — это round-trip к брокеру, и при тысяче запросов в секунду вы просто заваливаете его созданием и удалением очередей. Используйте direct reply-to или одну постоянную очередь ответов на клиента.
- Ждать без таймаута.
while self.response is None: process_data_events()без дедлайна — это вечный цикл, если воркер умер. Веб-запрос будет висеть, пока не отвалится браузер. - Сервер молчит при ошибке. Исключение улетело в лог, ответ не отправлен, клиент ждёт таймаут. Всегда возвращайте ответ — успешный или с полем
error. - Запрос без
expiration. Воркеры лежали ночь, утром поднялись и радостно отрендерили тысячу PDF, которых уже никто не ждёт. - Блокирующий вызов внутри веб-обработчика. Каждый запрос пользователя держит поток веб-сервера занятым на всё время рендера. Тридцать пользователей — тридцать занятых потоков. Именно так брокерный RPC превращает масштабируемое приложение в неработающее.
Итоги
- RPC поверх RabbitMQ строится на двух свойствах:
reply_to(куда отвечать) иcorrelation_id(на какой вопрос). - Сервер ничего не помнит: он копирует оба поля из запроса в ответ. Состояние ожидания живёт на клиенте.
amq.rabbitmq.reply-to(direct reply-to) избавляет от создания временных очередей; требуетauto_ack=True.- Таймаут на стороне клиента и
expirationна сообщении обязательны, иначе система висит на первом же упавшем воркере. - Ошибку сервер обязан вернуть ответом, а не съесть в логе.
- Синхронный ответ от обычного сервиса — это работа для HTTP. Брокерный RPC берите ради очереди, backpressure и воркеров без адреса.