Не должно быть более подробного руководства по RabbitMQ!

задняя часть

@[toc] С августа многие учебники по RabbitMQ периодически публикуются. Недавно я нашел время, чтобы организовать их. В будущем может быть видеоурок, так что следите за обновлениями.

1. Промежуточное программное обеспечение общего сообщения, большой ПК

Когда дело доходит до промежуточного программного обеспечения сообщений, по оценкам, каждый может более или менее говорить об этом, ActiveMQ, RabbitMQ, RocketMQ, Kafka и т. д. и различных протоколах, таких как JMS, AMQP и т. д. Однако эти промежуточное программное обеспечение сообщений имеют свои особенности. Какой из них выбрать в разработке? Сегодня Сун Гэ пришел разобраться со своими друзьями.

1.1 Несколько протоколов

Давайте поговорим о некоторых распространенных протоколах промежуточного программного обеспечения сообщений.

1.1.1 JMS

1.1.1.1 Введение в JMS

Сначала поговорим о JMS.

Полное название JMS — служба сообщений Java, что похоже на JDBC. В отличие от JDBC, JMS — это интерфейс службы сообщений JavaEE. Существуют две основные версии JMS:

  • 1.1
  • 2.0.

По сравнению с двумя, последний в основном упрощает код для отправки и получения сообщений.

Учитывая, что промежуточное программное обеспечение сообщений является очень распространенным инструментом, JavaEE разработала для этого специальную спецификацию JMS.

Однако, как и JDBC, как спецификация, JMS представляет собой просто набор интерфейсов и не содержит конкретных реализаций.Если мы хотим использовать JMS, нам обычно нужны соответствующие реализации, точно так же, как для использования JDBC требуются соответствующие драйверы.

1.1.1.2 Модель JMS

Служба сообщений JMS поддерживает две модели сообщений:

  • Одноранговая модель или модель очереди
  • опубликовать/подписать модель

В одноранговой модели или модели очереди производитель публикует сообщения в определенной очереди, а потребитель читает сообщения из этой очереди. Здесь производитель знает очередь потребителя и отправляет сообщения непосредственно в соответствующую очередь. Это одноранговая модель обмена сообщениями, которая вкратце выглядит следующим образом:

  1. Только один потребитель получит сообщение.
  2. Производитель не должен находиться в рабочем состоянии, когда потребитель потребляет сообщение, и потребитель также не должен находиться в рабочем состоянии, когда сообщение отправлено, то есть производитель и потребитель сообщения полностью отделены друг от друга. .
  3. Каждое успешно обработанное сообщение подписывается потребителем сообщения.

Модель «издатель/подписчик» поддерживает публикацию сообщений в определенной теме сообщений, и потребители могут определять интересующие их темы. Это модель сообщений «точка-точка».

  • Несколько потребителей могут потреблять сообщения.
  • Между издателем и подписчиком существует временная зависимость, издателю необходимо создать подписку, чтобы клиент мог подписаться; подписчик должен оставаться в сети, чтобы получать сообщения; конечно, если подписчик создает постоянную подписку, то сообщения, опубликованные производитель сообщения, пока подписчик не подключен, будет повторно опубликован, когда подписчик повторно подключится.
1.1.1.3 Реализация JMS

Промежуточное программное обеспечение сообщений с открытым исходным кодом, которое поддерживает JMS:

  • Kafka
  • Apache ActiveMQ
  • HornetQ для сообщества JBoss
  • Joram
  • MantaRay от Коридана
  • OpenJMS

Некоторое коммерческое промежуточное ПО для обмена сообщениями с поддержкой JMS:

  • WebLogic Server JMS
  • EMS
  • GigaSpaces
  • iBus
  • IONA JMS
  • IQManager (приобретен Sun Microsystems в августе 2005 г.)
  • JMS+
  • Nirvana
  • SonicMQ
  • WebSphere MQ

Многие из них были раскопаны археологами Songge.На самом деле Kafka и ActiveMQ могут быть теми, с которыми мы больше контактируем в нашей повседневной разработке.

1.1.2 AMQP

1.1.2.1 Введение в AMQP

Другой протокол, связанный с промежуточным программным обеспечением сообщений, — AMQP.

Спрос на Message Queue имеет давнюю историю.В 1980-х годах такие компании, как Goldman Sachs, внедрили продукты Teknekron в финансовые операции.В то время программное обеспечение Message Queue называлось информационной шиной (TIB). TIB был принят телекоммуникационными и коммуникационными компаниями, а Reuters приобрело Teknekron. После этого IBM разработала MQSeries, а Microsoft разработала Microsoft Message Queue (MSMQ). Проблема с этими коммерческими поставщиками MQ заключается в привязке к поставщику и высоких ценах. В 2001 году служба сообщений Java пыталась решить проблему блокировки и интерактивности, но она оказалась более громоздкой для приложений.

Поэтому в 2004 году JPMorgan Chase и iMatrix начали разработку открытого стандарта Advanced Message Queuing Protocol (AMQP). В 2006 году была выпущена спецификация AMQP. В 2007 году был выпущен RabbitMQ 1.0, разработанный Rabbit Technologies на основе стандарта AMQP.

Последняя версия RabbitMQ — 3.5.7, основанная на AMQP 0-9-1.

В протоколе AMQP отправка и получение сообщений включают следующие понятия:

  • Брокер: приложение, которое получает и распространяет сообщения RabbitMQ, который мы используем каждый день, — это брокер сообщений.
  • Виртуальный хост: предназначен для мультиарендности и безопасности, делит основные компоненты AMQP на виртуальную группу, аналогичную концепции пространства имен в сети. Когда несколько разных пользователей используют услуги, предоставляемые одним и тем же RabbitMQ, несколько виртуальных хостов могут быть разделены, и каждый пользователь создает свой собственный виртуальный хост.exchange/queueПодождите, эта Сун Гэ уже писала специальную статью, Портал:Как понять VirtualHost в RabbitMQ.
  • Соединение: TCP-соединение между издателем/потребителем и брокером. Операция отключения будет выполняться только на стороне клиента, и брокер не отключится, если не произойдет сбой сети или проблема со службой брокера.
  • Канал: если соединение устанавливается каждый раз при доступе к RabbitMQ, накладные расходы на установление TCP-соединения будут огромными, а эффективность будет низкой при большом количестве сообщений. Канал — это логическое соединение, установленное в Connection. Если приложение поддерживает многопоточность, обычно каждый поток создает отдельный канал для связи. Метод AMQP содержит идентификатор канала, чтобы помочь клиенту и брокеру сообщений идентифицировать канал, поэтому каналы полностью изолированы. Будучи облегченным соединением, Channel значительно снижает нагрузку операционной системы на установление TCP-соединения.Как использовать страницу управления RabbitMQОн также подробно описан в статье.
  • Обмен: Сообщение поступает на первую остановку Брокера, согласно правилам распределения, соответствует ключу маршрутизации в таблице запросов и распределяет сообщение в очередь. Распространенными типами являются: прямой (точка-точка), тематический (публикация и подписка) и разветвленный (широковещательный).
  • Очередь: Сообщение, наконец, отправляется сюда, чтобы дождаться, пока Потребитель заберет его. Сообщение может быть скопировано в несколько очередей одновременно.
  • Binding: виртуальное соединение между Exchange и Queue, привязка может содержать ключ маршрутизации. Информация о привязке сохраняется в таблице запросов в Exchange в качестве основы для распространения сообщений.
1.1.2.2 Реализация AMQP

Давайте рассмотрим некоторые конкретные продукты промежуточного программного обеспечения сообщений, которые реализуют протокол AMQP.

  • Apache Qpid
  • Apache ActiveMQ
  • RabbitMQ

Может быть, некоторые друзья задаются вопросом, почему существует ActiveMQ? На самом деле ActiveMQ поддерживает не только JMS, но и AMQP.

Кроме того, есть хорошо известный RocketMQ, созданный Али.Это пользовательский протокол.Сообщество также предоставляет JMS, но он не очень зрелый.Song Ge уточнит позже.

1.1.3 MQTT

Небольшие партнеры, занимающиеся разработкой IoT, должны часто вступать в контакт с этим протоколом. MQTT (Message Queuing Telemetry Transport, Message Queuing Telemetry Transport) — это протокол обмена мгновенными сообщениями, разработанный IBM. В настоящее время он кажется одним из наиболее важных протоколов в разработке IoT. ., этот протокол поддерживает все платформы и может подключать практически все сетевые элементы к внешнему миру. Он используется в качестве протокола связи для датчиков и приводов (например, для подключения домов к Интернету через Twitter). Он имеет преимущества простого формата, небольшое потребление полосы пропускания и поддержка мобильных устройств.Терминальная связь, поддержка PUSH, подходит для встроенных систем.

1.1.4 XMPP

XMPP (расширяемый протокол обмена сообщениями и присутствия) — это протокол на основе XML, который в основном используется для обмена мгновенными сообщениями (IM) и обнаружения присутствия в Интернете и подходит для почти мгновенных операций между серверами. По своей сути потоковая передача на основе XML представляет собой протокол, который в конечном итоге может позволить пользователям Интернета отправлять мгновенные сообщения кому-либо еще в Интернете, даже с использованием других операционных систем и браузеров. Его преимуществами являются общая открытость, сильная совместимость, расширяемость и высокая безопасность, а недостатком является то, что формат кодирования XML занимает большую полосу пропускания.

1.1.5 JMS Vs AMQP

Для нас, инженеров Java, протоколы JMS и AMQP должны быть теми, с которыми мы больше общаемся каждый день.Поскольку JMS и AMQP являются протоколами, в чем разница между ними? Взгляните на следующую картинку:

Эта картина очень четкая, не буду многословен.

1.2. Важные продукты

1.2.1 ActiveMQ

ActiveMQ – это подпроект Apache. Он реализуется провайдером JMS, который полностью поддерживает спецификации JMS1.1 и J2EE1.4. Расширенные сценарии приложений могут быть эффективно реализованы с помощью небольшого объема кода, и он поддерживает подключаемые транспортные протоколы. Такие как:in-VM, TCP, SSL, NIO, UDP, multicast, JGroups and JXTA transports.

ActiveMQ поддерживает часто используемые клиенты на нескольких языках, таких как C++, Java, .Net, Python, Php, Ruby и т. д.

ActiveMQ теперь разделен на две версии:

  • ActiveMQ Classic
  • ActiveMQ Artemis

ActiveMQ Classic здесь — это оригинальный ActiveMQ, а ActiveMQ Artemis разработан на основе кода сервера HornetQ, подаренного RedHat.Два кода совершенно разные.Последний поддерживает JMS2.0 и использует асинхронный ввод-вывод на основе Netty, что значительно улучшает производительность., Что еще более удивительно, так это то, что последний не только поддерживает протокол JMS, но также поддерживает протокол AMQP, STOMP и MQTT.Можно сказать, что последний имеет очень богатый игровой процесс.

Поэтому при его использовании рекомендуется выбирать ActiveMQ Artemis напрямую.

1.2.2 RabbitMQ

RabbitMQ является наиболее важным продуктом в системе AMQP. Он разработан и реализован на основе языка Erlang. По оценкам, многие люди были замучены установкой RabbitMQ. Сонг Гэ рекомендует устанавливать RabbitMQ напрямую с помощью Docker, что избавляет от беспокойства и усилия (официальная учетная запись отвечает на докер в фоновом режиме и имеет учебник).

RabbitMQ поддерживает AMQP, XMPP, SMTP, STOMP и другие протоколы с мощными функциями, подходящими для разработки на уровне предприятия.

Давайте взглянем на структурную схему RabbitMQ:

Что касается RabbitMQ, Song Ge недавно опубликовал более десяти руководств, поэтому я не буду здесь многословен.

1.2.3 RocketMQ

RocketMQ — промежуточное программное обеспечение для распределенных сообщений с открытым исходным кодом от Alibaba.Первоначальное название было Metaq.Он был переименован в RocketMQ с версии 3.0.Это набор MQ, реализованный Alibaba с использованием языка Java со ссылкой на идею дизайна Kafka. RocketMQ интегрирует несколько продуктов MQ (Notify, Metaq) в Alibaba, поддерживает только основные функции, удаляет все другие зависимости времени выполнения и обеспечивает простейшие основные функции.На этой основе он сотрудничает с другими продуктами Alibaba с открытым исходным кодом для реализации MQ в различных сценариях. Архитектура в настоящее время в основном используется для систем торговли ордерами.

RocketMQ имеет следующие характеристики:

  • Строгий порядок сообщений гарантируется.
  • Обеспечивает фильтрацию сообщений.
  • Обеспечивает расширенный режим извлечения сообщений.
  • Эффективная горизонтальная масштабируемость абонента.
  • Механизм подписки на сообщения в реальном времени.
  • Возможность накопления сообщений на уровне миллиардов

Для инженеров Java это также часто используемый MQ.

1.2.4 Kafka

Kafka — это платформа обработки потоков с открытым исходным кодом под Apache, написанная на Scala и Java. Kafka — это высокопроизводительная распределенная система обмена сообщениями с публикацией и подпиской, которая обрабатывает потоковые данные для всех действий (просмотр веб-страниц, поиск и другие действия пользователя) потребителей на веб-сайте. Цель Kafka — унифицировать онлайн- и офлайн-обработку сообщений с помощью механизма параллельной загрузки Hadoop и предоставлять сообщения в реальном времени через кластеры.

Кафка имеет следующие особенности:

  • Быстрое сохранение: благодаря последовательному чтению и записи диска и механизму нулевого копирования сохранение сообщений может выполняться с системными издержками O(1).
  • Высокая пропускная способность: на обычном сервере может быть достигнута пропускная способность 10 Вт/с.
  • Высокое накопление: потребители в рамках темы могут находиться в автономном режиме в течение длительного времени, а объем накопления сообщений велик.
  • Полностью распределенная система: Брокер, Производитель и Потребитель изначально и автоматически поддерживают распределение, а более сложная балансировка нагрузки может быть автоматически достигнута с помощью Zookeeper.
  • Поддерживает параллельную загрузку данных Hadoop.

В разработке больших данных вы можете часто вступать в контакт с Kafka, так же как и в Java-разработке, но относительно реже.

1.2.5 ZeroMQ

ZeroMQ известна как самая быстрая система очередей сообщений.Она специально разработана для сценариев с высокой пропускной способностью/малой задержкой и часто используется в финансовых приложениях, ориентируясь на сценарии передачи данных в реальном времени. ZeroMQ — это не отдельная служба, а встроенная библиотека.Он инкапсулирует такие функции, как сетевое взаимодействие, очередь сообщений, планирование потоков и т. д., и предоставляет краткий API для верхнего уровня.Приложение реализует высокопроизводительное сетевое взаимодействие, загружая библиотеку файлы и вызов функций API.

Возможности ZeroMQ:

  • Модель очереди без блокировки: для канала обмена данными между межпотоковыми взаимодействиями (клиент и сеанс) используется алгоритм очереди без блокировки CAS, асинхронные события регистрируются на обоих концах канала, а сообщения считываются или записываются. в канал. , события чтения и записи будут запускаться автоматически.
  • Алгоритм пакетной обработки: адаптивная оптимизация выполняется для пакетных сообщений, которые могут получать и отправлять сообщения пакетами.
  • Привязка потоков в многоядерном режиме без переключения ЦП: в отличие от традиционного многопоточного параллельного режима, семафора или критической секции, ZeroMQ в полной мере использует преимущества многоядерности, каждое ядро ​​привязано к запуску рабочего потока, избегая конфликт между несколькими потоками Накладные расходы на переключение ЦП.

1.2.6 Прочее

Кроме того, Redis также можно использовать в качестве очереди сообщений.Song Ge также опубликовал статью, в которой рассказывается об использовании Redis в качестве общей очереди сообщений и очереди отложенных сообщений.

1.3 Сравнение

Наконец, давайте сравним промежуточное ПО каждого сообщения с изображением.

Друзья, ответьте на mqpkmq в фоновом режиме официального аккаунта, вы можете получить ссылку на эту таблицу Excel.

2. Страница управления RabbitMQ

Считается, что веб-страница управления RabbitMQ использовалась многими небольшими партнерами. Если вы нажмете на нее, вы, возможно, поймете, что она означает. Однако в духе совершенства Songge все еще хочет освежить детали этой страницы управления.

2.1 Обзор

Во-первых, страница веб-управления выглядит следующим образом:

Сначала есть шесть вкладок:

  1. Обзор: Здесь вы можете получить обзор общей ситуации с RabbitMQ.Если это кластер, вы также можете просмотреть ситуацию с каждым узлом в кластере. Включая информацию о сопоставлении портов RabbitMQ и т. д., можно просмотреть на этой вкладке.
  2. Соединения: на этой вкладке находятся производители и потребители, подключенные к RabbitMQ.
  3. Каналы: здесь представлена ​​информация о «канале». Что касается связи между «каналом» и «соединением», Сонге подробно расскажет вам об этом позже.
  4. Обмен: Здесь отображается вся информация об обмене.
  5. Очередь: здесь отображается вся информация об очереди.
  6. Администратор: здесь отображается вся информация о пользователе.

Верхний правый угол — это время обновления страницы, которое по умолчанию обновляется каждые 5 секунд, и отображаются все виртуальные хосты.

Это общая ситуация всей страницы управления, далее мы представим их по очереди.

2.2 Overview

Обзор разделен на следующие функциональные модули:

Они есть:

Итого:

Итоги содержат количество сообщений, которые необходимо использовать, количество сообщений, которые необходимо подтвердить, общее количество сообщений и различные скорости обработки сообщений (скорость отправки, скорость подтверждения, скорость записи на диск и т. д.).

Узлы:

Узлы — это фактически некоторые машины, поддерживающие работу RabbitMQ, которые эквивалентны узлам кластера.

Нажмите на каждый узел, чтобы просмотреть подробную информацию об узле.

Статистика оттока:

