Как гарантировать, что в очередь RabbitMQ не будут попадать сообщения от других микросервисов?

«Как гарантировать, что в очередь RabbitMQ не будут попадать сообщения от других микросервисов?» — вопрос из категории Брокеры сообщений, который задают на 26% собеседований Node.js Разработчик. Ниже — развёрнутый ответ с разбором ключевых моментов.

Ответ

В Node.js-микросервисах я использую несколько подходов для изоляции очередей RabbitMQ:

1. Эксклюзивные очереди (exclusive queues)

const amqp = require('amqplib');

async function setupPrivateQueue() {
  const connection = await amqp.connect('amqp://localhost');
  const channel = await connection.createChannel();

  // Очередь будет доступна только этому соединению
  const { queue } = await channel.assertQueue('', {
    exclusive: true, // ключевой параметр
    durable: false
  });

  console.log('Private queue created:', queue);
  return { connection, channel, queue };
}

2. Использование отдельных виртуальных хостов (VHost) и прав доступа

// Подключение к выделенному VHost
const connection = await amqp.connect('amqp://service_user:password@localhost/my_service_vhost');

// В RabbitMQ настройки через CLI:
// rabbitmqctl add_vhost my_service_vhost
// rabbitmqctl set_permissions -p my_service_vhost service_user ".*" ".*" ".*"

3. Валидация сообщений и заголовков

const { queue } = await channel.assertQueue('service.tasks', { durable: true });

channel.consume(queue, (msg) => {
  const content = JSON.parse(msg.content.toString());

  // Проверяем источник сообщения
  if (msg.properties.headers?.source !== 'my_service') {
    console.warn('Message from unauthorized source, rejecting');
    channel.nack(msg); // отклоняем сообщение
    return;
  }

  // Валидация структуры
  if (!content.taskId || !content.type) {
    console.error('Invalid message format');
    channel.nack(msg);
    return;
  }

  // Обработка сообщения
  processTask(content);
  channel.ack(msg);
});

4. Использование exchange с routing keys

// Каждый сервис публикует в свой exchange
await channel.assertExchange('service_a.exchange', 'direct', { durable: true });
await channel.assertQueue('service_a.tasks', { durable: true });
await channel.bindQueue('service_a.tasks', 'service_a.exchange', 'tasks');

// Публикация с указанием source
channel.publish('service_a.exchange', 'tasks', Buffer.from(JSON.stringify(payload)), {
  persistent: true,
  headers: { source: 'service_a', version: '1.0' }
});

На практике я комбинирую эти подходы: использую отдельные VHost для production/staging, настраиваю права доступа через RabbitMQ management plugin и всегда добавляю валидацию сообщений в consumer.