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 / gRPCRPC через брокер
Сетевых хоповОдин: клиент — серверДва: клиент — брокер — сервер, и столько же обратно
Точка отказаСерверСервер и брокер: упал брокер — не работают даже живые сервисы
Таймауты, ретраиВстроены, есть в любой библиотекеПишете руками, включая дедупликацию
Коды ошибокСтандартные статусыСвой протокол ошибок в теле ответа
ТрассировкаРаботает из коробкиНужно самому таскать 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 и воркеров без адреса.
Проверьте себя
1. Зачем в RPC поверх RabbitMQ нужен correlation_id, если у клиента уже есть своя очередь для ответов?
AЧтобы брокер знал, в какую очередь положить ответ
BЧтобы сопоставить ответ с конкретным запросом, когда их отправлено несколько параллельно
CЧтобы сервер мог подтвердить сообщение через basic_ack
DЧтобы включить дедупликацию сообщений на стороне брокера
2. Сервису нужен синхронный ответ от другого сервиса, у которого есть обычный сетевой адрес и нагрузка без всплесков. Что выбрать?
ARPC через RabbitMQ — так надёжнее, брокер сохранит запрос
BОбычный HTTP-вызов: брокер добавит лишний хоп и станет точкой отказа
CRPC через RabbitMQ — иначе не получится балансировать нагрузку
DFanout-обменник с ожиданием первого ответа