Это не так просто перевести, это показывает скорость создания/закрытия Соединения, Канала и Очереди.

Порты и контексты:

Это показывает информацию о сопоставлении портов и информацию о веб-контексте.

  • 5672 — это коммуникационный порт RabbitMQ.
  • 15672 — это порт страницы веб-администрирования.
  • 25672 — порт связи кластера.

Export definitions && Импорт определений:

Последние два могут импортировать и экспортировать некоторую информацию о конфигурации текущего экземпляра:

2.3 Connections

Основным отображением здесь является информация о подключенном в данный момент RabbitMQ, будь то производитель сообщений или потребитель сообщений, пока он подключен, он будет отображаться здесь.

Обратите внимание, что AMQP 0-9-1 в протоколе относится к номеру версии протокола AMQP.

Другие свойства имеют следующие значения:

  • Имя пользователя: имя пользователя, используемое для текущего подключения.
  • Состояние: состояние текущего соединения: «работает» означает «работает», «бездействует» означает «бездействует».
  • SSL/TLS: указывает, следует ли использовать ssl для подключения.
  • Каналы: общее количество каналов, созданных текущим подключением.
  • От клиента: количество отправленных пакетов в секунду.
  • Клиенту: пакетов, полученных в секунду.

Щелкните имя соединения, чтобы просмотреть сведения о каждом соединении.

В деталях вы можете просмотреть количество каналов и другую подробную информацию для каждого соединения, а также принудительно закрыть соединение.

2.4 Channels

В этом месте отображается информация канала:

Так что же такое канал?

Соединение (IP) может иметь несколько каналов, как показано на рисунке выше, всего есть два соединения, но всего 12 каналов.

Соединение может иметь несколько каналов, и эти несколько каналов реализуются через многопоточность.В общем случае мы создаем в канале очереди, коммутаторы и т. д.

Канал производителя, как правило, немедленно закрывается, потребитель всегда слушает, и канал почти всегда будет существовать.

Значения вышеуказанных параметров следующие:

  • Канал: название канала.
  • Имя пользователя: Имя пользователя, используемое для входа на этот канал.
  • Модель: режим подтверждения канала, C означает подтверждение, T означает транзакцию.
  • Состояние: Текущее состояние канала, «работает» означает «работает», «бездействует» означает «бездействует».
  • Неподтверждено: общее количество сообщений, подлежащих подтверждению.
  • Prefetch: Prefetch представляет максимальное количество неподтвержденных сообщений, которое может выдержать каждый потребитель. Короче говоря, он используется для указания, сколько сообщений потребитель может получить от RabbitMQ за раз и кэшировать их в потребителе. RabbitMQ прекратит доставку новых сообщений потребителю, пока не отправит сообщение, которое будет подтверждено. В общем, потребители несут ответственность за непрерывную обработку сообщений, постоянное подтверждение, а затем до тех пор, пока количество неподтвержденных сообщений меньше, чем количество получателей предварительной выборки *, RabbitMQ будет продолжать доставлять сообщения.
  • Unacker: общее количество сообщений, которые необходимо подтвердить.
  • публикация: скорость, с которой производитель сообщений отправляет сообщения.
  • подтверждение: скорость, с которой производитель сообщения подтверждает сообщение.
  • unroutable (drop): указывает, что сообщение не было получено и было удалено.
  • доставить/получить: скорость, с которой потребители сообщений получают сообщения.
  • ack: скорость, с которой потребители сообщений подтверждают сообщения.

2.5 Exchange

В этом месте отображается информация о коммутаторе:

Здесь будет отображаться различная информация о переключателе.

Тип указывает тип переключателя.

Характеристики имеют два значения D и I.

D означает, что переключатель является постоянным, и свойства переключателя хранятся внутри сервера.Когда сервер MQ неожиданно или отключается, при перезапуске RabbitMQ нет необходимости вручную или выполнять код для создания переключателя. переключатель будет создан автоматически, что эквивалентно постоянному существованию.

Я имею в виду, что этот обмен не может использоваться производителями сообщений для отправки сообщений, а используется только для привязки между обменами.

Message rate in представляет скорость, с которой приходят сообщения. Скорость исходящих сообщений указывает скорость, с которой выходят сообщения.

нажмите нижеAdd a new exchangeМожно создать новый коммутатор.

2.6 Queue

Эта вкладка используется для отображения очереди сообщений:

Значения каждого из них следующие:

  • Имя: указывает имя очереди сообщений.
  • Тип: Указывает тип очереди сообщений.Помимо классического на рисунке выше, есть еще один тип сообщения, Quorum. Два отличия заключаются в следующем:

  • Функции: указывает функции очереди сообщений, а D представляет постоянство очереди сообщений.
  • Состояние: указывает состояние текущей очереди, «работает» означает «работает», «неактивный» означает «бездействует».
  • Готово: указывает общее количество сообщений, которые необходимо использовать.
  • Unacked: указывает общее количество сообщений, на которые нужно ответить.
  • Всего: Указывает общее количество сообщений Ready+Unacked.
  • входящие: указывает скорость входящих сообщений.
  • доставить/получить: указывает скорость, с которой получаются сообщения.
  • ack: указывает скорость подтверждения сообщения.

Щелкните Добавить новую очередь ниже, чтобы добавить новую очередь сообщений.

Щелкните имя каждой очереди сообщений, чтобы войти в очередь сообщений. После входа в очередь сообщений можно выполнить дальнейшие операции с очередью сообщений, например:

  • Привязать очередь сообщений к определенному обмену.
  • отправлять сообщения.
  • Получить сообщение.
  • Переместить сообщение (требуется поддержка плагина).
  • Удалить очередь сообщений.
  • Пустые сообщения в очереди сообщений.
  • ...

Как показано ниже:

2.7 Admin

Вот некоторые операции по управлению пользователями, как показано ниже:

Значение каждого атрибута следующее:

  • Имя: указывает имя пользователя.
  • Теги: Указывает теги ролей, можно выбрать только один.
  • Может получить доступ к виртуальным хостам: указывает виртуальные хосты, которым разрешен вход.
  • Имеет пароль: указывает, установил ли пользователь пароль.

Две общие операции — это управление пользователями и виртуальными хостами.

Нажмите Добавить пользователя ниже, чтобы добавить нового пользователя.При добавлении пользователя вам необходимо установить теги для пользователя, который на самом деле является ролью пользователя, следующим образом:

  • никто:

Не могу получить доступ к плагину управления

  • управление:

Все, что пользователь может делать с AMQP Список виртуальных хостов, на которые вы можете войти через AMQP Просмотр очередей, обменов и привязок на ваших собственных виртуальных хостах Просмотр и закрытие собственных каналов и связей Просмотр «глобальной» статистики о ваших собственных виртуальных хостах, включая активность других пользователей на этих виртуальных хостах.

  • политик:

все, что может сделать руководство Просмотр, создание и удаление политик и параметров собственных виртуальных хостов.

  • мониторинг:

все, что может сделать руководство Список всех виртуальных хостов, включая те, на которые они не могут войти Просмотр подключений и каналов других пользователей Просмотр данных на уровне узла, таких как кластеризация и использование памяти Просмотр реальной глобальной статистики обо всех виртуальных хостах

  • администратор:

все, что может сделать политик и мониторинг Создание и удаление виртуальных хостов Просмотр, создание и удаление пользователей Просмотр разрешений на создание и удаление Закрыть соединения других пользователей

  • самозванец

Эмулятор, не удается войти в консоль администратора.

Кроме того, здесь также можно выполнить операцию виртуального хоста, и в следующих подразделах вы познакомитесь с виртуальным хостом.

3. RabbitMQ семь способов отправки и получения сообщений

В этом разделе я поделюсь с вами семью формами обмена сообщениями RabbitMQ. Посмотри.

В большинстве случаев мы можем использовать RabbitMQ в среде Spring Boot или Spring Cloud, поэтому в этой статье я в основном поделюсь с вами использованием RabbitMQ с точки зрения этих двух аспектов.

3.1 Введение в архитектуру RabbitMQ

Картинка стоит тысячи слов, а именно:

1587705504342

Эта диаграмма включает в себя следующие понятия:

  1. Производитель (издатель): Публикуйте сообщения на бирже (Exchange) в RabbitMQ.
  2. Обмен (Exchange): устанавливает соединение с производителем и получает сообщения от производителя.
  3. Потребитель (Consumer): прослушивает сообщения в Очереди в RabbitMQ.
  4. Очередь: Exchange распределяет сообщения в указанную очередь, а очередь взаимодействует с потребителями.
  5. Маршруты: правила для коммутатора для пересылки сообщений в очереди.

3.2 Подготовка

Как мы все знаем, RabbitMQ — это продукт из лагеря AMQP, Spring Boot предоставляет AMQP зависимость автоматической конфигурации spring-boot-starter-amqp, поэтому сначала создайте проект Spring Boot и добавьте эту зависимость следующим образом:

После успешного создания проекта настройте базовую информацию о подключении RabbitMQ в application.properties следующим образом:

spring.rabbitmq.host=localhost
spring.rabbitmq.username=guest
spring.rabbitmq.password=guest
spring.rabbitmq.port=5672

Затем настройте RabbitMQ.В RabbitMQ все сообщения, отправленные производителями сообщений, будут перераспределяться Exchange, и Exchange будет распределять сообщения по разным очередям в соответствии с разными стратегиями.

На официальном сайте RabbitMQ представлены следующие формы рассылки сообщений:

Вот семь типов, из которых седьмой является подтверждением сообщения.Сообщение подтверждает, что Songge ранее отправляла соответствующие статьи.Портал:

Итак, здесь я в основном представляю первые шесть способов отправки и получения сообщений.

3.3 Обмен сообщениями

3.3.1 Hello World

Хм? Почему здесь нет переключателя? На самом деле это переключатель по умолчанию, нам нужно предоставить производителя, очередь и потребителя. Схема распространения сообщения выглядит следующим образом:

Давайте посмотрим на реализацию кода:

Давайте сначала посмотрим на определение очереди:

@Configuration
public class HelloWorldConfig {

    public static final String HELLO_WORLD_QUEUE_NAME = "hello_world_queue";

    @Bean
    Queue queue1() {
        return new Queue(HELLO_WORLD_QUEUE_NAME);
    }
}

Давайте посмотрим на определение потребителя сообщений:

@Component
public class HelloWorldConsumer {
    @RabbitListener(queues = HelloWorldConfig.HELLO_WORLD_QUEUE_NAME)
    public void receive(String msg) {
        System.out.println("msg = " + msg);
    }
}

Отправка сообщения:

@SpringBootTest
class RabbitmqdemoApplicationTests {

    @Autowired
    RabbitTemplate rabbitTemplate;


    @Test
    void contextLoads() {
        rabbitTemplate.convertAndSend(HelloWorldConfig.HELLO_WORLD_QUEUE_NAME, "hello");
    }

}

В настоящее время фактически используется прямой обмен по умолчанию (DirectExchange).Стратегия маршрутизации DirectExchange заключается в привязке очереди сообщений к DirectExchange.Когда сообщение поступает в DirectExchange, оно будет перенаправлено в то же сообщение.routing keyНапример, в той же очереди, если имя очереди сообщений «hello-queue», то сообщение с ключом маршрутизации «hello-queue» будет получено очередью сообщений.

3.3.2 Work queues

Ситуация такова:

Один производитель, один обмен по умолчанию (DirectExchange), одна очередь, два потребителя, как показано ниже:

Очередь соответствует нескольким получателям.По умолчанию сообщения равномерно распределяются по очереди, а сообщения распределяются по разным получателям. Потребители могут настроить свои возможности параллелизма, чтобы улучшить возможности использования сообщений, или настроить ручные подтверждения, чтобы решить, следует ли использовать сообщение.

Давайте сначала посмотрим на конфигурацию возможностей параллелизма, а именно:

@Component
public class HelloWorldConsumer {
    @RabbitListener(queues = HelloWorldConfig.HELLO_WORLD_QUEUE_NAME)
    public void receive(String msg) {
        System.out.println("receive = " + msg);
    }
    @RabbitListener(queues = HelloWorldConfig.HELLO_WORLD_QUEUE_NAME,concurrency = "10")
    public void receive2(String msg) {
        System.out.println("receive2 = " + msg+"------->"+Thread.currentThread().getName());
    }
}

Видно, что для второго потребителя я настроил параллелизм равным 10. В это время для второго потребителя будет 10 дочерних потоков для одновременного потребления сообщений.

Запустите проект, и вы увидите в общей сложности 11 потребителей в фоновом режиме RabbitMQ.

На этом этапе, если производитель отправит 10 сообщений, все они будут использованы одновременно.

Сообщение отправляется следующим образом:

@SpringBootTest
class RabbitmqdemoApplicationTests {

    @Autowired
    RabbitTemplate rabbitTemplate;


    @Test
    void contextLoads() {
        for (int i = 0; i < 10; i++) {
            rabbitTemplate.convertAndSend(HelloWorldConfig.HELLO_WORLD_QUEUE_NAME, "hello");
        }
    }

}

Журнал потребления сообщений выглядит следующим образом:

Как видите, все сообщения потребляются первым потребителем. Но друзья должны быть внимательны, так бывает не всегда (попробуйте несколько раз, и вы увидите разницу), сообщение также может быть потреблено первым потребителем (просто потому, что второй потребитель имеет десять потоков, запущенных вместе, поэтому второй потребитель потребляет большую часть сообщений).

Конечно, потребители сообщений также могут включить ручное подтверждение, чтобы они могли решить, использовать ли сообщения, отправленные RabbitMQ.Способ настройки ручного подтверждения следующий:

spring.rabbitmq.listener.simple.acknowledge-mode=manual

Код потребления выглядит следующим образом:

@Component
public class HelloWorldConsumer {
    @RabbitListener(queues = HelloWorldConfig.HELLO_WORLD_QUEUE_NAME)
    public void receive(Message message,Channel channel) throws IOException {
        System.out.println("receive="+message.getPayload());
        channel.basicAck(((Long) message.getHeaders().get(AmqpHeaders.DELIVERY_TAG)),true);
    }

    @RabbitListener(queues = HelloWorldConfig.HELLO_WORLD_QUEUE_NAME, concurrency = "10")
    public void receive2(Message message, Channel channel) throws IOException {
        System.out.println("receive2 = " + message.getPayload() + "------->" + Thread.currentThread().getName());
        channel.basicReject(((Long) message.getHeaders().get(AmqpHeaders.DELIVERY_TAG)), true);
    }
}

В этот момент второй потребитель отклоняет все сообщения, а первый потребитель потребляет все сообщения.

Это касается рабочих очередей.

3.3.3 Publish/Subscribe

Давайте посмотрим на модель публикации-подписки, которая выглядит следующим образом:

Один производитель, несколько потребителей, у каждого потребителя своя очередь, производитель отправляет сообщение не напрямую в очередь, а на биржу, каждая очередь привязана к бирже, и сообщение, отправленное производителем, проходит через биржу, достичь очереди и достичь цели получения сообщения несколькими потребителями. Следует отметить, что если сообщение отправляется на Exchange без привязки к очереди, сообщение будет потеряно, поскольку Exchange не имеет возможности хранить сообщения в RabbitMQ, только очередь может хранить сообщения, как показано ниже. ниже:

В данном случае у нас есть четыре переключателя на выбор, а именно:

  • Direct
  • Fanout
  • Topic
  • Header

Позвольте мне привести вам простой пример.

3.3.3.1 Direct

Стратегия маршрутизации DirectExchange заключается в привязке очереди сообщений к DirectExchange. Когда сообщение поступает в DirectExchange, оно будет переадресовано в очередь с тем же ключом маршрутизации, что и сообщение. Например, если имя очереди сообщений «hello- очередь", то ключ маршрутизации "hello-queue" сообщения будут получены этой очередью сообщений. DirectExchange настраивается следующим образом:

@Configuration
public class RabbitDirectConfig {
    public final static String DIRECTNAME = "javaboy-direct";
    @Bean
    Queue queue() {
        return new Queue("hello-queue");
    }
    @Bean
DirectExchange directExchange() {
        return new DirectExchange(DIRECTNAME, true, false);
    }
    @Bean
    Binding binding() {
        return BindingBuilder.bind(queue())
                .to(directExchange()).with("direct");
    }
}
  • Сначала создайте очередь сообщений Queue, а затем создайте объект DirectExchange.Три параметра: имя, действительно ли оно по-прежнему действует после перезапуска и удаляется ли оно, если оно не используется в течение длительного времени.
  • Создайте объект Binding, чтобы связать вместе Exchange и Queue.
  • Конфигурация двух bean-компонентов, DirectExchange и Binding, может быть опущена, то есть, если используется DirectExchange, можно настроить только один экземпляр Queue.

Давайте снова посмотрим на потребителя:

@Component
public class DirectReceiver {
    @RabbitListener(queues = "hello-queue")
    public void handler1(String msg) {
        System.out.println("DirectReceiver:" + msg);
    }
}

Метод, указанный в аннотации @RabbitListener, является методом потребления сообщения, а параметр метода — полученным сообщением. Затем добавьте объект RabbitTemplate в класс модульного теста для отправки сообщения следующим образом:

@RunWith(SpringRunner.class)
@SpringBootTest
public class RabbitmqApplicationTests {
    @Autowired
    RabbitTemplate rabbitTemplate;
    @Test
    public void directTest() {
        rabbitTemplate.convertAndSend("hello-queue", "hello direct!");
    }
}

Окончательный результат выполнения следующий:

3.3.3.2 Fanout

Стратегия обмена данными FanoutExchange заключается в пересылке всех сообщений, поступающих на FanoutExchange, во все привязанные к нему очереди. В этой стратегии ключ маршрутизации не будет иметь никакого эффекта. FanoutExchange настраивается следующим образом:

@Configuration
public class RabbitFanoutConfig {
    public final static String FANOUTNAME = "sang-fanout";
    @Bean
    FanoutExchange fanoutExchange() {
        return new FanoutExchange(FANOUTNAME, true, false);
    }
    @Bean
    Queue queueOne() {
        return new Queue("queue-one");
    }
    @Bean
    Queue queueTwo() {
        return new Queue("queue-two");
    }
    @Bean
    Binding bindingOne() {
        return BindingBuilder.bind(queueOne()).to(fanoutExchange());
    }
    @Bean
    Binding bindingTwo() {
        return BindingBuilder.bind(queueTwo()).to(fanoutExchange());
    }
}

Здесь сначала создается FanoutExchange, и значение параметров такое же, как и при создании DirectExchange, затем создаются две Очереди, а затем обе Очереди привязываются к FanoutExchange. Затем создайте двух потребителей следующим образом:

@Component
public class FanoutReceiver {
    @RabbitListener(queues = "queue-one")
    public void handler1(String message) {
        System.out.println("FanoutReceiver:handler1:" + message);
    }
    @RabbitListener(queues = "queue-two")
    public void handler2(String message) {
        System.out.println("FanoutReceiver:handler2:" + message);
    }
}

Два потребителя потребляют сообщения в двух очередях сообщений соответственно, а затем отправляют сообщения в модульном тесте следующим образом:

@RunWith(SpringRunner.class)
@SpringBootTest
public class RabbitmqApplicationTests {
    @Autowired
    RabbitTemplate rabbitTemplate;
    @Test
    public void fanoutTest() {
        rabbitTemplate
        .convertAndSend(RabbitFanoutConfig.FANOUTNAME, 
                null, "hello fanout!");
    }
}

Обратите внимание, что это не требуется при отправке сообщения здесьroutingkey, указавexchangeВот и все,routingkeyможет напрямую передатьnull.

Окончательный журнал выполнения выглядит следующим образом:

3.3.3.3 Topic

TopicExchange — более сложная, но гибкая стратегия маршрутизации. В TopicExchange Queue привязана к TopicExchange с помощью ключа маршрутизации. Когда сообщение поступает в TopicExchange, TopicExchange направляет сообщение в одну или несколько очередей в соответствии с ключом маршрутизации сообщения. TopicExchange настраивается следующим образом:

@Configuration
public class RabbitTopicConfig {
    public final static String TOPICNAME = "sang-topic";
    @Bean
    TopicExchange topicExchange() {
        return new TopicExchange(TOPICNAME, true, false);
    }
    @Bean
    Queue xiaomi() {
        return new Queue("xiaomi");
    }
    @Bean
    Queue huawei() {
        return new Queue("huawei");
    }
    @Bean
    Queue phone() {
        return new Queue("phone");
    }
    @Bean
    Binding xiaomiBinding() {
        return BindingBuilder.bind(xiaomi()).to(topicExchange())
                .with("xiaomi.#");
    }
    @Bean
    Binding huaweiBinding() {
        return BindingBuilder.bind(huawei()).to(topicExchange())
                .with("huawei.#");
    }
    @Bean
    Binding phoneBinding() {
        return BindingBuilder.bind(phone()).to(topicExchange())
                .with("#.phone.#");
    }
}
  • Сначала создайте TopicExchange с теми же параметрами, что и раньше. Затем создайте три очереди, первая очередь используется для хранения сообщений, связанных с «xiaomi», вторая очередь используется для хранения сообщений, связанных с «huawei», а третья очередь используется для хранения сообщений, связанных с «телефоном».
  • Привяжите три очереди к TopicExchange соответственно. «xiaomi.#» в первой привязке указывает, что ключ маршрутизации сообщения, который начинается с «xiaomi», будет направлен в очередь с именем «xiaomi». «huawei.#» в первой привязке указывает, что ключ маршрутизации сообщения, начинающегося с «huawei», будет перенаправлен в очередь с именем «huawei», а «#.phone.#» в третьей привязке указывает, что все сообщения, содержащие «телефон» в routingkey будет перенаправлен в очередь с именем «телефон».

Затем создайте трех потребителей для трех очередей следующим образом:

@Component
public class TopicReceiver {
    @RabbitListener(queues = "phone")
    public void handler1(String message) {
        System.out.println("PhoneReceiver:" + message);
    }
    @RabbitListener(queues = "xiaomi")
    public void handler2(String message) {
        System.out.println("XiaoMiReceiver:"+message);
    }
    @RabbitListener(queues = "huawei")
    public void handler3(String message) {
        System.out.println("HuaWeiReceiver:"+message);
    }
}

Затем отправьте сообщение в модульном тесте следующим образом:

@RunWith(SpringRunner.class)
@SpringBootTest
public class RabbitmqApplicationTests {
    @Autowired
    RabbitTemplate rabbitTemplate;
    @Test
    public void topicTest() {
        rabbitTemplate.convertAndSend(RabbitTopicConfig.TOPICNAME,
"xiaomi.news","小米新闻..");
        rabbitTemplate.convertAndSend(RabbitTopicConfig.TOPICNAME,
"huawei.news","华为新闻..");
        rabbitTemplate.convertAndSend(RabbitTopicConfig.TOPICNAME,
"xiaomi.phone","小米手机..");
        rabbitTemplate.convertAndSend(RabbitTopicConfig.TOPICNAME,
"huawei.phone","华为手机..");
        rabbitTemplate.convertAndSend(RabbitTopicConfig.TOPICNAME,
"phone.news","手机新闻..");
    }
}

Согласно конфигурации в RabbitTopicConfig, первое сообщение будет направлено в очередь с именем «xiaomi», второе сообщение будет направлено в очередь с именем «huawei», а третье сообщение будет направлено в очередь с именем «huawei». очередь с именами «xiaomi» и «телефон», четвертое сообщение будет перенаправлено в очередь с именем «huawei» и с именем «телефон», а последнее сообщение будет направлено в очередь с именем «телефон» в очереди.

3.3.3.4 Header

HeadersExchange — это менее используемая стратегия маршрутизации. HeadersExchange направляет сообщения в разные очереди в соответствии с их заголовками. Эта стратегия также не зависит от ключа маршрутизации. Конфигурация выглядит следующим образом:

@Configuration
public class RabbitHeaderConfig {
    public final static String HEADERNAME = "javaboy-header";
    @Bean
    HeadersExchange headersExchange() {
        return new HeadersExchange(HEADERNAME, true, false);
    }
    @Bean
    Queue queueName() {
        return new Queue("name-queue");
    }
    @Bean
    Queue queueAge() {
        return new Queue("age-queue");
    }
    @Bean
    Binding bindingName() {
        Map<String, Object> map = new HashMap<>();
        map.put("name", "sang");
        return BindingBuilder.bind(queueName())
                .to(headersExchange()).whereAny(map).match();
    }
    @Bean
    Binding bindingAge() {
        return BindingBuilder.bind(queueAge())
                .to(headersExchange()).where("age").exists();
    }
}

Большая часть конфигурации здесь такая же, как и предыдущая.Различие в основном отражается в конфигурации Binding.В первом методе bindingName, гдеAny означает, что до тех пор, пока в заголовке сообщения есть один заголовок, соответствующий ключу /value в карте, сообщение будет отправлено в message.Routing to the Queue с именем «имя-очередь», где также может использоваться метод All, указывающий, что все заголовки сообщения должны совпадать. whereAny и whereAll на самом деле соответствуют свойству, называемому x-match. Конфигурация в bindingAge означает, что пока заголовок сообщения содержит возраст, независимо от значения возраста, сообщение будет перенаправлено в очередь с именем «age-queue».

Затем создайте двух потребителей сообщений:

@Component
public class HeaderReceiver {
    @RabbitListener(queues = "name-queue")
    public void handler1(byte[] msg) {
        System.out.println("HeaderReceiver:name:"
                + new String(msg, 0, msg.length));
    }
    @RabbitListener(queues = "age-queue")
    public void handler2(byte[] msg) {
        System.out.println("HeaderReceiver:age:"
                + new String(msg, 0, msg.length));
    }
}

Обратите внимание, что параметры здесь принимаются в виде массивов байтов. Затем создайте метод отправки сообщения в модульном тесте, где отправка сообщения также не зависит от ключа маршрутизации, как показано ниже:

@RunWith(SpringRunner.class)
@SpringBootTest
public class RabbitmqApplicationTests {
    @Autowired
    RabbitTemplate rabbitTemplate;
    @Test
    public void headerTest() {
        Message nameMsg = MessageBuilder
                .withBody("hello header! name-queue".getBytes())
                .setHeader("name", "sang").build();
        Message ageMsg = MessageBuilder
                .withBody("hello header! age-queue".getBytes())
                .setHeader("age", "99").build();
        rabbitTemplate.send(RabbitHeaderConfig.HEADERNAME, null, ageMsg);
        rabbitTemplate.send(RabbitHeaderConfig.HEADERNAME, null, nameMsg);
    }
}

Здесь создаются два сообщения. У двух сообщений разные заголовки. Сообщения с разными заголовками будут отправлены в разные очереди.

Окончательный эффект выполнения выглядит следующим образом:

3.3.4 Routing

В этом случае:

Один производитель, один коммутатор, две очереди и два потребителя.После создания Exchange производитель привязывает соответствующую очередь по RoutingKey и при отправке сообщения может указать конкретный RoutingKey сообщения.

Как показано ниже:

Это для маршрутизации сообщений в соответствии с ключом маршрутизации.Я не буду приводить здесь пример.Вы можете обратиться к сводке 3.3.1.

3.3.5 Topics

В этом случае:

Один производитель, один обмен, две очереди, два потребителя, производитель создает Обмен Темы и привязывает его к очереди. Эту привязку можно сделать с помощью*и#ключевое слово, указатьRoutingKeyСодержание, обратите внимание на формат при написанииxxx.xxx.xxxнаписать.

Как показано ниже:

Я не буду приводить пример этого, в предыдущем разделе 3.3.3 уже приводился пример, поэтому я не буду повторяться.

3.3.6 RPC

RPC - это форма отправки и получения сообщений. Сонг Гэ только что написал статью два дня назад и представил ее всем. Я не буду здесь много говорить. Портал:

3.3.7 Publisher Confirms

Этот вид отправки подтверждает, что Сун Гэ ранее писал статьи на эту тему, и портал:

4. RabbitMQ реализует RPC

Когда дело доходит до RPC (протокол удаленного вызова процедур), в умах друзей всплывают оценки RESTful API, Dubbo, WebService, Java RMI, CORBA и т. д.

На самом деле RabbitMQ также предоставляет нам функции RPC, и он очень прост в использовании.

Song Ge использует простой случай, чтобы поделиться с вами тем, как Spring Boot + RabbitMQ реализует простой вызов RPC.

Уведомление

Что касается реализации RPC-вызовов в RabbitMQ, у некоторых друзей могут быть некоторые недоразумения, думая, что это непросто? Создайте две очереди сообщений queue_1 и queue_2.Сначала клиент отправляет сообщение в queue_1, сервер слушает сообщение в queue_1 и обрабатывает его после получения, после завершения обработки сервер отправляет сообщение в очередь queue_2, а затем клиент слушает queue_2.Сообщение в очереди, так что известен результат обработки сервера.

Этот метод не невозможен, это немного хлопотно! RabbitMQ предоставляет готовые решения, которые можно использовать напрямую, что очень удобно. Дальше будем учиться вместе.

4.1 Архитектура

Давайте сначала посмотрим на простую архитектурную схему:

Эта картинка делает проблему очень ясной:

  1. Во-первых, клиент отправляет сообщение.По сравнению с обычными сообщениями, это сообщение имеет еще два ключевых содержимого: одно — корреляция_ид, которое представляет собой уникальный идентификатор сообщения, а другое — reply_to, которое представляет имя очереди ответов на сообщения. .
  2. Сервер получает сообщения из очереди отправки сообщений и обрабатывает соответствующую бизнес-логику, после завершения обработки отправляет результат обработки в очередь обратного вызова, указанную в параметре reply_to.
  3. Клиент читает сообщение из очереди обратного вызова и может знать, как выглядит выполнение сообщения.

Эта ситуация на самом деле очень подходит для обработки асинхронных вызовов.

4.2 Практика

Далее давайте посмотрим, как это работает на конкретном примере.

4.2.1 Развитие клиента

Во-первых, давайте создадим проект Spring Boot с именем производителя, в качестве производителя сообщений, добавим зависимости web и rabbitmq при создании, как показано ниже:

После успешного создания проекта сначала настройте базовую информацию RabbitMQ в application.properties следующим образом:

spring.rabbitmq.host=localhost
spring.rabbitmq.port=5672
spring.rabbitmq.username=guest
spring.rabbitmq.password=guest
spring.rabbitmq.publisher-confirm-type=correlated
spring.rabbitmq.publisher-returns=true

Первые четыре строки этой конфигурации понятны, поэтому я не буду их повторять. Следующие две строки: Во-первых, настраиваем метод подтверждения сообщения. Подтверждаем его через корреляцию. быть включены в будущие сообщения.Только через корреляция_id мы можем связать отправленное сообщение с возвращаемым значением. Последняя строка конфигурации — включить возврат ошибки отправки.

Затем мы предоставляем класс конфигурации следующим образом:

/**
 * @author 江南一点雨
 * @微信公众号 江南一点雨
 * @网站 http://www.itboyhub.com
 * @国际站 http://www.javaboy.org
 * @微信 a_java_boy
 * @GitHub https://github.com/lenve
 * @Gitee https://gitee.com/lenve
 */
@Configuration
public class RabbitConfig {

    public static final String RPC_QUEUE1 = "queue_1";
    public static final String RPC_QUEUE2 = "queue_2";
    public static final String RPC_EXCHANGE = "rpc_exchange";

    /**
     * 设置消息发送RPC队列
     */
    @Bean
    Queue msgQueue() {
        return new Queue(RPC_QUEUE1);
    }

    /**
     * 设置返回队列
     */
    @Bean
    Queue replyQueue() {
        return new Queue(RPC_QUEUE2);
    }

    /**
     * 设置交换机
     */
    @Bean
    TopicExchange exchange() {
        return new TopicExchange(RPC_EXCHANGE);
    }

    /**
     * 请求队列和交换器绑定
     */
    @Bean
    Binding msgBinding() {
        return BindingBuilder.bind(msgQueue()).to(exchange()).with(RPC_QUEUE1);
    }

    /**
     * 返回队列和交换器绑定
     */
    @Bean
    Binding replyBinding() {
        return BindingBuilder.bind(replyQueue()).to(exchange()).with(RPC_QUEUE2);
    }


    /**
     * 使用 RabbitTemplate发送和接收消息
     * 并设置回调队列地址
     */
    @Bean
    RabbitTemplate rabbitTemplate(ConnectionFactory connectionFactory) {
        RabbitTemplate template = new RabbitTemplate(connectionFactory);
        template.setReplyAddress(RPC_QUEUE2);
        template.setReplyTimeout(6000);
        return template;
    }


    /**
     * 给返回队列设置监听器
     */
    @Bean
    SimpleMessageListenerContainer replyContainer(ConnectionFactory connectionFactory) {
        SimpleMessageListenerContainer container = new SimpleMessageListenerContainer();
        container.setConnectionFactory(connectionFactory);
        container.setQueueNames(RPC_QUEUE2);
        container.setMessageListener(rabbitTemplate(connectionFactory));
        return container;
    }
}

В этом классе конфигурации мы настраиваем очередь отправки сообщений msgQueue и очередь возврата сообщений replyQueue соответственно, а затем привязываем эти две очереди к обмену сообщениями. Это нормальная работа RabbitMQ, тут и говорить нечего.

В Spring Boot за отправку сообщений отвечает инструмент RabbitTemplate, По умолчанию система автоматически предоставляет этот инструмент, но здесь нам нужно перенастроить инструмент.В основном, чтобы добавить очередь возврата для отправки сообщений, и, наконец, нам нужно установить прослушиватель для очереди возврата.

Хорошо, тогда мы можем начать отправлять определенные сообщения:

/**
 * @author 江南一点雨
 * @微信公众号 江南一点雨
 * @网站 http://www.itboyhub.com
 * @国际站 http://www.javaboy.org
 * @微信 a_java_boy
 * @GitHub https://github.com/lenve
 * @Gitee https://gitee.com/lenve
 */
@RestController
public class RpcClientController {

    private static final Logger logger = LoggerFactory.getLogger(RpcClientController.class);

    @Autowired
    private RabbitTemplate rabbitTemplate;

    @GetMapping("/send")
    public String send(String message) {
        // 创建消息对象
        Message newMessage = MessageBuilder.withBody(message.getBytes()).build();

        logger.info("client send:{}", newMessage);

        //客户端发送消息
        Message result = rabbitTemplate.sendAndReceive(RabbitConfig.RPC_EXCHANGE, RabbitConfig.RPC_QUEUE1, newMessage);

        String response = "";
        if (result != null) {
            // 获取已发送的消息的 correlationId
            String correlationId = newMessage.getMessageProperties().getCorrelationId();
            logger.info("correlationId:{}", correlationId);

            // 获取响应头信息
            HashMap<String, Object> headers = (HashMap<String, Object>) result.getMessageProperties().getHeaders();

            // 获取 server 返回的消息 id
            String msgId = (String) headers.get("spring_returned_message_correlation");

            if (msgId.equals(correlationId)) {
                response = new String(result.getBody());
                logger.info("client receive:{}", response);
            }
        }
        return response;
    }
}

Код в этом блоке на самом деле является обычным кодом, я выберу несколько ключевых узлов, чтобы рассказать о них:

  1. Отправка сообщения вызывает метод sendAndReceive, который имеет собственное возвращаемое значение — сообщение, возвращаемое сервером.
  2. В сообщении, возвращаемом сервером, информация заголовка содержит поле spring_returned_message_correlation, которое является корреляцией_идентификатором при отправке сообщения.Посредством корреляции_идентификатора при отправке сообщения и значением поля spring_returned_message_correlation в заголовке возвращенного сообщения мы можем сравнить содержимое возвращаемого сообщения с отправленным. Сообщения связываются вместе, чтобы подтвердить, что возвращенное содержимое относится к отправленному сообщению.

Это разработка всего клиента, по сути, ядром является вызов метода sendAndReceive. Несмотря на простоту вызова, работы по подготовке все же достаточно. Например, если мы не настроим корреляцию в application.properties, в отправленном сообщении не будет Корреляция_id, поэтому невозможно сопоставить возвращаемое содержимое сообщения с содержимым отправленного сообщения.

4.2.2 Разработка сервера

Давайте посмотрим на разработку на стороне сервера.

Во-первых, создайте проект Spring Boot с именем потребитель.Зависимости, добавляемые проектом, соответствуют тем, которые создаются при разработке клиента, поэтому я не буду вдаваться в подробности.

Затем настройте файл конфигурации application.properties.Конфигурация этого файла также совпадает с конфигурацией в клиенте и здесь повторяться не будет.

Затем предоставляется класс конфигурации RabbitMQ. Этот класс конфигурации относительно прост. Просто настройте очередь сообщений и привяжите ее к обмену сообщениями следующим образом:

/**
 * @author 江南一点雨
 * @微信公众号 江南一点雨
 * @网站 http://www.itboyhub.com
 * @国际站 http://www.javaboy.org
 * @微信 a_java_boy
 * @GitHub https://github.com/lenve
 * @Gitee https://gitee.com/lenve
 */
@Configuration
public class RabbitConfig {

    public static final String RPC_QUEUE1 = "queue_1";
    public static final String RPC_QUEUE2 = "queue_2";
    public static final String RPC_EXCHANGE = "rpc_exchange";

    /**
     * 配置消息发送队列
     */
    @Bean
    Queue msgQueue() {
        return new Queue(RPC_QUEUE1);
    }

    /**
     * 设置返回队列
     */
    @Bean
    Queue replyQueue() {
        return new Queue(RPC_QUEUE2);
    }

    /**
     * 设置交换机
     */
    @Bean
    TopicExchange exchange() {
        return new TopicExchange(RPC_EXCHANGE);
    }

    /**
     * 请求队列和交换器绑定
     */
    @Bean
    Binding msgBinding() {
        return BindingBuilder.bind(msgQueue()).to(exchange()).with(RPC_QUEUE1);
    }

    /**
     * 返回队列和交换器绑定
     */
    @Bean
    Binding replyBinding() {
        return BindingBuilder.bind(replyQueue()).to(exchange()).with(RPC_QUEUE2);
    }
}

Наконец, давайте посмотрим на потребление сообщений:

@Component
public class RpcServerController {
    private static final Logger logger = LoggerFactory.getLogger(RpcServerController.class);
    @Autowired
    private RabbitTemplate rabbitTemplate;

    @RabbitListener(queues = RabbitConfig.RPC_QUEUE1)
    public void process(Message msg) {
        logger.info("server receive : {}",msg.toString());
        Message response = MessageBuilder.withBody(("i'm receive:"+new String(msg.getBody())).getBytes()).build();
        CorrelationData correlationData = new CorrelationData(msg.getMessageProperties().getCorrelationId());
        rabbitTemplate.sendAndReceive(RabbitConfig.RPC_EXCHANGE, RabbitConfig.RPC_QUEUE2, response, correlationData);
    }
}

Логика здесь проще:

  1. Сервер сначала получает сообщение и распечатывает его.
  2. Сервер извлекает Correlation_id из исходного сообщения.
  3. Сервер вызывает метод sendAndReceive для отправки сообщения в очередь RPC_QUEUE2 с параметром корреляции_id.

После отправки сообщения с сервера клиент получит результат, возвращенный сервером.

Хорошо, все готово.

4.2.3 Тестирование

Далее проведем простой тест.

Сначала запустите RabbitMQ.

Затем запустите производителя и потребителя соответственно, а затем вызовите интерфейс производителя в postman для тестирования следующим образом:

Видно, что обратная информация от сервера получена.

Давайте посмотрим на рабочий журнал производителя:

Видно, что после отправки сообщения информация, возвращаемая потребителем, также принимается.

Как видите, потребитель также получает сообщение от клиента.

5. Срок действия сообщения RabbitMQ

Срок действия сообщений в RabbitMQ истекает, если они не используются в течение длительного времени? У друзей, которые пользовались RabbitMQ, могут возникнуть такие вопросы, Сонг Гэ приехал, чтобы обсудить эту проблему со всеми.

5.1 По умолчанию

Сначала давайте рассмотрим случай по умолчанию.

По умолчанию срок действия сообщений не истекает, то есть когда мы отправляем сообщения в будние дни, если мы не установим какие-либо соответствующие параметры для истечения срока действия сообщения, то срок действия сообщения не истечет, даже если сообщение не будет использовано, оно всегда будет сохранено .в очереди.

В этом случае мне не нужно демонстрировать конкретный код.Предыдущие статьи Song Ge в основном похожи на то, что касается RabbitMQ.

5.2 TTL

TTL (Time-To-Live), время, в течение которого сообщение сохраняется, то есть срок действия сообщения. Если мы хотим, чтобы сообщения имели время жизни, мы можем добиться этого, установив TTL. Если время жизни сообщения превышает TTL и сообщение не было отправлено, сообщение становится死信死信а также死信队列, Сонг Гэ представит вам позже.

Есть два разных способа установить TTL:

  1. При объявлении очереди мы можем установить срок действия сообщения в свойствах очереди, чтобы все сообщения, поступающие в очередь, имели одинаковый срок действия.
  2. Срок действия сообщения устанавливается при отправке сообщения, поэтому разные сообщения имеют разные сроки действия.

А если оба установлены?

что короче.

Когда мы устанавливаем срок действия сообщения, сообщение будет удалено из очереди по истечении срока его действия (введите очередь недоставленных сообщений, которая будет такой же в следующем тексте и не будет помечена), но есть некоторые отличия в соответствующие сроки удаления двух методов:

  1. Для первого метода, когда очередь сообщений устанавливает время истечения срока действия, сообщение будет удалено по истечении срока его действия, поскольку сообщение сохраняется в очереди сообщений после входа в RabbitMQ, а заголовок очереди — это самое раннее сообщение с истечением срока действия, поэтому Только RabbitMQ. Запланированная задача требуется для сканирования сообщений с истекшим сроком действия из заголовка и их непосредственного удаления, если таковые имеются.
  2. Для второго метода, когда срок действия сообщения истечет, оно не будет удалено сразу, а будет удалено, когда сообщение должно быть доставлено потребителю, поскольку во втором методе время истечения срока действия каждого сообщения отличается. истек, для этого необходимо просмотреть все сообщения в очереди.Когда сообщений много, производительность возрастает.Поэтому для второго метода сообщение удаляется, когда оно должно быть доставлено в потребитель.

После введения TTL давайте взглянем на конкретное использование.

Далее весь код Songge будет брать AMPQ, упакованный в Spring Boot, в качестве примера для объяснения.

5.2.1 Срок действия одного сообщения

Давайте сначала посмотрим на время истечения срока действия одного сообщения.

Сначала создайте проект Spring Boot и введите зависимости Web и RabbitMQ следующим образом:

Затем настройте информацию о подключении RabbitMQ в application.properties следующим образом:

spring.rabbitmq.host=127.0.0.1
spring.rabbitmq.port=5672
spring.rabbitmq.username=guest
spring.rabbitmq.password=guest
spring.rabbitmq.virtual-host=/

Затем немного настройте очередь сообщений:

@Configuration
public class QueueConfig {

    public static final String JAVABOY_QUEUE_DEMO = "javaboy_queue_demo";
    public static final String JAVABOY_EXCHANGE_DEMO = "javaboy_exchange_demo";
    public static final String HELLO_ROUTING_KEY = "hello_routing_key";

    @Bean
    Queue queue() {
        return new Queue(JAVABOY_QUEUE_DEMO, true, false, false);
    }

    @Bean
    DirectExchange directExchange() {
        return new DirectExchange(JAVABOY_EXCHANGE_DEMO, true, false);
    }

    @Bean
    Binding binding() {
        return BindingBuilder.bind(queue())
                .to(directExchange())
                .with(HELLO_ROUTING_KEY);
    }
}

Этот класс конфигурации в основном делает три вещи: настраивает очередь сообщений, настраивает коммутатор и связывает их вместе.

  1. Сначала настройте очередь сообщений, создайте новую очередь: первый параметр — это имя очереди сообщений; второй параметр указывает, является ли сообщение постоянным; третий параметр указывает, является ли очередь сообщений монопольной, обычно мы устанавливаем для него значение false, т. е. является неэксклюзивным; четвертый параметр указывает, что если в очереди нет подписанных потребителей, очередь будет автоматически удалена, что обычно применимо к временным очередям.
  2. Настройте коммутатор DirectExchange.
  3. Свяжите обмены и очереди вместе.

Эта конфигурация должна быть очень простой, здесь нечего объяснять, здесь есть эксклюзивность, Сун Гэ сказал немного больше здесь:

Что касается эксклюзивности, если для нее установлено значение true, доступ к очереди сообщений может получить только создавшее ее соединение, а другие соединения не могут получить доступ к очереди сообщений.Если вы попытаетесь повторно объявить или получить доступ к эксклюзивной очереди в другом сообщит об ошибке блокировки ресурса. С другой стороны, для монопольных очередей при разрыве соединения очередь сообщений также будет автоматически удалена (независимо от того, объявлена ​​очередь постоянной или нет).

Затем предоставьте интерфейс отправки сообщений следующим образом:

@RestController
public class HelloController {
    @Autowired
    RabbitTemplate rabbitTemplate;

    @GetMapping("/hello")
    public void hello() {
        Message message = MessageBuilder.withBody("hello javaboy".getBytes())
                .setExpiration("10000")
                .build();
        rabbitTemplate.convertAndSend(QueueConfig.JAVABOY_QUEUE_DEMO, message);
    }
}

При создании объекта Message мы можем установить время истечения срока действия сообщения Здесь время истечения срока действия сообщения установлено на 10 секунд.

Вот и все!

Далее запускаем проект и тестируем отправку сообщения. Когда сообщение отправлено успешно, поскольку потребителя нет, сообщение не будет использовано. Откройте страницу управления RabbitMQ и перейдите на вкладку «Очереди».Через 10 секунд мы обнаружим, что сообщение исчезло:

Это просто!

Чтобы установить время истечения срока действия для одного сообщения, вы можете установить срок действия сообщения при его отправке.

5.2.2 Срок действия сообщения очереди

Установите время истечения срока действия сообщения для очереди следующим образом:

@Bean
Queue queue() {
    Map<String, Object> args = new HashMap<>();
    args.put("x-message-ttl", 10000);
    return new Queue(JAVABOY_QUEUE_DEMO, true, false, false, args);
}

После завершения настройки модифицируем логику отправки сообщения следующим образом:

@RestController
public class HelloController {
    @Autowired
    RabbitTemplate rabbitTemplate;

    @GetMapping("/hello")
    public void hello() {
        Message message = MessageBuilder.withBody("hello javaboy".getBytes())
                .build();
        rabbitTemplate.convertAndSend(QueueConfig.JAVABOY_QUEUE_DEMO, message);
    }
}

Как видите, сообщение можно отправить в обычном режиме, не устанавливая время истечения срока действия сообщения.

ОК, запустите проект, отправьте сообщение для тестирования. Просмотрите страницу управления RabbitMQ следующим образом:

Видно, что свойства Feature очереди сообщений равны D и TTL, D указывает на то, что сообщение в очереди сообщений является постоянным, а TTL указывает на истечение срока действия сообщения.

Обновите страницу через 10 секунд и обнаружите, что количество сообщений вернулось к 0.

Это необходимо для установки времени истечения срока действия сообщения для очереди сообщений.После установки все сообщения, поступающие в очередь, имеют срок действия.

5.2.3 Особые случаи

Существует также особый случай, который заключается в том, чтобы установить время истечения TTL сообщения равным 0, что означает, что если сообщение не может быть использовано немедленно, оно будет немедленно отброшено.Эта функция может частично заменить немедленный параметр, ранее поддерживаемый RabbitMQ3. .0 Это связано с тем, что непосредственный параметр будет иметь метод basic.return для возврата тела сообщения в случае сбоя доставки (эта функция может быть реализована с использованием очереди недоставленных сообщений).

Конкретный код Songge не будет демонстрировать, это должно быть проще.

5.3 Очередь недоставленных сообщений

Некоторые друзья не могут не спросить, куда делись удаленные сообщения? Его действительно удалили? Ни ни ни! Это связано с очередью недоставленных сообщений.Далее, давайте взглянем на очередь недоставленных сообщений.

5.3.1 Переключатель недоставленных сообщений

Обмен недоставленными письмами, Dead-Letter-Exchange — это DLX.

Обмен недоставленными письмами используется для получения недоставленных сообщений (Dead Message), так что же такое недоставленное сообщение? Есть несколько ситуаций, в которых обычное сообщение становится недоставленным сообщением:

  • Сообщение отклонено (Basic.Reject/Basic.Nack) с параметром requeue, установленным на false
  • срок действия сообщения истек
  • очередь достигает максимальной длины

Когда сообщение становится недоставленным сообщением в очереди, оно будет отправлено в DLX в это время, а очередь сообщений, связанная с DLX, называется очередью недоставленных сообщений.

DLX тоже обычный свитч по сути.Мы можем указать DLX для любой очереди.Когда в очереди есть мертвое письмо, RabbitMQ автоматически опубликует мертвое письмо в DLX, которое затем будет перенаправлено на другую привязку.в очереди DLX (т. е. очередь недоставленных сообщений).

5.3.2 Очередь недоставленных сообщений

Это легко понять, очередь, связанная с переключателем недоставленных сообщений, является очередью недоставленных сообщений.

5.3.3 Практика

Давайте рассмотрим простой пример.

Сначала создадим переключатель недоставленных сообщений, затем создадим очередь недоставленных сообщений и свяжем переключатель недоставленных сообщений и очередь недоставленных сообщений вместе:

public static final String DLX_EXCHANGE_NAME = "dlx_exchange_name";
public static final String DLX_QUEUE_NAME = "dlx_queue_name";
public static final String DLX_ROUTING_KEY = "dlx_routing_key";

/**
 * 配置死信交换机
 *
 * @return
 */
@Bean
DirectExchange dlxDirectExchange() {
    return new DirectExchange(DLX_EXCHANGE_NAME, true, false);
}
/**
 * 配置死信队列
 * @return
 */
@Bean
Queue dlxQueue() {
    return new Queue(DLX_QUEUE_NAME);
}
/**
 * 绑定死信队列和死信交换机
 * @return
 */
@Bean
Binding dlxBinding() {
    return BindingBuilder.bind(dlxQueue())
            .to(dlxDirectExchange())
            .with(DLX_ROUTING_KEY);
}

По сути, это ничем не отличается от обычных обменов и обычных очередей сообщений.

Затем настройте переключатель недоставленных сообщений для очереди сообщений следующим образом:

@Bean
Queue queue() {
    Map<String, Object> args = new HashMap<>();
    //设置消息过期时间
    args.put("x-message-ttl", 0);
    //设置死信交换机
    args.put("x-dead-letter-exchange", DLX_EXCHANGE_NAME);
    //设置死信 routing_key
    args.put("x-dead-letter-routing-key", DLX_ROUTING_KEY);
    return new Queue(JAVABOY_QUEUE_DEMO, true, false, false, args);
}

Всего два параметра:

  • x-dead-letter-exchange: настроить обмен недоставленными письмами.
  • x-dead-letter-routing-key: настройка недоставленных писемrouting_key.

Это настроено.

Сообщения, отправленные в эту очередь сообщений, в будущем будут отправлены в DLX, если возникнут такие проблемы, как отказ, отклонение или истечение срока действия, а затем попадут в очередь сообщений, привязанную к DLX.

Использование очереди недоставленных сообщений ничем не отличается от потребления обычной очереди сообщений:

@RabbitListener(queues = QueueConfig.DLX_QUEUE_NAME)
public void dlxHandle(String msg) {
    System.out.println("dlx msg = " + msg);
}

Это легко~

6. RabbitMQ реализует очередь задержки

Существуют различные запланированные задачи.Обычные запланированные задачи, такие как резервное копирование журнала, мы можем выполнять резервное копирование в 3 часа ночи каждый день.Такого рода запланированные задачи с фиксированным временем можно легко реализовать с помощью выражений cron.Есть также некоторые специальные.Временная задача , покажите всем бомбу замедленного действия в фильме, и она взорвется через 3 минуты.Такую временную задачу непросто описать с помощью cron, потому что время начала не определено, и мы иногда сталкиваемся с похожими потребностями во время разработки.Например:

  • В проектах электронной коммерции, когда мы размещаем заказ, оплата обычно занимает 20 или 30 минут, в противном случае заказ войдет в логику обработки исключений и будет отменен, затем войдет в логику обработки исключений, это можно рассматривать как задержку очередь.
  • Купила умную кастрюлю, из которой можно варить кашу, кладу все ингредиенты в кастрюлю перед уходом на работу, а потом выставляю время и минуты начала варки каши, чтобы после работы выпить ароматную кашу , то эта каша. Инструкция также может рассматриваться как отложенная задача, которая ставится в очередь отсрочки и выполняется, когда время истекло.
  • Система бронирования встреч компании уведомит всех пользователей, которые забронируют встречу, за полчаса до ее начала после того, как встреча будет успешно забронирована.
  • Если заказ на работу по обеспечению безопасности не обрабатывается более 24 часов, компания автоматически свяжется с группой WeChat предприятия, чтобы напомнить соответствующему ответственному лицу.
  • После того, как пользователь разместит заказ на вынос, брат на вынос получит напоминание о том, что тайм-аут истекает, когда до тайм-аута остается 10 минут.
  • ...

Есть много сценариев, в которых нам нужны очереди с задержкой.

В этой статье RabbitMQ используется в качестве примера, чтобы рассказать об игровом процессе отложенных очередей.

В общем, есть два способа реализовать запланированные задачи на RabbitMQ:

  • Используйте механизм истечения срока действия сообщения и частной очереди сообщений, который поставляется с RabbitMQ, для реализации запланированных задач.
  • Использование плагина RabbitMQ RabbitMQ_delayed_message_exchange для реализации запланированных задач является более простым решением.

Давайте рассмотрим два использования по отдельности.

6.1 Использование плагинов

6.1.1 Установите плагин

Сначала нам нужно загрузить плагин rabbitmq_delayed_message_exchange, который является проектом с открытым исходным кодом на GitHub, мы можем загрузить его напрямую:

Выбирайте подходящую вам версию, здесь я выбираю последнюю версию 3.9.0.

После завершения загрузки выполните следующую команду в командной строке, чтобы скопировать загруженный файл в контейнер Docker:

docker cp ./rabbitmq_delayed_message_exchange-3.9.0.ez some-rabbit:/plugins

Здесь первый параметр — это адрес файла на хосте, а второй параметр — это место для копирования в контейнер.

Затем выполните следующую команду, чтобы войти в контейнер RabbitMQ:

docker exec -it some-rabbit /bin/bash

После входа в контейнер выполните следующую команду, чтобы включить плагин:

rabbitmq-plugins enable rabbitmq_delayed_message_exchange

После успешной активации вы также можете использовать следующую команду для просмотра всех установленных подключаемых модулей, чтобы увидеть, есть ли какие-либо подключаемые модули, которые мы только что установили, следующим образом:

rabbitmq-plugins list

Полный процесс выполнения команды выглядит следующим образом:

ОК, после завершения настройки выполняемexitкоманда для выхода из контейнера RabbitMQ. Затем приступайте к кодированию.

6.1.2 Обмен сообщениями

Затем начните отправлять и получать сообщения.

Во-первых, мы создаем проект Spring Boot и вводим зависимости Web и RabbitMQ следующим образом:

После успешного создания проекта настройте базовую информацию RabbitMQ в application.properties следующим образом:

spring.rabbitmq.host=localhost
spring.rabbitmq.password=guest
spring.rabbitmq.username=guest
spring.rabbitmq.virtual-host=/

Затем предоставьте класс конфигурации RabbitMQ:

@Configuration
public class RabbitConfig {
    public static final String QUEUE_NAME = "javaboy_delay_queue";
    public static final String EXCHANGE_NAME = "javaboy_delay_exchange";
    public static final String EXCHANGE_TYPE = "x-delayed-message";

    @Bean
    Queue queue() {
        return new Queue(QUEUE_NAME, true, false, false);
    }

    @Bean
    CustomExchange customExchange() {
        Map<String, Object> args = new HashMap<>();
        args.put("x-delayed-type", "direct");
        return new CustomExchange(EXCHANGE_NAME, EXCHANGE_TYPE, true, false,args);
    }
    
    @Bean
    Binding binding() {
        return BindingBuilder.bind(queue())
                .to(customExchange()).with(QUEUE_NAME).noargs();
    }
}

Основное отличие здесь в том, что определение переключателя другое, и друзьям нужно обратить внимание.

Переключатель, который мы используем здесь, — это CustomExchange, который является переключателем, предоставленным в Spring.При создании CustomExchange существует пять параметров, и их значения следующие:

  • Переключить имя.
  • Тип переключателя, это место фиксированное.
  • Является ли переключатель постоянным.
  • Удалять ли коммутатор, если к коммутатору не привязана очередь.
  • Другие параметры.

В последнем параметре args указывается тип рассылки сообщений коммутатора.Это всем известный тип direct,fanout,topic и header.Какой тип используется,и каким образом коммутатор будет распространять сообщения в дальнейшем.

Затем мы создаем еще одного потребителя сообщений:

@Component
public class MsgReceiver {
    private static final Logger logger = LoggerFactory.getLogger(MsgReceiver.class);
    @RabbitListener(queues = RabbitConfig.QUEUE_NAME)
    public void handleMsg(String msg) {
        logger.info("handleMsg,{}",msg);
    }
}

Просто распечатайте содержимое сообщения.

Затем напишите метод модульного тестирования для отправки сообщения:

@SpringBootTest
class MqDelayedMsgDemoApplicationTests {

    @Autowired
    RabbitTemplate rabbitTemplate;
    @Test
    void contextLoads() throws UnsupportedEncodingException {
        Message msg = MessageBuilder.withBody(("hello 江南一点雨"+new Date()).getBytes("UTF-8")).setHeader("x-delay", 3000).build();
        rabbitTemplate.convertAndSend(RabbitConfig.EXCHANGE_NAME, RabbitConfig.QUEUE_NAME, msg);
    }

}

Установите время задержки сообщения в заголовке сообщения.

Что ж, затем запустите проект Spring Boot, затем запустите метод модульного тестирования, чтобы отправить сообщение, окончательный журнал печати консоли выглядит следующим образом:

Из лога видно, что задержка сообщений реализована.

6.2 DLX реализует очередь с задержкой

6.2.1 Идея реализации очереди с задержкой

Идея реализации очереди с задержкой тоже очень проста, т.е.Предыдущая статьяТо, что мы называем DLX (обмен недоставленными письмами) + TTL (время ожидания сообщения).

Мы можем думать об очередях недоставленных сообщений как об очередях с задержкой.

Конкретно:

Если сообщение необходимо отложить на 30 минут, мы устанавливаем срок действия этого сообщения на 30 минут и настраиваем переключатель недоставленных сообщений и недоставленных сообщений для этого сообщения.routing_key, и не устанавливайте потребителя для этой очереди сообщений, то через 30 минут сообщение попадает в очередь недоставленных сообщений, потому что оно не потребляется потребителем.В это время у нас есть потребитель, «сидящий на корточках» в очереди недоставленных сообщений. После попадания в очередь недоставленных сообщений она немедленно расходуется.

Это идея реализации очереди с задержкой, она очень простая?

6.2.2 Случай

Далее Сонг Гэ продемонстрирует конкретную реализацию очереди с задержкой на простом примере.

Сначала подготовьте запущенный RabbitMQ.

Затем мы создаем проект Spring Boot и вводим зависимости RabbitMQ:

Затем настройте базовую информацию о подключении RabbitMQ в application.properties:

spring.rabbitmq.host=localhost
spring.rabbitmq.username=guest
spring.rabbitmq.password=guest
spring.rabbitmq.port=5672

Далее настроим две очереди сообщений: обычную очередь и очередь недоставленных сообщений:

@Configuration
public class QueueConfig {
    public static final String JAVABOY_QUEUE_NAME = "javaboy_queue_name";
    public static final String JAVABOY_EXCHANGE_NAME = "javaboy_exchange_name";
    public static final String JAVABOY_ROUTING_KEY = "javaboy_routing_key";
    public static final String DLX_QUEUE_NAME = "dlx_queue_name";
    public static final String DLX_EXCHANGE_NAME = "dlx_exchange_name";
    public static final String DLX_ROUTING_KEY = "dlx_routing_key";

    /**
     * 死信队列
     * @return
     */
    @Bean
    Queue dlxQueue() {
        return new Queue(DLX_QUEUE_NAME, true, false, false);
    }

    /**
     * 死信交换机
     * @return
     */
    @Bean
    DirectExchange dlxExchange() {
        return new DirectExchange(DLX_EXCHANGE_NAME, true, false);
    }

    /**
     * 绑定死信队列和死信交换机
     * @return
     */
    @Bean
    Binding dlxBinding() {
        return BindingBuilder.bind(dlxQueue()).to(dlxExchange())
                .with(DLX_ROUTING_KEY);
    }

    /**
     * 普通消息队列
     * @return
     */
    @Bean
    Queue javaboyQueue() {
        Map<String, Object> args = new HashMap<>();
        //设置消息过期时间
        args.put("x-message-ttl", 1000*10);
        //设置死信交换机
        args.put("x-dead-letter-exchange", DLX_EXCHANGE_NAME);
        //设置死信 routing_key
        args.put("x-dead-letter-routing-key", DLX_ROUTING_KEY);
        return new Queue(JAVABOY_QUEUE_NAME, true, false, false, args);
    }

    /**
     * 普通交换机
     * @return
     */
    @Bean
    DirectExchange javaboyExchange() {
        return new DirectExchange(JAVABOY_EXCHANGE_NAME, true, false);
    }

    /**
     * 绑定普通队列和与之对应的交换机
     * @return
     */
    @Bean
    Binding javaboyBinding() {
        return BindingBuilder.bind(javaboyQueue())
                .to(javaboyExchange())
                .with(JAVABOY_ROUTING_KEY);
    }
}

Хотя этот код конфигурации немного длиннее, принцип на самом деле прост.

  • Конфигурацию можно разделить на две группы, первая группа настраивает очередь недоставленных сообщений, а вторая группа настраивает обычную очередь. Каждая группа состоит из очереди сообщений, обмена сообщениями и привязки.
  • При настройке очереди сообщений укажите очередь недоставленных сообщений для очереди сообщений.Незнакомые друзья могут прочитать предыдущую статью Портал:Срок действия сообщений в RabbitMQ истекает?.
  • При настройке времени истечения срока действия сообщения в очереди единицей времени по умолчанию являются миллисекунды.

Затем мы настраиваем потребителя для очереди недоставленных сообщений следующим образом:

@Component
public class DlxConsumer {
    private static final Logger logger = LoggerFactory.getLogger(DlxConsumer.class);

    @RabbitListener(queues = QueueConfig.DLX_QUEUE_NAME)
    public void handle(String msg) {
        logger.info(msg);
    }
}

Распечатайте сообщение по мере его получения.

Вот и все.

Стартовый проект.

Наконец, мы отправляем сообщение в модульном тесте:

@SpringBootTest
class DelayQueueApplicationTests {

    @Autowired
    RabbitTemplate rabbitTemplate;

    @Test
    void contextLoads() {
        System.out.println(new Date());
        rabbitTemplate.convertAndSend(QueueConfig.JAVABOY_EXCHANGE_NAME, QueueConfig.JAVABOY_ROUTING_KEY, "hello javaboy!");
    }

}

По этому поводу сказать нечего, это обычная отправка сообщения, через 10 секунд это сообщение будет распечатано в потребителе очереди недоставленных сообщений.

Что ж, это две идеи, по которым мы используем RabbitMQ в качестве очереди задержки ~ Заинтересованные друзья могут попробовать ~

7. Надежность отправки RabbitMQ

Микросервисы могут быть спроектированы как микросервисы, управляемые сообщениями, а реагирующие системы также могут быть основаны на промежуточном программном обеспечении сообщений.С этой точки зрения промежуточное программное обеспечение сообщений действительно важно при разработке интернет-приложений.

Взяв в качестве примера RabbitMQ, Сонг Гэ пришел поговорить с вами о надежности отправки сообщения в середине сообщения.

Обратите внимание, что в следующем содержании я в основном обсуждаю с вами, как гарантировать, что производитель сообщения успешно отправляет сообщение, и не затрагивает проблему потребления сообщения.

7.1 Механизм отправки сообщений RabbitMQ

Как мы все знаем, отправка сообщения в RabbitMQ вводит понятие Exchange (обмена), отправка сообщения сначала поступает на биржу, а затем по установленным правилам маршрутизации биржа направляет сообщение в разные Queues (очереди) , а затем потреблять разными потребителями.

Общий процесс таков, поэтому для обеспечения надежности отправки сообщения оно в основном подтверждается с двух сторон:

  1. Сообщение успешно доставлено на Exchange
  2. Сообщение успешно доставлено в очередь

Если эти два шага могут быть подтверждены, то мы можем считать, что сообщение отправлено успешно.

Если есть проблема на любом из этих двух шагов, значит, сообщение не было успешно доставлено. В это время нам, возможно, придется повторить попытку и использовать другие средства для повторной отправки сообщения. После нескольких попыток, если сообщение все еще не может быть доставлено, оно может быть необходимо вмешательство человека.

После приведенного выше анализа мы можем подтвердить, что для успешной отправки сообщения нам нужно сделать всего три вещи:

  1. Убедитесь, что сообщение поступает в Exchange.
  2. Убедитесь, что сообщение поступает в очередь.
  3. Включите запланированные задачи для доставки сообщений, которые не отправляются регулярно.

7.2 Усилия RabbitMQ

Три шага предложенных выше, третий шаг нужно реализовать самим, а первые два шага у RabbitMQ есть готовые решения.

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

  1. Механизм открытых транзакций
  2. Механизм подтверждения отправителя

Это два разных решения. Их нельзя включить одновременно. Можно выбрать только одно из них. Если оба включены одновременно, будет выдано сообщение об ошибке:

Давайте рассмотрим их отдельно. Все следующие кейсы разработаны в Spring Boot, и соответствующий исходный код можно скачать в конце статьи.

7.2.1 Включите механизм транзакций

Способ открытия механизма транзакций RabbitMQ в Spring Boot следующий:

Во-первых, вам необходимо предоставить диспетчер транзакций следующим образом:

@Bean
RabbitTransactionManager transactionManager(ConnectionFactory connectionFactory) {
    return new RabbitTransactionManager(connectionFactory);
}

Затем сделайте две вещи в генераторе сообщений: добавьте аннотации транзакций и установите канал связи в режим транзакций:

@Service
public class MsgService {
    @Autowired
    RabbitTemplate rabbitTemplate;

    @Transactional
    public void send() {
        rabbitTemplate.setChannelTransacted(true);
        rabbitTemplate.convertAndSend(RabbitConfig.JAVABOY_EXCHANGE_NAME,RabbitConfig.JAVABOY_QUEUE_NAME,"hello rabbitmq!".getBytes());
        int i = 1 / 0;
    }
}

Обратите внимание на две вещи:

  1. Добавьте способ отправки сообщения@TransactionalАннотации помечают транзакции.
  2. Вызовите метод setChannelTransacted со значением true, чтобы включить режим транзакций.

Вот и все.

В приведенном выше случае у нас есть 1/0 в конце, что должно вызвать исключение во время выполнения, мы можем попытаться запустить метод и обнаружить, что сообщение не было отправлено успешно.

Когда мы включаем режим транзакций, у производителей RabbitMQ есть еще четыре шага для отправки сообщений:

  1. Клиент делает запрос на перевод канала в режим транзакций.
  2. Сервер отвечает, соглашаясь перевести канал в режим транзакций.
  3. Клиент отправляет сообщение.
  4. Клиент фиксирует транзакцию.
  5. Сервер отвечает, подтверждая транзакцию.

Вышеуказанные шаги, кроме третьего шага, уже есть, а остальные шаги лишние без причины. Итак, вы видите, что режим транзакций на самом деле немного неэффективен, и это не оптимальное решение. Мы можем подумать о том, какие проекты будут использовать промежуточное программное обеспечение сообщений? Вообще говоря, это некоторые проекты с высокой степенью параллелизма, и производительность параллелизма особенно важна в настоящее время.

Поэтому RabbitMQ также предоставляет механизм подтверждения отправителя (подтверждение издателя), чтобы убедиться, что сообщение отправлено успешно.Таким образом, производительность намного выше, чем в режиме транзакций.Давайте посмотрим.

7.2.2 Механизм подтверждения отправителя

7.2.2.1 Обработка одиночного сообщения

Во-первых, мы удаляем код о транзакции прямо сейчас, а затем настраиваем и включаем механизм подтверждения отправителя сообщения в application.properties следующим образом:

spring.rabbitmq.publisher-confirm-type=correlated
spring.rabbitmq.publisher-returns=true

Первая строка — это обратный вызов подтверждения, который настраивает прибытие сообщения на биржу, а вторая строка — это обратный вызов, который настраивает сообщение для поступления в очередь.

Конфигурация атрибута первой строки имеет три значения:

  1. нет: указывает, что режим подтверждения выпуска отключен, что является значением по умолчанию.
  2. коррелированный: указывает метод обратного вызова, который будет запущен после успешной публикации сообщения на бирже.
  3. простой: похож на коррелированный и поддерживаетwaitForConfirms()иwaitForConfirmsOrDie()вызов метода.

Далее нам нужно открыть два монитора, конкретная конфигурация выглядит следующим образом:

@Configuration
public class RabbitConfig implements RabbitTemplate.ConfirmCallback, RabbitTemplate.ReturnsCallback {
    public static final String JAVABOY_EXCHANGE_NAME = "javaboy_exchange_name";
    public static final String JAVABOY_QUEUE_NAME = "javaboy_queue_name";
    private static final Logger logger = LoggerFactory.getLogger(RabbitConfig.class);
    @Autowired
    RabbitTemplate rabbitTemplate;
    @Bean
    Queue queue() {
        return new Queue(JAVABOY_QUEUE_NAME);
    }
    @Bean
    DirectExchange directExchange() {
        return new DirectExchange(JAVABOY_EXCHANGE_NAME);
    }
    @Bean
    Binding binding() {
        return BindingBuilder.bind(queue())
                .to(directExchange())
                .with(JAVABOY_QUEUE_NAME);
    }

    @PostConstruct
    public void initRabbitTemplate() {
        rabbitTemplate.setConfirmCallback(this);
        rabbitTemplate.setReturnsCallback(this);
    }

    @Override
    public void confirm(CorrelationData correlationData, boolean ack, String cause) {
        if (ack) {
            logger.info("{}:消息成功到达交换器",correlationData.getId());
        }else{
            logger.error("{}:消息发送失败", correlationData.getId());
        }
    }

    @Override
    public void returnedMessage(ReturnedMessage returned) {
        logger.error("{}:消息未成功路由到队列",returned.getMessage().getMessageProperties().getMessageId());
    }
}

Относительно этого класса конфигурации скажу следующее:

  1. Определить класс конфигурации, реализоватьRabbitTemplate.ConfirmCallbackиRabbitTemplate.ReturnsCallbackДва интерфейса, эти два интерфейса, обратный вызов первого используется для определения того, что сообщение прибыло на биржу, а второй вызывается, когда сообщение не удается перенаправить в очередь.
  2. Определите метод initRabbitTemplate и добавьте аннотацию @PostConstruct, а также настройте два обратных вызова для rabbitTemplate в этом методе.

Вот и все.

Далее мы тестируем отправку сообщения.

Сначала мы пытаемся отправить сообщение на несуществующий обмен следующим образом:

rabbitTemplate.convertAndSend("RabbitConfig.JAVABOY_EXCHANGE_NAME",RabbitConfig.JAVABOY_QUEUE_NAME,"hello rabbitmq!".getBytes(),new CorrelationData(UUID.randomUUID().toString()));

Обратите внимание, что первый параметр — это строка, а не переменная, переключатель не существует, и консоль выдаст следующую ошибку:

Дальше даем реальный обмен, но несуществующую очередь, вот так:

rabbitTemplate.convertAndSend(RabbitConfig.JAVABOY_EXCHANGE_NAME,"RabbitConfig.JAVABOY_QUEUE_NAME","hello rabbitmq!".getBytes(),new CorrelationData(UUID.randomUUID().toString()));

Обратите внимание, что второй параметр теперь является строкой, а не переменной.

Видно, что хотя сообщение успешно дошло до биржи, оно не было успешно перенаправлено в очередь (поскольку очереди не существует).

Это отправка сообщения, давайте рассмотрим пакетную отправку сообщений.

7.2.2.2 Пакетная обработка сообщений

Если это пакетная обработка сообщений, контроль обратного вызова для успешной отправки такой же, и здесь повторяться не будет.

Это шаблон подтверждения издателя.

По сравнению с транзакциями пропускная способность сообщений в этом режиме будет значительно улучшена.

7.3 Повторить попытку при ошибке

Есть два случая неудачной повторной попытки: первый — MQ вообще не найден, а другой — MQ найден, но сообщение не отправлено.

Давайте рассмотрим два типа повторных попыток отдельно.

7.3.1 Встроенный механизм повторных попыток

Механизм транзакций и механизм подтверждения отправителя, упомянутые выше, позволяют отправителю подтвердить успешную отправку сообщения. Если отправитель не может подключиться к MQ с самого начала, в Spring Boot также есть соответствующий механизм повторных попыток, но этот механизм повторных попыток не имеет ничего общего с самим MQ. Это делается с помощью механизма повторных попыток в Spring. Конкретная конфигурация выглядит следующим образом: следует:

spring.rabbitmq.template.retry.enabled=true
spring.rabbitmq.template.retry.initial-interval=1000ms
spring.rabbitmq.template.retry.max-attempts=10
spring.rabbitmq.template.retry.max-interval=10000ms
spring.rabbitmq.template.retry.multiplier=2

Смысл конфигурации сверху вниз:

  • Включить механизм повтора.
  • Интервал запуска повтора.
  • Максимальное количество повторных попыток.
  • Максимальный интервал между попытками.
  • Множитель интервального времени. (Настроенный здесь множитель интервала равен 2, затем первый интервал равен 1 секунде, второй интервал повторных попыток равен 2 секундам, третий раз равен 4 секундам и т. д.)

После завершения настройки снова запустите проект Spring Boot, а затем отключите MQ. Если вы попытаетесь отправить сообщение в это время, оно не будет отправлено, что приведет к автоматическому повтору.

7.3.2 Повторная попытка обслуживания

Повторная попытка обслуживания в основном предназначена для случая, когда сообщение не достигает коммутатора.

Если сообщение не доходит до обмена успешно, согласно нашему объяснению во втором разделе, будет запущен обратный вызов ошибки отправки сообщения, В этом обратном вызове мы можем поднять шум!

Общая идея такова:

  1. Сначала создайте таблицу для записи сообщений, отправляемых промежуточному программному обеспечению, например:

Каждый раз, когда сообщение отправляется, запись добавляется в базу данных. Поля здесь хорошо понятны, и я добавлю три:

  • статус: Указывает статус сообщения. Имеется три значения: 0, 1 и 2 означают, что сообщение отправляется, сообщение было успешно отправлено и сообщение не удалось отправить.
  • tryTime: указывает время первой повторной попытки сообщения (после отправки сообщения оно не было успешно отправлено во время tryTime, и в это время можно начать повторную попытку).
  • count: указывает количество повторных попыток сообщения.

Другие поля хорошо изучены, поэтому я не буду вдаваться в подробности по отдельности.

  1. Когда сообщение отправлено, мы сохраняем запись об отправке сообщения в таблице и устанавливаем статус состояния на 0, а tryTime на 1 минуту позже.
  2. В методе обратного вызова с подтверждением, если получен обратный вызов, свидетельствующий об успешной отправке сообщения, статус сообщения устанавливается равным 1 (для сообщения при отправке устанавливается msgId, а msgId используется для уникальной блокировки сообщение, когда сообщение отправлено успешно. ).
  3. Кроме того, откройте временную задачу, и временная задача будет отправляться в базу данных каждые 10 с для извлечения сообщений из базы данных, в частности, для извлечения тех записей, статус которых равен 0, а время tryTime истекло.После получения этих сообщений сначала определите превысило ли количество повторных попыток 3 раза, если более 3 раз, изменить статус сообщения на 2, что означает, что сообщение не может быть отправлено и не будет повторной попытки. Для записей с не более чем 3 повторными попытками сообщение отправляется повторно и значение его счетчика +1.

Общая идея та же, что и выше. Сонге не будет давать здесь код. Отправка почты в vhr Сонге обрабатывается таким образом. Полный код можно найти в проекте vhr (github.com/lenve/vhr).

Конечно, у этого подхода есть два недостатка:

  1. Переход к базе данных может замедлить Qos MQ, но иногда нам не нужны высокие Qos MQ, поэтому приложение зависит от конкретной ситуации.
  2. Согласно изложенной выше идее, одно и то же сообщение может быть отправлено несколько раз, но это не проблема, мы можем решить проблему идемпотентности при потреблении сообщений.

Конечно, каждый также должен обращать внимание на то, чтобы сообщение было отправлено на 100% успешно, в зависимости от конкретной ситуации.

Что ж, вот некоторые распространенные проблемы с производителями сообщений и соответствующие решения.В следующем разделе Сонг Гэ обсудит с вами, как обеспечить успешное потребление сообщений и решить проблему идемпотентности.

8. Надежность потребления RabbitMQ

В последнем разделе Сонг Гэ рассказал всем о том, как высокая доступность MQ обеспечивает успешную отправку сообщений.Для обеспечения успешной отправки сообщений используются различные конфигурации, и даже в некоторых крайних случаях одно и то же сообщение может отправляться повторно. в любом случае, сообщение наконец-то отправлено.Если вы не читали предыдущую статью, то рекомендуется сначала прочитать ее, а потом изучить эту статью:

Сегодня поговорим о потреблении сообщений и посмотрим, как обеспечить успешное потребление сообщений и обеспечить идемпотентность.

8.1 Две идеи потребления

В общем, есть две разные идеи для потребления сообщений RabbitMQ:

  • Push: MQ активно отправляет сообщениетолкатьДля потребителей этот метод требует, чтобы потребитель установил буфер для кэширования сообщений.Для потребителей всегда есть куча сообщений, которые необходимо обработать в памяти, поэтому этот метод более эффективен, что также есть в большинстве приложений в настоящее время. , используемый метод потребления.
  • Pull (вытягивание): потребитель активно вытягивает из MQВытащитьСообщения, этот метод не очень эффективен, но иногда, если серверу необходимо получать сообщения пакетами, этот метод можно использовать.

Я приведу пример в обоих случаях.

Давайте сначала посмотрим на push:

Этот способ более распространен, т.@RabbitListenerАннотации для обозначения потребителей следующим образом:

@Component
public class ConsumerDemo {
    @RabbitListener(queues = RabbitConfig.JAVABOY_QUEUE_NAME)
    public void handle(String msg) {
        System.out.println("msg = " + msg);
    }
}

Этот метод срабатывает, когда в очереди прослушивания есть сообщение.

Посмотрим еще раз на тягу:

@Test
public void test01() throws UnsupportedEncodingException {
    Object o = rabbitTemplate.receiveAndConvert(RabbitConfig.JAVABOY_QUEUE_NAME);
    System.out.println("o = " + new String(((byte[]) o),"UTF-8"));
}

Вызовите метод receiveAndConvert. Параметром метода является имя очереди. После выполнения метода сообщение будет извлечено из MQ. Если возвращаемое значение метода равно null, это означает, что в очереди нет сообщения. Метод ReceiveAndConvert имеет перегруженный метод, вы можете указать время ожидания в перегруженном методе, например 3 секунды. В этот момент, предполагая, что в очереди больше нет сообщений, метод receiveAndConvert будет заблокирован на 3 секунды. Если в течение 3 секунд в очереди есть новое сообщение, он вернется. Через 3 секунды, если по-прежнему нет новое сообщение в очереди, оно вернет null.Этот период ожидания Если не установлен, он по умолчанию равен 0.

Это два разных режима потребления сообщений.

Если вам нужно постоянно получать сообщения из очереди сообщений, вы можете использовать режим push; если вы просто потребляете сообщение, вы можете использовать режим pull.Не переводите режим pull в бесконечный цикл и не подписывайтесь на сообщения в замаскированной форме, что серьезно повлияет на производительность RabbitMQ.

8.2 Два способа обеспечить успешное потребление

существуетПредыдущая статьяВ , мы делаем все возможное, чтобы сообщение могло быть отправлено успешно.На самом деле, для успешного использования сообщения официальный предоставляет соответствующий механизм, давайте посмотрим.

Чтобы гарантировать, что сообщения могут надежно достигать потребителей сообщений, RabbitMQ предоставляет механизм подтверждения потребления сообщений. Когда потребители потребляют сообщения, они могут указать параметр autoAck, чтобы указать метод подтверждения потребления сообщений.

  • Когда autoAck имеет значение false, даже если потребитель получил сообщение, RabbitMQ не удалит сообщение немедленно, а пометит сообщение для удаления после того, как потребитель явным образом ответит на сигнал подтверждения, а затем снова Delete.
  • Когда autoAck имеет значение true, потребитель сообщения автоматически установит отправленное сообщение как подтверждение, а затем удалит сообщение (из памяти или с диска), даже если сообщение не дошло до потребителя.

Давайте посмотрим на картинку:

Как показано выше, на странице веб-управления RabbitMQ:

  • Готово указывает количество сообщений, которые необходимо использовать.
  • Unacked указывает количество сообщений, которые были отправлены потребителю, но еще не получили подтверждение от потребителя.

Здесь мы можем наблюдать подтверждение потребления сообщений с уровня пользовательского интерфейса.

Когда мы устанавливаем для autoAck значение false, для RabbitMQ потребление делится на две части:

  • сообщения, которые будут потребляться
  • Сообщения, которые были доставлены потребителям, но еще не были подтверждены потребителями

Другими словами, когда для autoAck установлено значение false, потребитель становится очень спокойным, у него будет достаточно времени, чтобы обработать сообщение, а затем вручную подтвердить, когда сообщение будет обработано нормально, тогда RabbitMQ будет думать, что потребление сообщения прошло успешно. Если RabbitMQ не получил обратной связи от клиента, а клиент в это время был отключен, тогда RabbitMQ только что поместит сообщение обратно в очередь и будет ждать следующего потребления.

Подводя итог, можно сказать, что обеспечение успешного использования сообщений — это не что иное, как подтверждение подтверждения вручную или автоматически, и ничего больше. Конечно, независимо от того, какой из этих двух, это может в конечном итоге привести к повторному потреблению сообщений, поэтому в целом нам также нужно решить проблему идемпотентности при обработке сообщений.

8.3 Отклонение сообщения

Когда клиент получает сообщение, он может принять это сообщение или отклонить его. Давайте посмотрим, как отказаться:

@Component
public class ConsumerDemo {
    @RabbitListener(queues = RabbitConfig.JAVABOY_QUEUE_NAME)
    public void handle(Channel channel, Message message) {
        //获取消息编号
        long deliveryTag = message.getMessageProperties().getDeliveryTag();
        try {
            //拒绝消息
            channel.basicReject(deliveryTag, true);
        } catch (IOException e) {
            e.printStackTrace();
        }
    }
}

После того, как потребитель получит сообщение, он может отказаться от потребления сообщения.Этапы отказа делятся на два шага:

  1. Получите номер сообщения deliveryTag.
  2. Вызовите метод basicReject, чтобы отклонить сообщение.

При вызове метода basicReject вторым параметром является requeue, то есть нужно ли повторно ставить в очередь. Если второй параметр истинен, отклоненное сообщение снова попадет в очередь сообщений и будет ждать следующего потребления, если второй параметр ложен, отклоненное сообщение будет отброшено, и его будут потреблять новые потребители.

Обратите внимание, что метод basicReject может отклонять только одно сообщение за раз.

8.4 Подтверждение сообщения

Подтверждение сообщения делится на автоматическое подтверждение и ручное подтверждение, рассмотрим их отдельно.

8.4.1 Автоматическое подтверждение

Давайте сначала посмотрим на автоматическое подтверждение, В Spring Boot по умолчанию потребление сообщений подтверждается автоматически.

Давайте посмотрим на следующий метод потребления сообщений:

@Component
public class ConsumerDemo {
    @RabbitListener(queues = RabbitConfig.JAVABOY_QUEUE_NAME)
    public void handle2(String msg) {
        System.out.println("msg = " + msg);
        int i = 1 / 0;
    }
}

Текущий класс внедряется в контейнер Spring через аннотацию @Component, а затем метод потребления сообщения помечается аннотацией @RabbitListener.По умолчанию метод потребления сообщения имеет свою собственную транзакцию, то есть если метод выдает исключение во время выполнения это будет Потребляемое сообщение вернется в очередь, чтобы дождаться следующего потребления.Если метод выполняется нормально, не вызывая исключения, сообщение считается потребленным.

8.4.2 Ручное подтверждение

Ручное подтверждение Я разделил его на два типа: ручное подтверждение в режиме push и ручное подтверждение в режиме pull.

8.4.2.1 Ручное подтверждение режима push

Чтобы включить ручное подтверждение, нам нужно сначала отключить автоматическое подтверждение.Способ отключения следующий:

spring.rabbitmq.listener.simple.acknowledge-mode=manual

Эта конфигурация означает изменение режима подтверждения сообщения на ручное подтверждение.

Далее давайте посмотрим на код в потребителе:

@RabbitListener(queues = RabbitConfig.JAVABOY_QUEUE_NAME)
public void handle3(Message message,Channel channel) {
    long deliveryTag = message.getMessageProperties().getDeliveryTag();
    try {
        //消息消费的代码写到这里
        String s = new String(message.getBody());
        System.out.println("s = " + s);
        //消费完成后,手动 ack
        channel.basicAck(deliveryTag, false);
    } catch (Exception e) {
        //手动 nack
        try {
            channel.basicNack(deliveryTag, false, true);
        } catch (IOException ex) {
            ex.printStackTrace();
        }
    }
}

Поместите то, что потребитель хочет сделать, вtry..catchв кодовом блоке.

Если сообщение успешно использовано в обычном режиме, выполнитеbasicAckПолное подтверждение.

Если потребление сообщения не удается, выполнитеbasicNackметод, чтобы сообщить RabbitMQ, что потребление сообщения не удалось.

Здесь задействованы два метода:

  • basicAck: это для ручного подтверждения того, что сообщение было успешно обработано.Этот метод имеет два параметра: первый параметр представляет собой идентификатор сообщения; второй параметр, множественный, если он равен false, означает, что только текущее сообщение успешно обработано. Если это правда, это означает, что все сообщения, которые не были подтверждены текущим потребителем до того, как текущее сообщение было успешно использовано.
  • basicNack: это для того, чтобы сообщить RabbitMQ, что текущее сообщение не было успешно использовано.Этот метод имеет три параметра: первый параметр представляет собой идентификатор сообщения, второй параметр множественный Если он равен false, это означает, что только потребление текущее сообщение отклонено.Если истинно, то Указывает, что все сообщения, которые не были подтверждены текущим потребителем до текущего сообщения, отвергаются; значение третьего параметра requeue такое же, как и предыдущего, независимо от того, были ли отклонены сообщения повторно ставятся в очередь.

Когда последний параметр в basicNack имеет значение false, это также связано с проблемой очереди недоставленных сообщений.Этот Songge напишет статью и подробно расскажет вам в будущем.

8.4.2.2 Ручное подтверждение режима Pull

Подтвердить режим pull вручную сложнее, в RabbitTemplate, инкапсулированном в Spring, нет соответствующего метода, поэтому мы должны использовать нативный метод, как показано ниже:

public void receive2() {
    Channel channel = rabbitTemplate.getConnectionFactory().createConnection().createChannel(false);
    long deliveryTag = 0L;
    try {
        GetResponse getResponse = channel.basicGet(RabbitConfig.JAVABOY_QUEUE_NAME, false);
        deliveryTag = getResponse.getEnvelope().getDeliveryTag();
        System.out.println("o = " + new String((getResponse.getBody()), "UTF-8"));
        channel.basicAck(deliveryTag, false);
    } catch (IOException e) {
        try {
            channel.basicNack(deliveryTag, false, true);
        } catch (IOException ex) {
            ex.printStackTrace();
        }
    }
}

Используемые здесь методы basicAck и basicNack такие же, как и в предыдущих, поэтому я не буду их повторять.

8.5 Проблема идемпотентности

Наконец, давайте поговорим об идемпотентности сообщений.

Представьте себе следующий сценарий:

После того, как потребитель завершает использование сообщения, он отправляет подтверждение подтверждения в RabbitMQ. В это время RabbitMQ не получает подтверждение из-за отключения сети или по другим причинам. Тогда RabbitMQ не будет удалять сообщение в это время. установлено После того, как соединение установлено, потребитель все еще будет получать сообщение снова, что приводит к повторному использованию сообщения. В то же время по аналогичным причинам при отправке сообщения одно и то же сообщение может быть отправлено дважды (см.Четыре стратегии для обеспечения надежности отправки сообщений RabbitMQ! Что вы используете?). По разным причинам мы должны иметь дело с проблемой идемпотентности при использовании сообщений.

Разобраться с проблемой идемпотентности несложно.В принципе, она решается с точки зрения бизнеса.Позвольте мне рассказать об идее.

Используя Redis, прежде чем потребитель получит сообщение, идентификатор сообщения теперь помещается в Redis, а метод хранения выглядит следующим образом:

  • id-0 (ведение бизнеса)
  • id-1 (успех в бизнесе)

Если проверка не удалась, когда RabbitMQ доставляет сообщение другим потребителям, сначала выполните setnx. Если ключ уже существует (указывая, что кто-то уже использовал сообщение ранее), получите его значение. Если он равен 0, текущий потребитель ничего не имеет. , если он равен 1, подтвердить напрямую.

Крайний случай: когда первый потребитель выполняет бизнес, возникает тупик.На основе setnx для ключа устанавливается время жизни. Производитель, при отправке сообщения указывайте messageId.

Конечно, это всего лишь простая идея для справки.

Сонг Гэ также занимался проблемой идемпотентности сообщений в проекте vhr.Заинтересованные партнеры могут просмотреть исходный код vhr (github.com/lenve/vhr), код находится на почтовом сервере.

9. Понимание виртуального хоста

Когда мы устанавливаем RabbitMQ в первый раз, мы можем управлять этим RabbitMQ через веб-страницу.По умолчанию пользователем по умолчанию, который мы используем в первый раз, является гость.

После успешного входа вы можете просмотреть всех пользователей на вкладке администратора:

Как видите, у каждого пользователя естьCan access virtual hostsатрибут, что означает этот атрибут?

Сегодня брат Сун пришел немного поболтать с вами.

9.1 Мультиарендность

В RabbitMQ есть понятие multi-tenancy, как его понимать?

Мы устанавливаем сервер RabbitMQ, и каждый сервер RabbitMQ может создавать множество виртуальных серверов сообщений.Эти виртуальные серверы сообщений — это то, что мы называем виртуальными хостами, обычно называемыми виртуальными хостами.

По сути, каждый виртуальный хост представляет собой независимый небольшой сервер RabbitMQ. Этот виртуальный хост будет иметь свою собственную очередь сообщений, обмен сообщениями, соответствующие отношения привязки и т. д., а также имеет свои собственные независимые полномочия. Очереди на разных виртуальных хостах и ​​коммутаторах не могут быть привязаны к друг друга, что обеспечивает безопасную работу и позволяет избежать конфликтов имен.

Нам не нужно смотреть конкретно на виртуальные хосты, они такие же, как обычные физические RabbitMQ.Различные виртуальные хосты могут обеспечивать логическое разделение, чтобы гарантировать, что разные очереди сообщений приложений могут работать безопасно и независимо.

Что я думаю об отношениях между vhost и RabbitMQ? RabbitMQ эквивалентен файлу Excel, а vhost — это лист в файле Excel, и все наши операции выполняются на определенном листе.

По сути, vhost — это концепция протокола AMQP.

9.2 Создание виртуальных хостов из командной строки

Давайте сначала посмотрим, как создать виртуальный хост из командной строки.

Поскольку RabbitMQ здесь установлен с докером, мы сначала входим в контейнер докера:

docker exec -it some-rabbit /bin/bash

Затем выполните следующую команду, чтобы создать/myvhвиртуальный хост:

rabbitmqctl add_vhost myvh

Окончательный результат выполнения следующий:

Затем вы можете просмотреть существующий виртуальный хост с помощью следующей команды:

rabbitmqctl list_vhosts

Разумеется, к этой команде также можно добавить две опции: имя и трассировка.Имя указывает имя виртуального хоста, а трассировка указывает, используется ли функция трассировки (трассировка может помочь отслеживать входящие и исходящие сообщения в RabbitMQ), как показано на следующем рисунке:

Vhost можно удалить с помощью следующей команды:

rabbitmqctl delete_vhost myvh

Когда виртуальный хост удаляется, все очереди сообщений, переключатели и отношения привязки, связанные с виртуальным хостом, будут удалены.

Настройте виртуальный хост для пользователя:

rabbitmqctl set_permissions -p myvh guest ".*" ".*" ".*"

Предыдущие параметры легко сказать, последние три".*"Значения следующие:

  • Пользователи имеют настраиваемые права доступа ко всем ресурсам (создание/удаление очередей сообщений, создание/удаление обменов и т. д.).
  • Пользователь имеет права на запись (сообщения) на все ресурсы.
  • Пользователь имеет разрешение на чтение всех ресурсов (потребление сообщений, пустая очередь и т. д.).

Запретить пользователю доступ к виртуальному хосту:

rabbitmqctl clear_permissions -p myvh guest

9.3 Создать виртуальный хост для страницы управления

Конечно, мы также можем управлять виртуальными хостами в Интернете:

На вкладке администратора нажмите «Виртуальные хосты» справа следующим образом:

Затем нажмите Добавить новый виртуальный хост ниже, чтобы добавить новый виртуальный хост:

После ввода виртуального хоста вы можете изменить его разрешения и удалить виртуальный хост, как показано ниже:

9.4 Управление пользователями

Поскольку vhost обычно появляется вместе с пользователем, здесь я также упомяну, кстати, связанные операции пользователя.

Добавьте пользователя с именем javaboy и паролем 123 следующим образом:

rabbitmqctl add_user javaboy 123

Пароль пользователя можно изменить с помощью следующей команды (измените пароль javaboy на 123456):

rabbitmqctl change_password javaboy 123456

Пароль пользователя можно проверить с помощью следующей команды:

rabbitmqctl authenticate_user javaboy 123456

Успех проверки и отказ проверки следующие:

Вы можете просмотреть всех текущих пользователей с помощью следующей команды:

Первый столбец — это имя пользователя, а второй — роль пользователя.

Что касается ролей пользователей, то я уже говорил о них в предыдущей статье, поэтому повторяться здесь не буду. Портал:Как использовать страницу управления RabbitMQ.

Команда для установки роли для пользователя выглядит следующим образом (установить роль администратора для javaboy):

rabbitmqctl set_user_tags javaboy administrator

Наконец, команда для удаления пользователя выглядит следующим образом:

rabbitmqctl delete_user javaboy

10. REST API

Что касается управления RabbitMQ, мы можем сделать это через веб-страницу.В предыдущей статье Song Ge я также сделал связанное введение для своих друзей:

Однако если мы установим плагин rabbitmq_management, то есть установим клиент веб-управления в RabbitMQ, то сможем управлять RabbitMQ через REST API.

Некоторые друзья могут спросить, какая польза от этого?

Если в нашем проекте используется графический инструмент, такой как Granglia или Graphite, и мы хотим зафиксировать текущее потребление/накопление сообщений в RabbitMQ, мы можем использовать REST API для запроса этой информации и передачи результатов запроса в новый графический инструмент. В то же время, поскольку REST API является HTTP-запросом, поддерживаемые клиенты также разнообразны. Пока он может отправлять HTTP-запрос, его можно использовать. Разве это не особенно удобно?

10.1 REST API

Могут быть некоторые друзья, которые не знают, что такое REST API, поэтому вот краткое введение:

REST (передача репрезентативного состояния) — это архитектурный стиль веб-программного обеспечения, это стиль, а не стандарт, сетевые службы, которые соответствуют или совместимы с этим архитектурным стилем, называются службами REST.

Службы REST лаконичны и иерархичны и обычно основаны на существующих широко популярных протоколах и стандартах, таких как HTTP, URI, XML и HTML. В REST ресурсы задаются URI, а операции добавления, удаления, изменения и запроса ресурсов могут быть реализованы через GET, POST, PUT, DELETE и другие методы, предоставляемые протоколом HTTP.

Использование REST может более эффективно использовать кеш для повышения скорости ответа, а состояние сеанса связи в REST поддерживается клиентом, что позволяет разным серверам обрабатывать разные запросы в серии запросов, тем самым улучшая масштабируемость сервера.

В проектах с разделением клиентской и серверной части хорошо спроектированная архитектура веб-приложения должна соответствовать стилю REST.

10.2 Откройте страницу веб-управления

Давайте поговорим о том, как открыть страницу веб-управления Вообще говоря, у нас есть два способа открыть страницу веб-управления:

  1. При установке RabbitMQ выберите напрямуюrabbitmq:3-managementЗеркало, команда установки следующая:
docker run -d --rm --hostname my-rabbit --name some-rabbit -p 15672:15672 -p 5672:5672 rabbitmq:3-management

Таким образом, установленный RabbitMQ может напрямую использовать веб-страницу управления.

  1. При установке выбрать нормальный нормальный образrabbitmq:3, команда установки выглядит следующим образом:
docker run -d --hostname my-rabbit --name some-rabbit2 -p 5673:5672 -p 25672:15672 rabbitmq:3

После того, как это установлено, нам нужно войти в контейнер, а затем вручную открыть плагин веб-управления, команда выглядит следующим образом:

docker exec -it some-rabbit2 /bin/bash
rabbitmq-plugins enable rabbitmq_management

Первая команда предназначена для входа в контейнер, а вторая команда включает подключаемый модуль веб-управления.Результаты выполнения следующие:

Откройте страницу веб-управления одним из двух указанных выше способов, после чего мы сможем использовать REST API.

10.3 Практика

Далее мы познакомимся с несколькими распространенными операциями REST API.

Мы можем отправлять запросы через инструмент CURL или мы можем отправлять запросы через POSTMAN, просто выберите тот, который вам нравится. Сонг Гэ покажет вам оба пути.

10.3.1 Просмотр статистики очереди

Например, если мы хотим просмотреть статистику очереди hello-queue под виртуальным хостом myvh, мы можем просмотреть ее следующим образом:

curl -i -u javaboy:123 http://localhost:15672/api/queues/myvh/hello-queue

-iУказывает, что отображается информация заголовка ответа.

Окончательный результат выполнения следующий:

Как видите, возвращаемая информация имеет заголовки ответа и JSON, но возвращаемый JSON не отформатирован, что выглядит немного неудобно.Если возвращаемые данные только JSON без заголовков ответа, то мы можем использовать python для завершения форматирования данных, следующим образом:

Как видите, возвращаемые данные отформатированы.

Конечно, мы также можем использовать POSTMAN для отправки этого запроса следующим образом:

Обратите внимание, что метод аутентификации — Basic Auth, и установлены правильные имя пользователя и пароль.

Запрос POSTMAN все же гораздо удобнее.

10.3.2 Создание очереди

Создайте очередь с именем javaboy-queue на виртуальном хосте /myvh и используйте метод запроса CURL следующим образом:

curl -i -u javaboy:123 -XPUT -H "Content-Type:application/json" -d '{"auto_delete":false,"durable":true}' http://localhost:15672/api/queues/myvh/javaboy-queue

Обратите внимание, что метод запроса — это запрос PUT, а параметры запроса — в формате JSON.В JSON есть две вещи: однаauto_deleteОзначает, будет ли очередь автоматически удаляться, если у очереди нет потребительских подписок (если это какая-то временная очередь, это свойство можно установить в true); другоеdurableЭто означает, является ли очередь постоянной (постоянная очередь, очередь все еще существует после перезапуска RabbitMQ), если вы создали очередь с кодом Java, эти два параметра легко понять, потому что, когда мы создаем очередь с кодом Java, два параметра также часто используется.

Конечно, мы также можем использовать POSTMAN для отправки запросов:

возвращение201 CreatedУказывает, что очередь была успешно создана.

Будьте осторожны, чтобы установить имя пользователя/пароль на вкладке Авторизация:

10.3.3 Просмотр информации о текущем соединении

Мы можем просмотреть текущую информацию о соединении с помощью следующего запроса:

Запрос выглядит следующим образом:

curl -i -u javaboy:123 http://localhost:15672/api/connections

Способы просмотра POSTMAN следующие:

10.3.4 Просмотр текущей информации о пользователе

curl -i -u javaboy:123 http://localhost:15672/api/users

POSTMAN просматривает информацию следующим образом:

10.3.5 Создание пользователя

Создайте пользователя с именем zhangsan с паролем 123 и ролью администратора.

CURL:

curl -i -u javaboy:123 -H "{Content-Type:application/json}" -d '{"password":"123","tags":"administrator"}' -XPUT http://localhost:15672/api/users/zhangsan

POSTMAN:

10.3.6 Настройка виртуальных хостов для новых пользователей

Установите пользователя с именем zhangsan на виртуальный хост с именем myvh:

curl -i -u javaboy:123 -H "{Content-Type:application/json}" -d '{"configure":".*","write":".*","read":".*"}' -XPUT http://localhost:15672/api/permissions/myvh/zhangsan

Параметр представляет собой конкретную информацию о разрешении:

Метод запроса POSTMAN выглядит следующим образом:

Что ж, Сонг Гэ даст вам здесь несколько примеров.Для использования других API друзья могут открыть страницу управления RabbitMQ, нажать кнопку HTTP API ниже, и в ней есть полный документ:

11. Общие рабочие команды

Ранее мы представили некоторые REST API, очень удобно вызывать эти REST API там, где удобно отправлять HTTP-запросы. Однако в некоторых местах, где неудобно отправлять HTTP-запросы, эти REST API не очень удобны в использовании, поэтому сегодня Song Ge представит еще одну игру RabbitMQ ---rabbitmqadmin.

11.1 rabbitmqadmin

Обычно мы делаем упражнения самостоятельно и обычно открываем страницу веб-управления RabbitMQ, однако в производственной среде часто нет страницы веб-управления, и MQ можно управлять только с помощью команд CLI.

На самом деле, хотя страница веб-управления удобна, она часто не так быстра, как CLI, и с помощью командной строки CLI мы можем выполнять больше настроек, например, узнавать ключевую информацию и предоставлять ее для централизованного мониторинга. система для срабатывания сигнализации.

Непосредственно управлять командной строкой CLI немного проблематично, RabbitMQ предоставляет инструмент управления CLI rabbitmqadmin , который на самом деле представляет собой скрипт, написанный на Python на основе HTTP API RabbitMQ. Поскольку вручную писать запросы для REST API довольно хлопотно, эти скрипты просто упрощают нам эту операцию и делают ее проще.

Чтобы использовать rabbitmqadmin, вы должны сначала установить его.

Если мы создадим контейнер RabbitMQ сrabbitmq:3-managementзеркало, то по умолчанию установлен rabbitmqadmin.

В противном случае нам может понадобиться установить rabbitmqadmin самостоятельно, метод установки очень прост,

Сначала убедитесь, что на вашем устройстве установлен Python., который является самым простым, потому что инструмент rabbitmqadmin представляет собой скрипт Python.

Затем откройте веб-страницу управления RabbitMQ и введите следующий адрес (моя страница управления привязана к 25672):

http://localhost:25672/cli/index.html

На открывшейся странице вы можете увидеть ссылку для скачивания rabbitmqadmin. После загрузки rabbitmqadmin дайте ему права на выполнение:

chmod +x rabbitmqadmin

После загрузки rabbitmqadmin мы можем открыть его напрямую с помощью Блокнота, который на самом деле представляет собой набор скриптов Python.

Этот процесс довольно сложен в эксплуатации, поэтому я рекомендую вам использовать его напрямую.rabbitmq:3-managementЗеркало, в один шаг.

11.2 Особенности rabbitmqadmin

  • Список обменов, очередей, привязок, виртуальных хостов, пользователей, разрешений, соединений и каналов.
  • Создавайте и удаляйте обмены, очереди, привязки, виртуальные хосты, пользователей и разрешения.
  • Публикуйте и получайте сообщения, а также сведения о сообщениях.
  • Закройте соединение и очистите очередь.
  • Импорт и экспорт конфигурации.

Затем Сун Гэ представит эти функции одну за другой своим друзьям.

11.3 Список различной информации

Посмотреть все переключатели:

rabbitmqadmin list exchanges

Посмотреть все очереди:

rabbitmqadmin list queues

Посмотреть все привязки:

rabbitmqadmin list bindings

Посмотреть все виртуальные хосты:

rabbitmqadmin list vhosts

Посмотреть всю информацию о пользователе:

rabbitmqadmin list users

Посмотреть всю информацию о разрешениях:

rabbitmqadmin list permissions

Посмотреть всю информацию о подключении:

rabbitmqadmin list connections

Посмотреть всю информацию о канале:

rabbitmqadmin list channels

11.4 Полный пример

Далее давайте воспользуемся rabbitmqadmin, чтобы написать полный пример обмена сообщениями.

Сначала создайте обмен с именем javaboy-exchange:

rabbitmqadmin declare exchange name=javaboy-exchange durable=true auto_delete=false type=direct

Различные параметры здесь легко понять, поэтому я не буду говорить больше.

Затем создайте очередь с именем javaboy-queue:

rabbitmqadmin declare queue name=javaboy-queue durable=true auto_delete=false

Затем создайте Binding для привязки коммутатора к очереди сообщений:

rabbitmqadmin declare binding source=javaboy-exchange destination=javaboy-queue routing_key=javaboy-routing

Здесь задействованы три понятия:

  • источник: источник фактически относится к коммутатору.
  • назначение: фактически назначение относится к очереди сообщений.
  • routing_key: это ключ маршрутизации.

Далее опубликуйте сообщение:

rabbitmqadmin publish routing_key=javaboy-queue payload="hello javaboy"

Параметры тут очень простые, тут и говорить нечего.

Просмотр сообщений в очереди (только просмотр, а не потребление, сообщение остается после прочтения):

rabbitmqadmin get queue=javaboy-queue

Пустые сообщения из очереди:

rabbitmqadmin purge queue name=javaboy-queue

11.5 Список команд

Шрифт в таблице мелковат, все отвечают на rabbitmqadmin в фоновом режиме паблик-аккаунта [Jiangnan A Little Rain], чтобы получить ссылку на документ Excel.

12. Система разрешений RabbitMQ

Независимо от того, создаем ли мы пользовательские объекты с помощью веб-страниц или инструментов командной строки, вновь созданные пользовательские объекты нельзя использовать напрямую. Нам нужно сначала поместить этого пользователя под виртуальный хост, а затем предоставить ему разрешения. С разрешением этот пользователь может использовать его в обычном режиме. .

Итак, сегодня давайте взглянем на систему разрешений в RabbitMQ и посмотрим, как эта система разрешений выглядит.

12.1 Введение в систему разрешений RabbitMQ

RabbitMQ реализовал систему разрешений в стиле ACL, начиная с версии 1.6. Могут быть некоторые друзья, которые не знают, что такое система разрешений в стиле ACL. Вы можете посмотреть две предыдущие статьи Songge:

В этой системе управления разрешениями в стиле ACL разрешено множество мелких элементов управления разрешениями, и такие разрешения, как чтение, запись и настройка, могут быть установлены для разных пользователей.

Здесь задействованы три разных разрешения:

  • Чтение: все операции, связанные с потреблением сообщений, включая очистку всей очереди сообщений.
  • Написать: опубликовать сообщение.
  • Конфигурация: Создание и удаление очередей сообщений, обменов и т. д.

Это краткое введение в систему разрешений RabbitMQ.

12.2 Соответствие между операциями и разрешениями

Далее на следующем рисунке показано соответствие между операциями и разрешениями:

Фоновый ответ публичной учетной записиrabbitmq_permissionДоступна электронная таблица Excel с этим графиком.

Какую команду выполнять и какие требуются разрешения, четко описано на этой картинке.

12.3 Команды управления разрешениями

Формат команды операции разрешения в RabbitMQ следующий:

rabbitmqctl set_permissions [-p vhosts] {user} {conf} {write} {read}

Здесь есть несколько параметров:

  • [-p vhost]: имя виртуального хоста, который предоставляет пользователю доступ, если не написано, по умолчанию/.
  • пользователь: Имя пользователя.
  • conf: На какие ресурсы пользователь имеет настраиваемые разрешения (поддерживаются регулярные выражения).
  • write: На каких ресурсах у пользователя есть разрешение на запись (поддерживаются регулярные выражения).
  • чтение: на каких ресурсах у пользователя есть разрешение на чтение (поддерживаются регулярные выражения).

Что касается того, что могут делать настраиваемые разрешения, что могут делать разрешения на запись и что могут делать разрешения на чтение, вы можете обратиться ко второму разделу, который здесь не будет повторяться.

Песня Ge, чтобы привести простой пример.

Предположим, у нас естьzhangsanuser, мы хотим, чтобы этот пользователь имел все права доступа к виртуальному хосту myvh, тогда наша рабочая команда выглядит следующим образом:

rabbitmqctl set_permissions -p myvh zhangsan ".*" ".*" ".*"

Результат выполнения следующий:

Затем выполните следующую команду, чтобы убедиться, что авторизация прошла успешно:

rabbitmqctl -p myvh list_permissions

Видно, что разрешения Чжан Саня были назначены на место.

В приведенной выше команде авторизации мы используем оба".*", Song Ge добавит этот подстановочный знак:

  • ".*": Это означает совпадение всех обменов и очередей.
  • "javaboy-.*": это означает совпадение имен, начинающихся сjavaboy-Коммутаторы и очереди в начале.
  • "": Это означает, что он не соответствует никаким очередям и обменам (вы можете использовать это, если хотите отозвать разрешение пользователя).

Мы можем использовать следующую команду, чтобы удалить разрешения пользователя на виртуальном хосте, например, удалить все разрешения zhangsan на myvh следующим образом:

rabbitmqctl clear_permissions -p myvh zhangsan

После выполнения мы можем пройтиrabbitmqctl -p myvh list_permissionsкоманда, чтобы проверить, вступил ли в силу результат выполнения.Окончательный эффект выполнения выглядит следующим образом:

Если у пользователя есть соответствующие разрешения на несколько виртуальных хостов, следуйте приведенным выше инструкциям.rabbitmqctl -p myvh list_permissionsКоманда может просматривать разрешения только на одном виртуальном хосте. В настоящее время мы можем использовать следующую команду для просмотраlisiРазрешения на всех vhosts:

rabbitmqctl list_user_permissions lisi

12.4 Работа со страницей веб-управления

Конечно, если вы не хотите вводить команды, вы также можете управлять разрешениями через веб-терминал управления.

На вкладке «Администратор» щелкните имя пользователя, чтобы установить разрешения для пользователя, как показано ниже:

Разрешения могут быть установлены, а разрешения также могут быть очищены.

Конечно, на веб-странице также есть тематические разрешения, представляющие собой новую функцию, начиная с RabbitMQ3.7, которая можетtopic exchangeРазрешения устанавливаются в основном для протокола STOMP или MQTT, мы редко используем эту конфигурацию в нашей повседневной разработке Java. Если пользователь не установил его, соответствующийtopic exchangeтакже всегда иметь разрешения.

13. Создание кластера RabbitMQ

Один RabbitMQ не может обеспечить высокую доступность.Для достижения высокой доступности он должен быть кластеризован.

Сегодня Brother Song расскажет вам о построении кластера RabbitMQ.

13.1 Два режима

Когда дело доходит до кластеров, первый вопрос может заключаться в следующем: если у меня есть кластер RabbitMQ, сохраняет ли каждый экземпляр в моем кластере сообщений копию?

На самом деле это включает два режима кластера RabbitMQ:

  • нормальный кластер
  • зеркальный кластер

13.1.1 Общий кластер

Обычный режим кластера заключается в развертывании RabbitMQ на нескольких серверах, каждый сервер запускает экземпляр RabbitMQ, и несколько экземпляров взаимодействуют друг с другом.

В это время созданная очередь Queue, ее метаданные (в основном некоторая информация о конфигурации Queue) будут синхронизированы во всех экземплярах RabbitMQ, но сообщения в очереди будут существовать только в одном экземпляре RabbitMQ и не будут синхронизированы с другими очередями. .

Когда мы потребляем сообщение, если оно подключено к другому экземпляру, то этот экземпляр определяет местоположение очереди с помощью метаданных, а затем обращается к экземпляру, в котором находится очередь, для извлечения данных и отправки их потребителю.

Этот тип кластера может улучшить пропускную способность сообщений RabbitMQ, но он не может гарантировать высокую доступность, поскольку после зависания экземпляра RabbitMQ доступ к сообщению становится невозможен.Если очередь сообщений является постоянной, она может продолжаться после восстановления экземпляра RabbitMQ. ; если очередь сообщений не сохраняется, сообщение теряется.

Общая блок-схема выглядит следующим образом:

13.1.2 Зеркальный кластер

Самая большая разница между ним и обычными кластерами заключается в том, что данные очереди и исходные данные больше не хранятся на одной машине отдельно, а на нескольких машинах одновременно. Другими словами, каждый экземпляр RabbitMQ имеет зеркальные данные (данные копии). Каждый раз, когда пишется сообщение, данные автоматически синхронизируются с несколькими экземплярами, так что после сбоя одной из машин у других машин остается копия данных для продолжения предоставления услуг, что обеспечивает высокую доступность.

Общая блок-схема выглядит следующим образом:

13.1.3 Типы узлов

В RabbitMQ есть два типа узлов:

  • Узел RAM: Узел памяти хранит все определения метаданных очередей, коммутаторов, привязок, пользователей, разрешений и виртуальных хостов в памяти, что дает преимущество в более быстром выполнении таких операций, как объявления коммутаторов и очередей.
  • Дисковый узел: Храните метаданные на диске.Система с одним узлом позволяет только узлам дискового типа предотвратить потерю информации о конфигурации системы при перезапуске RabbitMQ.

RabbitMQ требует по крайней мере один дисковый узел в кластере, все остальные узлы могут быть узлами памяти, и когда узел присоединяется к кластеру или покидает его, он должен уведомить по крайней мере один дисковый узел об изменении. Если единственный дисковый узел в кластере выходит из строя, кластер продолжает работать, но никакие другие операции (добавление, удаление, изменение, проверка) не могут быть выполнены до тех пор, пока узел не восстановится. Чтобы обеспечить надежность информации о кластере или если вы не уверены, использовать ли дисковые узлы или узлы памяти, рекомендуется использовать дисковые узлы напрямую.

13.2 Создание общего кластера

13.2.1 Предварительные знания

Общая структура понятна, а потом будем строить кластер. Давайте начнем с обычного кластера, и мы будем использовать докер для его создания.

Прежде чем строить, есть два предварительных знания, которые вам необходимо знать:

  1. При построении кластера значение cookie Erlang в узлах должно быть согласованным.По умолчанию файл находится в /var/lib/rabbitmq/.erlang.cookie.Когда мы используем docker для создания контейнера RabbitMQ, мы можем установить соответствующий значение cookie для него.
  2. RabbitMQ подключает сервисы через имена хостов, и необходимо гарантировать, что каждое имя хоста может быть пропинговано. Вы можете вручную добавить соответствия имени хоста и IP, отредактировав файл /etc/hosts. Если имя хоста не может быть пропинговано, служба RabbitMQ не запустится (если мы создаем кластер RabbitMQ на разных серверах, нам нужно обратить внимание на этот момент. В следующем резюме 2.2 мы будем использовать ссылку подключения контейнера Docker для реализуют связь между контейнерами (доступ, немного другой).

13.2.2 Начало сборки

Выполните следующие команды, чтобы создать три контейнера RabbitMQ:

docker run -d --hostname rabbit01 --name mq01 -p 5671:5672 -p 15671:15672 -e RABBITMQ_ERLANG_COOKIE="javaboy_rabbitmq_cookie" rabbitmq:3-management
docker run -d --hostname rabbit02 --name mq02 --link mq01:mylink01 -p 5672:5672 -p 15672:15672 -e RABBITMQ_ERLANG_COOKIE="javaboy_rabbitmq_cookie" rabbitmq:3-management
docker run -d --hostname rabbit03 --name mq03 --link mq01:mylink02 --link mq02:mylink03 -p 5673:5672 -p 15673:15672 -e RABBITMQ_ERLANG_COOKIE="javaboy_rabbitmq_cookie" rabbitmq:3-management

Результаты приведены ниже:

Три узла теперь запущены, обратите внимание, что в mq02 и mq03 соответственно используется--linkПараметры для подключения к контейнеру. Если вы не понимаете этот параметр, вы можете ответить на докер в фоновом режиме общедоступной учетной записи Jiangnan Yiyuyu. Об этом рассказывается во вводном руководстве по докеру, написанном Сонг Гэ. Я не буду здесь многословен. Также обратите внимание, что контейнер mq03 должен иметь возможность подключаться как к mq01, так и к mq02.

Далее входим в контейнер mq02, сначала проверяем файл hosts, видно, что настроенное нами подключение к контейнеру вступило в силу:

В дальнейшем в контейнере mq02 можно будет получить доступ к контейнеру mq01 через mylink01 или rabbit01.

Далее приступаем к настройке кластера.

Выполните следующие команды, чтобы добавить контейнер mq02 в кластер:

rabbitmqctl stop_app
rabbitmqctl join_cluster rabbit@rabbit01
rabbitmqctl start_app

Далее мы можем проверить состояние кластера, введя следующую команду:

rabbitmqctl cluster_status

Как видите, в кластере уже две ноды.

Далее аналогичным образом добавляем mq03 в кластер:

rabbitmqctl stop_app
rabbitmqctl join_cluster rabbit@rabbit01
rabbitmqctl start_app

Далее мы можем просмотреть информацию о кластере:

Как видите, на данный момент в кластере уже есть три узла.

Фактически, в настоящее время мы также можем просматривать информацию о кластере через веб-страницу.На домашней веб-странице трех экземпляров RabbitMQ мы можем видеть следующее содержимое:

13.2.3 Тестирование кода

Далее давайте кратко протестируем этот кластер.

Мы создаем родительский проект с именем mq_cluster_demo, а затем создаем в нем два подпроекта.

Первый подпроект называется provider, который является производителем сообщений.При создании он вводит зависимости Web и RabbitMQ следующим образом:

Затем настройте applicaiton.properties следующим образом (обратите внимание на конфигурацию кластера):

spring.rabbitmq.addresses=localhost:5671,localhost:5672,localhost:5673
spring.rabbitmq.username=guest
spring.rabbitmq.password=guest

Далее предоставляется простая очередь, а именно:

@Configuration
public class RabbitConfig {
    public static final String MY_QUEUE_NAME = "my_queue_name";
    public static final String MY_EXCHANGE_NAME = "my_exchange_name";
    public static final String MY_ROUTING_KEY = "my_queue_name";

    @Bean
    Queue queue() {
        return new Queue(MY_QUEUE_NAME, true, false, false);
    }

    @Bean
    DirectExchange directExchange() {
        return new DirectExchange(MY_EXCHANGE_NAME, true, false);
    }

    @Bean
    Binding binding() {
        return BindingBuilder.bind(queue())
                .to(directExchange())
                .with(MY_ROUTING_KEY);
    }
}

Об этом говорить нечего, это базовый контент, далее проверим отправку сообщения в юнит-тесте:

@SpringBootTest
class ProviderApplicationTests {

    @Autowired
    RabbitTemplate rabbitTemplate;

    @Test
    void contextLoads() {
        rabbitTemplate.convertAndSend(null, RabbitConfig.MY_QUEUE_NAME, "hello 江南一点雨");
    }

}

После успешной отправки сообщения на стороне веб-управления RabbitMQ мы увидим сообщение, отображаемое во всех трех экземплярах RabbitMQ, но на самом деле само сообщение существует только в одном экземпляре RabbitMQ.

Далее мы создадим потребитель сообщений.Зависимости и конфигурация потребителя сообщений точно такие же, как у производителя сообщений.Я не буду повторяться.Добавим в потребитель сообщений приемник сообщений:

@Component
public class MsgReceiver {

    @RabbitListener(queues = RabbitConfig.MY_QUEUE_NAME)
    public void handleMsg(String msg) {
        System.out.println("msg = " + msg);
    }
}

Когда потребитель сообщений успешно запущен, в этом методе принимается только одно сообщение, что еще раз подтверждает, что созданный нами кластер RabbitMQ в порядке.

13.2.4 Бэктестинг

Затем Сун Гэ приведет два контрпримера, чтобы доказать, что сообщение не синхронизировано с другими экземплярами RabbitMQ.

Убедитесь, что три экземпляра RabbitMQ находятся в состоянии запуска, закройте Потребителя, а затем отправьте сообщение через провайдера.После успешной передачи закройте экземпляр mq01, а затем запустите экземпляр Потребителя.В это время Потребитель Экземпляр не будет потреблять сообщения, но сообщит об ошибке, в которой говорится, что экземпляр mq01. Если соединение не может быть установлено, этот пример показывает, что сообщение находится на mq01 и не синхронизировано с двумя другими MQ. Напротив, если провайдер успешно отправляет сообщение, мы не закрываем экземпляр mq01, а закрываем экземпляр mq02, тогда вы обнаружите, что потребление сообщения не затрагивается.

13.3 Создание зеркального кластера

Так называемый режим зеркального кластера не требует дополнительного построения, нам нужно только настроить очередь как зеркальную очередь.

Данную конфигурацию можно настроить через веб-страницу или через командную строку, рассмотрим их отдельно.

13.3.1 Зеркальная очередь конфигурации веб-страницы

Давайте сначала посмотрим, как настроить зеркальную очередь на веб-странице.

Щелкните вкладку «Администратор», затем щелкните «Политики» справа, затем щелкнитеAdd/update a policy,Как показано ниже:

Затем добавьте стратегию, как показано ниже:

Значение каждого параметра следующее:

  • Имя: имя политики.
  • Шаблон: соответствующий шаблон для очереди (регулярное выражение).
  • Определение: определение изображения, есть три основных параметра: ha-mode, ha-params, ha-sync-mode.
    • ha-mode: определяет режим зеркальной очереди.Допустимые значения: все, точно и узлы. где all означает зеркалирование на всех узлах кластера (это значение по умолчанию); точно означает зеркалирование на указанном количестве узлов, количество узлов задается ха-парамами; узлы означает зеркалирование на указанном узле, указываются имена узлов через ха-парамс.
    • ha-params: параметры, необходимые для режима ha-mode.
    • ha-sync-mode: режим синхронизации сообщений в очереди. Допустимые значения — автоматический и ручной.
  • priority — необязательный параметр, указывающий приоритет политики.

После завершения настройки нажмите кнопкуadd/update policyкнопку, чтобы завершить добавление стратегии, следующим образом:

После завершения добавления мы можем выполнить простой тест.

Сначала убедитесь, что все три RabbitMQ запущены, а затем используйте указанный выше провайдер для отправки сообщения в очередь сообщений.

Закройте экземпляр mq01 после отправки.

Затем запустите потребитель.В это время обнаруживается, что потребитель может завершить потребление сообщений (обратите внимание на отличие от предыдущего обратного теста), что означает, что зеркальная очередь была успешно построена.

13.3.2 Зеркальная очередь конфигурации командной строки

Формат конфигурации командной строки следующий:

rabbitmqctl set_policy [-p vhost] [--priority priority] [--apply-to apply-to] {name} {pattern} {definition}

Возьмем простой пример конфигурации:

rabbitmqctl set_policy -p / --apply-to queues my_queue_mirror "^" '{"ha-mode":"all","ha-sync-mode":"automatic"}'

14. Закончить цветок~