пиши на фронт
Прошло много времени с тех пор, как я в последний раз публиковал статью. На самом деле, я не переставал писать за это время. Я просто занят поиском работы и окончанием школы. Обновите его ~
- Эта статья немного длинная, поэтому она разделена на две части.
- PS: Статья по Java Knowledge Q&A на Github не перестала писаться и будет обновляться одна за другой в последнее время.
Каталог статей:
1. По-простому
1.1 Что такое промежуточное ПО?
Определение IDC (Internet Data Center): промежуточное ПО — это независимая сервисная программа системного программного обеспечения.Распределенное прикладное программное обеспечение использует это программное обеспечение для совместного использования ресурсов между различными технологиями.Промежуточное ПО находится в операционной системе клиентского сервера.Управление вычислительными ресурсами и сетевыми коммуникациями.
Во-первых,Промежуточное программное обеспечение — это общий термин для определенного типа программного обеспечения, а не для конкретного программного обеспечения.. ЭтоМежду платформой (аппаратным обеспечением операционной системы) и приложениемОн скрывает различные сложности базовой операционной системы и снижает техническую нагрузку разработчиков.В то же время его дизайн не нацелен на конкретную цель, а предоставляет функциональные модульные сервисы с универсальными характеристиками.Эти сервисы имеют Стандартные программные интерфейсы и протоколы. также могут иметь разные реализации в зависимости от платформы.
Популярный пример (только для справки, не полностью соответствует):
- Я открываю кофейню, и есть n поставщиков кофейных зерен, таких как ABC и т. д., но я должен выбрать зерна по доступной цене и хорошего качества, но рынок колеблется под влиянием многих факторов, возможно, мой текущий выбор, который не лучший вариант через некоторое время. Поэтому я специально нашла рыночного посредника и попросила его помочь мне побеспокоиться об этом киоске.Я просто объяснила вам цену и требования к качеству, нужно только найти, а процесс меня совершенно не беспокоит. Концепция этого посредника аналогична промежуточному программному обеспечению.
1.1.1 Понятие распределенного (дополнительного)
Этот абзац взят из статьи, которую я написал ранее о начале работы с Dubbo.
Определения в Baidu и Wiki относительно профессиональные и малопонятные.В большинстве блогов или учебных пособий часто используются определения из «Принципов и парадигм распределенных систем», а именно:«Распределенная система — это набор независимых компьютеров, которые представляются пользователям как единая связанная система».
Давайте используем некоторое пространство, чтобы объяснить, что называется распределенным
1.1.1.1 Что такое централизованная система
Говоря о распределении, следует отметить, что«Централизованная система», эта концепция лучше всего понятна, это установка функций, программ и т. д. на одно и то же устройство, и это хост-устройство предоставляет услуги внешнему миру.
Для самого простого примера: вы берете хост ПК, конвертируете его в простой сервер и после настройки различного содержимого устанавливаете в него MySQL, веб-сервер, FTP, Nginx и т. д. После того, как проект упакован и развернут, сервисы могут быть предоставлены внешнему миру.Однако, как только возникает проблема с машиной, будь то программная или аппаратная, вся система будет серьезно вовлечена в ошибки.Яйца кладутся в одну корзину, и все они будут биты, если их надо бить.
1.1.12 Что такое распределенная система
Поскольку у централизованной системы есть такая проблема, которая влияет на все тело, одна из функций распределенной системы, естественно, состоит в том, чтобы решить такую проблему.Как мы знаем из определения, распределенная система в смысле опыта пользователя, точно так же, как традиционная единая система, какие-то изменения вносятся внутри самой системы, а ощущения у пользователей не очень
Например, крупные платформы электронной коммерции, такие как Taobao и Jingdong, имеют десятки тысяч хостов, иначе они не могут обрабатывать большой объем данных и запросов, какие конкретно подразделения и операции, о них мы поговорим ниже, а для мы, пользователи, не нуждаемся и не хотим заботиться об этом, мы все еще можем просто думать, что столкнулись с «хозяином» «Таобао».
Таким образом, относительно профессиональное заявление о распределении выглядит так (детализация процесса)Две или более программы выполняются на разных хост-процессах соответственно, и они координируют друг с другом выполнение общих функций, тогда систему, состоящую из этих программ, можно назвать распределенной системой.
- Обе программы одинаковые - распределенные
- Это все разные программы - кластеризация
Промежуточное программное обеспечение сообщений, как следует из его названия, является промежуточным программным обеспечением, используемым для обработки служб, связанных с сообщениями.Он обеспечивает канал для связи и взаимодействия между системами.Например, отправителю нужно только передать информацию, которая должна быть передана промежуточному программному обеспечению сообщения, и отправленный протокол, режим, сеть, сбой и другие проблемы в процессе отправки обрабатываются промежуточным программным обеспечением, поэтому оно отвечает за обеспечение надежной передачи информации.
Следовательно, промежуточное программное обеспечение сообщений — это технология, используемая для получения, хранения и отправки данных, которая предоставляет различные функции, которые могут реализовать высокую доступность и высокую надежность сообщений, а также обеспечить хороший механизм отказоустойчивости. Программа может значительно помочь в оккупации системных ресурсов и повышении эффективности передачи.
- Разные MQ имеют разные характеристики и направление, в котором они лучше.Нельзя сказать, кто хороший, а кто плохой, только кто больше подходит.
1.2.1 Сценарий приложения очереди сообщений
В соответствии с потребностями бизнеса он может иметь множество сценариев применения, таких как развязка, пиковое сглаживание, широковещание и т. д. Давайте рассмотрим два сценария, чтобы разобраться в простом процессе.
1.2.1.1 Разделение бизнеса
Недавно я думал о покупке нескольких книг для чтения. Возьмем пример покупки книги и размещения заказа. Когда я нажимаю, чтобы купить, может быть такая серия выполнения бизнес-логики, ① Вычесть емкость инвентаря ② Сгенерировать заказ ③ Оплатить ④ Обновить статус заказа ⑤ Отправить СМС об успешной покупке ⑥ Обновить статус экспресс-забора товара. На начальном этапе мы можем полностью выполнять эти услуги синхронно, но для повышения эффективности в дальнейшем мы можем разделить задачи, которые необходимо выполнить немедленно, и задачи, которые можно выполнить позже, например ⑤ Отправить сообщение об успешной покупке ⑥ Обновить статус экспресс-доставки товара, можно считать разным исполнением. После выполнения основного процесса эти службы, которые могут быть замедлены, можно определить как выполненные, отправив сообщение в MQ, чтобы гарантировать, что процесс завершится первым. Затем извлекайте сообщения MQ или MQ активно отправляет для асинхронного выполнения других служб.
1.2.1.2 Снятие пиков и заполнение впадин
Например, чтобы отправить сообщение-анонс с прочитанным и непрочитанным идентификатором, необходимо написать такое сообщение-анонс для каждого пользователя, например, хранить его в MongoDB, даже если MongoDB не может его поддерживать, он может написать миллионы или десятки миллионы записей мгновенно, поэтому вы можете рассмотреть возможность использования очереди сообщений. Например, в серверной системе Java мы можем использовать асинхронный многопоточный метод для отправки сообщений в очередь сообщений MQ, чтобы веб-система не занимала обычные операции CRUD базы данных при публикации сообщений объявлений. Системные сообщения хранятся в очереди сообщений, мы просто используем ее для сокращения пиков и заполнения впадин, а системные сообщения в конечном итоге будут храниться в базе данных. Таким образом, мы можем спроектировать таким образом, когда пользователь входит в систему, использовать асинхронный поток для получения системного сообщения пользователя из очереди сообщений MQ, а затем сохранить системное сообщение в базе данных и, наконец, сообщение в очереди сообщений. MQ автоматически удаляется. Из-за входа пользователя в непиковые часы задача записи сообщений в базу данных также становится непиковой записью.
1.3 Что такое RabbitMQ
RabbitMQ — это система очередей сообщений с открытым исходным кодом, написанная на языке Erlang и соответствующая протоколу AMQP.Он поддерживает несколько клиентов (языков) и используется для хранения и пересылки сообщений в распределенных системах.Он обладает высокой доступностью, высокой масштабируемостью и легким доступом.Удобство использования и другие характеристики.
Для получения более подробной информации, пожалуйста, посетите официальный сайт:
Короче говоря, это обычная очередь сообщений, и ее характеристики будут объяснены по порядку позже.Мы сначала начинаем с части загрузки и установки записи, а затем используем ее.
2. Загрузка и установка
Вообще говоря, методы установки включают ручную установку и установку Docker.В большинстве сценариев используется установка Docker, но в качестве этапа обучения, если вы не особо беспокоитесь, неплохо изучить ручную установку.
Примечание. Доступны как облачные серверы, так и виртуальные машины.Продемонстрированная версия Linux — CentOS 7.9.
2.1 Ручная установка
2.1.1 Процесс загрузки и установки
Примечание: вы можете загрузить и установить непосредственно через yum в Linux.Здесь вы выбираете загрузку файла на свой собственный хост Windows, а затем загружаете его в Linux через FTP для прямой установки. Можно избежать некоторых проблем с загрузкой, вызванных сетью на виртуальной машине.
-
Сначала откройте каталог загрузки на официальном сайте, а затем выберите версию в соответствии с вашей версией Linux.
-
Поскольку RabbitMQ написан на языке Erlang, вам также необходимо предоставить среду Erlang, а затем загрузить Erlang.
- адрес:www.erlang-solutions.com/downloads
- A: Скорость доступа к этому веб-сайту очень низкая, пожалуйста, наберитесь терпения, иначе вам придется повесить лестницу.
- B: Версия Erlang должна соответствовать RabbitMQ (как показано, существуют максимальные и минимальные ограничения версии)
- Адрес просмотра версии:woohoo.rabbitcurrent.com/what-while blue…
- Здесь выбраны RabbitMQ 3.8.14 и Erlang 23.2.3
- адрес:www.erlang-solutions.com/downloads
[root@centos7 rabbitmq]# ls
esl-erlang_23.2.3-1_centos_7_amd64.rpm rabbitmq-server-3.8.14-1.el7.noarch.rpm
[root@centos7 rabbitmq]# pwd
/usr/local/bin/rabbitmq
- Установите Erlang, Socat и RabbitMQ
- Erlang и Socat зависят от RabbitMQ.
# 安装 Erlang,安装后执行 erl -v 显示版本号则代表成功
rpm -ivh esl-erlang_23.2.3-1_centos_7_amd64.rpm
# 安装 Socat 这里没有下载源文件,而是直接通过 yum 在线安装,因为它并不大
yum install -y socat
# 安装 RabbitMQ
rpm -ivh rabbitmq-server-3.8.14-1.el7.noarch.rpm
- После завершения установки запустите службу, чтобы проверить, можно ли успешно запустить RabbitMQ.
# 启动服务
systemctl start rabbitmq-server
# 开机自启
systemctl enable rabbitmq-server
# 停止服务
systemctl stop rabbitmq-server
# 查看服务状态
systemctl status rabbitmq-server.service
Как показано на рисунке, установка начинается успешно
Если установка неверна, см.:
- Установка RabbitMQ в Linux для решения различных проблем
- Rabbitmq ОШИБКА: ошибка epmd для хоста deb: адрес (не удается подключиться к хосту/порту) решение
2.1.2 Настройка управления веб-интерфейсом
Описанная выше установка фактически завершена, но RabbitMQ предоставляет нам веб-интерфейс управления, который недоступен по умолчанию и требует установки.
- Установите плагин веб-управления и перезапустите службу.
# 安装命令
rabbitmq-plugins enable rabbitmq_management
# 重启服务
systemctl restart rabbitmq-server
- Обязательно откройте порт 15672 брандмауэра Linux, иначе он будет недоступен.На этапе обучения вы даже можете перейти к команде запроса, чтобы отключить брандмауэр
- Соответствующий сервер (Alibaba Cloud, Tencent Cloud и т. д.) должен открыть порт 15672 в группе безопасности.
- Доступ к Linux IP:15672 , например.
http://192.168.122.1:15672
# 查询 15672 是否开放,一般默认都是 no
firewall-cmd --query-port=15672/tcp
# 开放指定端口 15672
firewall-cmd --add-port=15672/tcp --permanent
# 重新载入
firewall-cmd --reload
# 再次查询,结果就是 yes 了
firewall-cmd --query-port=15672/tcp
-
Добавить учетную запись для удаленного входа
- RabbitMQ имеет учетную запись и пароль по умолчанию для гостя, но доступ к ним возможен только с локального хоста.
# 新增用户 用户名和密码都是 admin
rabbitmqctl add_user admin admin
- Добавить разрешения для учетных записей удаленного входа
- администратор (суперадминистратор): вход в консоль, просмотр всей информации, управление пользователями, управление политиками
- мониторинг (monitor): войти в консоль, просмотреть всю информацию
- policymaker: войти в консоль, указать политики
- управление (обычный администратор): войти в консоль
# 设置用户分配操作权限,admin 用户的权限为 administrator
rabbitmqctl set_user_tags admin administrator
- Добавить права доступа к ресурсам для пользователей
- Поскольку администратор уже является суперадминистратором, вы можете сделать это без назначения разрешений на ресурсы, и это будет сделано по умолчанию.
# 命令格式为: set_permissions [-p <vhostpath>] <user> <conf> <write> <read>
# 这里即为 admin 用户开启 配置文件和读写的权限
rabbitmqctl set_permissions -p / admin ".*"".*"".*"
- Доступ к Linux IP:15672 , например.
http://192.168.122.1:15672, введите имя пользователя и пароль администратора, только что установленные- Как показано на рисунке: доступ успешен
2.1.2.1 Команда резюме
- Добавить пользователя:
rabbitmqctl add_user <username> <password> - изменить пароль:
rabbitmqctl change_password <username> <newpass> - удалить пользователей:
rabbitmqctl delete_user <username> - список пользователей:
rabbitmqctl list_users - Установить роли пользователей:
rabbitmqctl set_user_tags <username> <tag1,tag2> - Удалить все роли пользователя:
rabbitmqctl set_user_tags <username> - Добавьте права доступа к ресурсам для пользователей:
set_permissions [-p <vhostpath>] <user> <conf> <write> <read>
Использование: введите rabbitmqctl, вам будут предложены возможные команды, а затем используйте rabbitmqctl hepl
2.1.3 Краткое введение в управление веб-интерфейсом
- Соединения: используется для управления производителями и потребителями после установления соединения с RabbitMQ.
- Каналы: после установления соединения будет сформирован канал, и доставка сообщения будет зависеть от канала.
- Обмены: используется для реализации маршрутизации сообщений.
- Очереди (queues): Очереди, в которых хранятся сообщения, сообщения ожидают использования и удаляются из очереди после потребления.
- Администратор (управление): используется для установки пользователя управления и соответствующих разрешений, как показано на следующем рисунке.
Теги используются для указания роли пользователя
- администратор (суперадминистратор): вход в консоль, просмотр всей информации, управление пользователями, управление политиками
- мониторинг (monitor): войти в консоль, просмотреть всю информацию
- policymaker: войти в консоль, указать политики
- управление (обычный администратор): войти в консоль
2.2 Установка докера
При установке RabbitMQ в Docker не нужно учитывать различные конфликты и проблемы несовместимости, такие как версии, среды и т. д. Это очень удобно.Виртуальная машина, которую я продемонстрировал, представляет собой голое железо CentOS 7.9, поэтому мы обновим yum, установим Docker и установим RabbitMQ и поговорим о шагах
2.2.1 Настройка юм
- Обновите юм до последней версии
# 更新 yum
yum update
# 检查yum依赖的几个包 yum-utils 提供 yum-config-manager 功能, 后面两个是 devicemapper 用到的
yum install -y yum-utils device-mapper-persistent-data lvm2
- Установите источник yum на Alibaba Cloud
yum-config-manager --add-repo http://mirrors.aliyun.com/docker-ce/linux/centos/docker-ce.repo
2.2.2 Установка докера
2.2.2.1 Шаги
- Установить докер с помощью yum
- docker-ce — это версия сообщества смысла, ee Enterprise Edition
yum install docker-ce -y
- Проверьте, прошла ли установка успешно, посмотрев версию
docker -v
- Ускорение образа Docker (здесь нужно заменить на свой)
sudo mkdir -p /etc/docker
sudo tee /etc/docker/daemon.json <<-'EOF'
{
"registry-mirrors": ["https://<你的ID>.mirror.aliyuncs.com"]
}
EOF
sudo systemctl daemon-reload
sudo systemctl restart docker
Иногда бывает сложно получить образы из DockerHub в Китае, но в это время вы можете настроить ускоритель образов. Docker официально и многие отечественные провайдеры облачных услуг предоставляют внутренние ускорители, такие как:
- ХКУСТ Зеркало:docker.mirrors.ustc.edu.cn/
- NetEase:hub-mirror.c.163.com/
- Али Клауд:https://.mirror.aliyuncs.com
- Облачный ускоритель Qiniu:reg-mirror.qiniu.com
Если после настройки определенного адреса ускорителя вы обнаружите, что изображение не может быть извлечено, переключитесь на другой адрес ускорителя. Все основные поставщики облачных услуг в Китае предоставляют услугу ускорения образов Docker.Рекомендуется выбрать соответствующую услугу ускорения изображений в соответствии с облачной платформой, на которой работает Docker.
Адрес получения образа Alibaba Cloud:Adult.console.aliyun.com/cai-Ханчжоу…
2.2.2.2 Общие команды Docker
2.2.2.2.1 Административные команды
- Просто запустите, остановите, перезапустите эти простые команды, также можно использовать службу, systemctl немного мощнее
# 启动 docker
systemctl docker start
# 停止 docker
systemctl docker stop
# 重启 docker
systemctl docker restart
# 查看 docker 状态
systemctl status docker
# 开机自启
systemctl enable docker
systemctl unenable docker
2.2.2.2.2 Команда «Зеркало»
# 导入镜像文件
docker load < xxx.tar.gz
# 查看安装的镜像
docker images
# 删除镜像
docker rmi 镜像名
2.2.3 Установите RabbitMQ (необязательно)
Примечание: будет лучше установить сразу с 2.2.3.2 в одном предложении
2.2.3.1 Пошаговая установка
- Получите зеркало RabbitMQ
docker pull rabbitmq:management
- Создать и запустить контейнер (указано в 3 в 3)
docker run -id --name 容器名 -p 15672:15672 -p 5672:5672 rabbitmq:management
2.2.3.2 Установка одним предложением
Вышеупомянутый метод установки заключается в том, чтобы сначала получить образ RabbitMQ, а затем начать установку. Здесь нет проблем. Возникнет проблема при создании, потому что мы хотим установить управление, которое является его веб-управлением. Если мы не сделаем некоторые обработки, он будет установлен по умолчанию.Пользователя нет, поэтому нужно зайти и настроить самому как и раньше, а Docker Hub привел пример нашей настройки, то есть с использованием-eДелегировать конфигурацию, использоватьRABBITMQ_DEFAULT_USERиRABBITMQ_DEFAULT_PASSНастроить имя пользователя и пароль
Для получения дополнительной информации, пожалуйста, ознакомьтесь с главами о настройке пользователя и пароля по умолчанию в Docker Hub.
registry.hub.docker.com/_/rabbitmq/
- выполнить установку
docker run -di --name myrabbitmq -e RABBITMQ_DEFAULT_USER=admin -e RABBITMQ_DEFAULT_PASS=admin -p 15672:15672 -p 5672:5672 -p 25672:25672 -p 61613:61613 -p 1883:1883 rabbitmq:management
- Проверьте состояние контейнера, чтобы убедиться, что операция прошла успешно.
# 查看容器运行状态
docker ps -a
# 启动
docker start 容器名
# 停止
docker stop 容器名
# 退出命令行,不停止
exit
# 进入到node容器(如果开启了 -t 的情况)
docker exec -it 容器名 bash
2.2.3.2.1 Введение параметров
Описание этих параметров поясняется ниже:
-
-i: указывает на запуск контейнера. -
-t: указывает способ (командную строку) зарезервировать взаимодействие для контейнера, т. е. выделить псевдотерминал. так часто вижу-itТакой матч. -
--name: Дайте контейнеру имя. -
-v: указывает отношение сопоставления каталогов (первый — это каталог хоста, второй — каталог, сопоставленный с хостом), вы можете использовать несколько-vСделайте несколько сопоставлений каталогов или файлов. Примечание. Рекомендуется выполнить сопоставление каталогов, изменить его на хосте, а затем поделиться им в контейнере. -
-d: Указывает, что контейнер демона создан для работы в фоновом режиме (так что контейнер не будет автоматически регистрироваться после создания контейнера, если вы добавите только-i-tДва параметра, которые будут автоматически входить в контейнер после создания), то есть бэкенд приостанавливается и запускается. -
-p: указывает сопоставление портов, первый — это порт хоста, а второй — сопоставленный порт в контейнере. можно использовать несколько-pВыполняйте сопоставление нескольких портов только после того, как сопоставление портов станет доступным для внешнего мира.
Позвольте мне привести Вам пример:
# 创建容器,把容器 3000 端口映射到宿主机 3000 端口,把/demo映射到宿主机的/demo face是我下载好的一个现成的镜像
docker run -d -it -p 3000:3000 -v /demo:/demo --name node face
# 例如,名为 node 的镜像中有一个需要执行的 python 程序,就可以通过如下命令进入刚才分配到的命令行中去执行这个程序
docker exec -it node bash
-
из-за использования
-tТаким образом, этот параметр может быть назначен псевдотерминалу черезdocker exec -it 容器名 bashвведите командную строку -
-vПосле того, как директория будет сопоставлена, после входа в контейнер также будет идентичная демо-папка, например, в ней может быть выполнена программа python
2.2.3.2.1 Описание порта
4369: Эрланг обнаружил порт
5672: клиентский порт связи.
15672: порт интерфейса администратора
25672: Внутренний порт связи между серверами.
61613: клиент STOMP без и с TLS
1883: клиент MQTT без и с включенным TLS
Более критичны 5672 и 15672
Для получения дополнительной информации о порте посетите официальную документацию веб-сайта.
Примечание. Если вы хотите подключиться удаленно, например, чтобы получить доступ к порту 15672 веб-страницы управления и порту 5672 подключения клиента Java, вы должны выполнить операцию открытия, иначе вы не сможете подключиться.
- Ниже приведен пример открытия порта 15672 на базе CentOS 7.9.
# 查询 15672 是否开放,一般默认都是 no
firewall-cmd --query-port=15672/tcp
# 开放指定端口 15672
firewall-cmd --add-port=15672/tcp --permanent
# 重新载入
firewall-cmd --reload
# 再次查询,结果就是 yes 了
firewall-cmd --query-port=15672/tcp
- Вот команда для отключения брандмауэра
systemctl disable firewalld
systemctl stop firewalld
3. Протокол и модель RabbitMQ
После завершения установки необходимо войти в тему, то есть использовать код Java или Springboot для реализации нескольких способов RabbitMQ, но если вы хотите хорошо разобраться в этих методах коммутации маршрутизации, вам необходимо понять его протокол и архитектуру модель.
3.1 Соглашение
3.1.1 Что такое соглашение?
Протокол, сокращение от сетевого протокола, сетевой протокол представляет собой набор соглашений, которые должны соблюдать обе стороны коммуникационного компьютера. Например, как установить соединение, как идентифицировать друг друга и так далее. Только при соблюдении этого соглашения компьютеры могут взаимодействовать друг с другом. Его три элемента: синтаксис, семантика, время.
Для того, чтобы данные перемещались от источника к получателю в сети, участники сетевого взаимодействия должны следовать одним и тем же правилам.Этот набор правил называется протоколом, что в конечном итоге отражается в формате пакетов данных, передаваемых по сети. сеть.
3.1.1.1 Три элемента сетевого протокола
- Синтаксис: структура и формат данных и управляющей информации, а также порядок, в котором данные появляются.
- Семантика: Объясняет значение каждой части управляющей информации и указывает, какую управляющую информацию необходимо отправить и какой ответ на выполненное действие.
- Время: подробное описание последовательности, в которой происходят события.
Люди ярко описывают эти три элемента: что делать, как делать и в каком порядке.
Например, HTTP-протокол
Синтаксис: HTTP определяет формат пакетов запросов и ответов. Семантика: клиент активно инициирует запрос, который называется запросом, а сервер возвращает данные, которые называются ответом. Время: запрос соответствует ответу, и есть запрос до ответа.
3.1.1.1.1 Вопрос интервью: Почему промежуточное ПО сообщений не использует протокол HTTP напрямую?
Для промежуточного программного обеспечения сообщений его основная ответственность заключается в том, чтобы отвечать за передачу данных, хранение, распространение, высокую производительность и простоту - это то, к чему мы стремимся, в то время как заголовки HTTP-запросов и заголовки ответов являются более сложными, включая файлы cookie. , шифрование и дешифрование данных, окно подоконник, код ответа и другие дополнительные функции, такие сложные функции нам не нужны.
At the same time, in most cases, HTTP is mostly short links. In the actual interaction process, a request to a response is likely to be interrupted. After the interruption, the persistence will not be performed, which will cause the loss of the запрос.这样就不利于消息中间件的业务场景,因为消息中间件可能是一个长期的获取信息的过程,出现问题和故障要对数据或消息执行持久化等,目的是为了保证消息和数据的高可靠和稳健的运行
3.1.2 Протокол AMQP RabbitMQ
Протокол, используемый RabbitMQ, — это AMQP (расширенный протокол организации очереди сообщений), который был предложен в 2003 году и впервые использовался для решения проблемы взаимодействия доставки сообщений между различными платформами в финансовой отрасли.
AMQP — это, точнее, двоичный протокол проводного уровня (протокол канала). В этом его существенное отличие от JMS: AMQP не ограничивается уровнем API, а напрямую определяет формат данных для сетевого обмена. Это делает поставщика (производителя), который реализует AMQP, по своей сути кроссплатформенным.
По сравнению с другими протоколами сообщений его характеристики:
- Поддержка распределенных транзакций
- Поддержка сохранения сообщений
- Высокопроизводительные и надежные преимущества обмена сообщениями
3.1.3 Модель архитектуры
Если вы хотите изучить конкретные режимы отправки следующих сообщений, вы должны четко понимать диаграмму модели, потому что эти методы являются выбором и уменьшением модели в той или иной степени.
-
Connection: сетевое соединение между приложением и брокером. -
Channel: Канал, то есть канал для передачи информации.Можно установить несколько Каналов, и каждый Канал представляет задачу сеанса.- Канал — это виртуальное соединение, установленное внутри TCP-соединения, и чтение и запись информации передаются через канал.Поскольку операционной системе очень дорого устанавливать и уничтожать TCP, концепция канала вводится для повторного использования. TCP-соединение.
-
Broker(Server): Определите объекты сервера очереди сообщений, например, здесь, сервер Rabbitmq. -
Virtual Host: Виртуальный хост.В одном брокере можно настроить несколько виртуальных хостов для изоляции разрешений разных пользователей.- Брокера можно понимать как всю службу базы данных, а Виртуальный хост — это ощущение каждой базы данных.Разные проекты могут соответствовать разным базам данных, включая бизнес-таблицу, к которой принадлежит проект, и так далее.
- В каждом виртуальном хосте может быть несколько обменов и очередей.
-
Exchange: Обмен, который получает сообщения, отправленные производителями, и отправляет их в очереди на основе ключей маршрутизации. -
Binding: виртуальное соединение между Exchange и Queue. Binding может включать в себя несколько ключей маршрутизации. -
Routing key: Правила маршрутизации, используйте его, чтобы подтвердить, что виртуальная машина направляет конкретное сообщение. -
Queue: очередь сообщений, которая является контейнером для сообщений, используемым для сохранения сообщений, каждое сообщение может быть передано в одну или несколько очередей, ожидающих потребления потребителями, то есть извлечения сообщения. -
Consumer: потребитель сообщения (программа, получающая сообщение).
4. Java реализует RabbitMQ
4.1 Окружающая среда Строительство
На официальном сайте представлено несколько моделей:ву ву ву.кролик в настоящее время.com/начать.…
До сих пор на официальном сайте было представлено в общей сложности моделей 7. Мы в основном представляем первые пять основных режимов.Некоторые люди классифицировали режимы Direct и Topic в режим маршрутизации, который также можно рассматривать как четыре типа.
4.1.1 Создание проекта Java
Сначала создайте проект Maven, в котором не используются скелеты, затем введите зависимости RabbitMQ, а также зависимости модульного тестирования.
<dependency>
<groupId>com.rabbitmq</groupId>
<artifactId>amqp-client</artifactId>
<version>5.10.0</version>
</dependency>
<dependency>
<groupId>junit</groupId>
<artifactId>junit</artifactId>
<version>4.11</version>
</dependency>
4.1.2 Создайте виртуальный хост (необязательно)
Здесь мы создаем новый виртуальный хост для обслуживания этого Java-проекта, вы также можете создать нового пользователя, а затем включить доступ к виртуальным хостам (то есть привязать виртуальный хост к пользователю). Мы по-прежнему используем admin (пользователь с правами администратора, которого я создал ранее) для демонстрации.
4.1.3 Создание класса инструмента подключения
Поскольку позже мы продемонстрируем различные примеры, и каждый код операции, такой как получение соединения, освобождение соединения, закрытие ресурсов и т. д., является согласованным, чтобы предотвратить избыточность кода, оптимизировать код и упростить его понимание, класс инструмента извлекается, так что все будут сосредоточены на простом сравнении различных реализаций.
- Вспомогательный класс RabbitMqUtil
public class RabbitMqUtil {
/**
* 主机名 即 Linux IP地址
*/
private static String host = "";
/**
* 端口号 客户端访问默认都是 5672
*/
private static int port = 0;
/**
* 虚拟主机 可以设置为默认的 / 或者自己创建出指定的虚拟主机
*/
private static String virtualHost = "";
/**
* 用户名
*/
private static String username = "";
/**
* 密码
*/
private static String password = "";
// 使用静态代码块为Properties对象赋值
static {
try {
//实例化对象
Properties properties = new Properties();
//获取properties文件的流对象
InputStream in = RabbitMqUtil.class.getClassLoader().getResourceAsStream("rabbitmq.properties");
properties.load(in);
// 分别获取 value
host = properties.getProperty("host");
port = Integer.parseInt(properties.getProperty("port"));
virtualHost = properties.getProperty("virtualHost");
username = properties.getProperty("username");
password = properties.getProperty("password");
} catch (Exception e) {
e.printStackTrace();
}
}
/**
* 获取连接
*
* @return 连接
*/
public static Connection getConnection() {
try {
// 创建连接工厂
ConnectionFactory connectionFactory = new ConnectionFactory();
// 设置连接 rabbitmq 主机
connectionFactory.setHost(host);
// 设置端口号
connectionFactory.setPort(port);
// 设置连接的虚拟主机(数据库的感觉)
connectionFactory.setVirtualHost(virtualHost);
// 设置访问虚拟主机的用户名和密码
connectionFactory.setUsername(username);
connectionFactory.setPassword(password);
// 返回一个新连接
return connectionFactory.newConnection();
} catch (Exception e) {
e.printStackTrace();
}
return null;
}
/**
* 关闭通道和释放连接
*
* @param channel channel
* @param connection connection
*/
public static void close(Channel channel, Connection connection) {
try {
if (channel != null) {
channel.close();
}
if (connection != null) {
connection.close();
}
} catch (Exception e) {
e.printStackTrace();
}
}
}
- properties
host=192.168.122.1
port=5672
virtualHost=/rabbitmq_maven_01
username=admin
password=admin
4.2 Пять методов реализации
инструкция:
- Строковое содержимое, такое как имя очереди, сообщение и т. д., лучше определять как переменную для передачи, в моей статье это написано прямо в параметрах, это магическое значение не очень красиво.
- Модульный тест Junit используется в производителе, но основная функция написана в потребителе. Это потому, что мы хотим, чтобы потребитель находился в состоянии непрерывного выполнения и ожидания. Использование Junit приведет к завершению программы после одного выполнения.
- В дополнение к записи в функции main вы также можете рассмотреть возможность использования sleep для ожидания или while(true), чтобы программа не завершалась напрямую.
4.2.1 Простой режим очереди (Hello Word)
-
Queue: Очередь сообщений, понимаемая как контейнер, производитель отправляет в нее сообщения, он хранит сообщения и ждет, пока потребители их потреблят. -
Consumer: потребитель сообщения (программа, получающая сообщение).
4.2.1.1 Как понять
Как показано на рисунке, в режиме простой очереди производитель, проходящий через очередь, соответствует потребителю. Его можно рассматривать как метод передачи «точка-точка».По сравнению со схемой модели в 3.1.3, основная особенность заключается в том, что Exchange (коммутатор) и routekey (ключ маршрутизации) не видны. Именно потому, что этот режим простой, поэтому не требует сложного условного распределения и т. д., поэтому пользователям не нужно явно рассматривать вопрос о коммутаторах и ключах маршрутизации.
- Однако следует отметить, что этот режим напрямую не подключает производителя к очереди, а использует переключатель по умолчанию, который отправит сообщение в очередь с тем же именем, что и ключ маршрута, который также является позицией routekey в коде. Причина для имени очереди заполнена
4.2.1.2 Реализация кода
4.2.1.2.1 Код продюсера
public class Producer {
@Test
public void sendMessage() throws IOException, TimeoutException {
// 通过工具类获取连接
Connection connection = RabbitMqUtil.getConnection();
// 获取连接通道
Channel channel = connection.createChannel();
// 通道绑定消息队列
channel.queueDeclare("queue1",false,false,false,null);
// 发布消息
channel.basicPublish("","queue1",null,"This is rabbitmq message 001 !".getBytes());
// 通过工具关闭channel和释放连接
RabbitMqUtil.close(channel,connection);
}
}
- Получить соединение через класс инструмента
- получить канал подключения: Согласно схеме модели в 3.1.3, производителю необходимо получить канал после установления соединения, чтобы получить доступ к следующим очередям обмена и т. д.
-
очередь сообщений привязки канала: Перед привязкой очереди следует привязать переключатель, но концепция переключателя в этом режиме скрыта, а за ним используется переключатель по умолчанию, поэтому очередь привязывается напрямую.
- Объяснение метода queueDeclare
- Параметр 1: очередь (название очереди), если очереди нет, то она будет создана автоматически.
- Параметр 2: устойчивый (независимо от того, является ли очередь постоянной), постоянство может гарантировать, что очередь все еще существует после перезапуска сервера.
- Параметр 3:эксклюзив (исключительная очередь) — указывает, является ли очередь эксклюзивной.Если этот пункт имеет значение true, очередь видна только тому соединению, которое объявило ее в первый раз, и автоматически удаляется при разрыве соединения.
- Параметр 4: autoDelete (автоматическое удаление), после того как последний потребитель израсходует сообщение, очередь автоматически удаляется.
- Параметр 5: аргументы (несет дополнительные атрибуты).
- Объяснение метода queueDeclare
-
опубликовать новость: Здесь можно указать способ отправки и содержимое очереди сообщений.Поскольку этот режим относительно прост, он не включает все параметры.Следующие режимы будут подробно описаны.
- Объяснение метода basicPublish
- Параметр 1: exchange (имя биржи).
- Параметр 2: routingKey (ключ маршрутизации), здесь заполните имя очереди, что можно понимать как отправку сообщения в очередь с тем же именем, что и у routekey.
- Параметр 3: реквизит (состояние управления сообщением), где вы можете контролировать сохранение сообщения.
- Параметры: MessageProperties.PERSISTENT_TEXT_PLAIN
- Параметр 4: body (тело сообщения), тип — массив байтов, тип нужно преобразовать.
- Объяснение метода basicPublish
- Закрывайте каналы и освобождайте соединения с помощью инструментов: сначала закройте канал, затем разорвите соединение.
4.2.1.2.2 Код потребителя
public class Consumer {
public static void main(String[] args) throws IOException, TimeoutException{
// 通过工具类获取连接
Connection connection = RabbitMqUtil.getConnection();
// 获取连接通道
Channel channel = connection.createChannel();
// 通道绑定消息队列
channel.queueDeclare("queue1", false, false, false, null);
// 消费消息
channel.basicConsume("queue1", true, new DefaultConsumer(channel) {
@Override
public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
System.out.println("new String(body): " + new String(body));
}
});
}
}
-
Получить соединение через класс инструмента
-
получить канал подключения
-
очередь сообщений привязки канала
-
Использование сообщений: используется для указания, какую очередь сообщений использовать, а также некоторые механизмы и обратные вызовы.
- Объяснение метода basicConsume
- Параметр 1: очередь (название очереди), то есть сообщение, какую очередь потреблять.
- Параметр 2: autoAck (автоматический ответ), чтобы запустить механизм автоматического подтверждения сообщения и удалить сообщение из очереди, пока оно используется.
- Параметр 3: callback (интерфейс обратного вызова при потреблении), тип обратного вызова Consumer, здесь используется DefaultConsumer, который является классом реализации Consumer. Среди них, переписав метод handleDelivery, можно получить содержимое данных потребления.Здесь в основном используется тело, то есть для просмотра тела сообщения.Остальные три параметра пока не используются.Если вы заинтересованы, вы можете сначала распечатать это общее понимание.
- Объяснение метода basicConsume
4.2.2 Режим рабочей очереди (рабочая очередь)
-
Producer: производитель сообщения (программа, отправившая сообщение). -
Queue: Очередь сообщений, понимаемая как контейнер, производитель отправляет в нее сообщения, он хранит сообщения и ждет, пока потребители их потреблят. -
Consumer: потребитель сообщения (программа, получающая сообщение).- Здесь мы предполагаем, что Потребитель 1, Потребитель 2, Потребитель 3 соответственно выполняют задачу не так быстро, как скорость потребителя, что привело бы к сосредоточению внимания на проблемах этого шаблона.
4.2.2.1 Как понять
Рабочий режим можно увидеть на рисунке, то есть на основе режима простой очереди добавляются несколько потребителей, то есть несколько потребителей привязаны к одной и той же очереди для совместного потребления, что может решить проблему в простом режим очереди Скорость производства намного больше скорости потребления, что приводит к накоплению сообщений.
- Поскольку сообщение исчезает после использования, нет необходимости беспокоиться о повторении задачи.
4.2.2.2 Реализация кода
Примечание. Существует два режима рабочей очереди.
- Режим опроса: каждый потребитель делит сообщение поровну
- Модель справедливого распределения (больше работы для тех, кто может): распределение по способностям, больше распределения с высокой скоростью обработки и меньше распределения с медленной скоростью обработки.
Сначала мы демонстрируем режим опроса, который может привести к справедливому режиму распределения в соответствии с его недостатками.
Ниже описаны только те части, которые отличаются от вышеперечисленных.В простом режиме эти основные методы были введены
4.2.2.2.1 Режим опроса — код производителя
public class Producer {
@Test
public void sendMessage() throws IOException, TimeoutException {
// 通过工具类获取连接
Connection connection = RabbitMqUtil.getConnection();
// 获取连接通道
Channel channel = connection.createChannel();
// 通道绑定消息队列
channel.queueDeclare("work", true, false, false, null);
for (int i = 1; i <= 20; i++) {
// 发布消息
channel.basicPublish("", "work", null, (i + "号消息").getBytes());
}
// 通过工具关闭channel和释放连接
RabbitMqUtil.close(channel, connection);
}
}
Процесс в основном такой же, как и в режиме простой очереди, с некоторыми незначительными изменениями.Производитель в основном добавляет слой циклов.Поскольку существует несколько потребителей, отправляется больше сообщений, и можно увидеть некоторые функции и проблемы.
4.2.2.2.2 Режим опроса — код потребителя
- Потребитель 1
public class Consumer1 {
public static void main(String[] args) throws IOException {
// 通过工具类获取连接
Connection connection = RabbitMqUtil.getConnection();
// 获取连接通道
final Channel channel = connection.createChannel();
// 通道绑定消息队列
channel.queueDeclare("work", true, false, false, null);
// 消费消息
channel.basicConsume("work", true, new DefaultConsumer(channel) {
@Override
public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
try {
Thread.sleep(2000);
} catch (InterruptedException e) {
e.printStackTrace();
}
System.out.println("消费者1号:消费-" + new String(body));
}
});
}
}
- Потребитель 2
public class Consumer2 {
public static void main(String[] args) throws IOException {
// 通过工具类获取连接
Connection connection = RabbitMqUtil.getConnection();
// 获取连接通道
final Channel channel = connection.createChannel();
// 通道绑定消息队列
channel.queueDeclare("work", true, false, false, null);
// 消费消息
channel.basicConsume("work", true, new DefaultConsumer(channel) {
@Override
public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
System.out.println("消费者2号:消费-" + new String(body));
}
});
}
Вышеупомянутые два потребителя включили автоматический ответ Ack в basicConsume, который будет подробно описан ниже.В то же время, в потребителе 1 добавлен оператор sleep 2s, чтобы имитировать, что потребитель 1 обрабатывает сообщения медленно, в то время как потребитель 2 обрабатывает сценарии. с быстрыми сообщениями.
результат операции:
- Consumer1
消费者1号:消费-1号消息
消费者1号:消费-3号消息
消费者1号:消费-5号消息
消费者1号:消费-7号消息
消费者1号:消费-9号消息
消费者1号:消费-11号消息
消费者1号:消费-13号消息
消费者1号:消费-15号消息
消费者1号:消费-17号消息
消费者1号:消费-19号消息
- Consumer2
消费者2号:消费-2号消息
消费者2号:消费-4号消息
消费者2号:消费-6号消息
消费者2号:消费-8号消息
消费者2号:消费-10号消息
消费者2号:消费-12号消息
消费者2号:消费-14号消息
消费者2号:消费-16号消息
消费者2号:消费-18号消息
消费者2号:消费-20号消息
Наблюдение за процессом выполнения: обнаружено, что, хотя каждый из двух потребителей обрабатывает половину сообщений в конце, и они распределяются по одному человеку на человека, скорость обработки потребителя № 2 высока, и все они обрабатывается сразу, а потребитель № 1, каждая обработка занимает 2 с., следовательно, он может обрабатываться только медленно, а потребитель № 2 находится в ситуации холостого расточительства.
Как перейти в режим справедливого распределения?
Это связано со вторым параметром в basicConsume, включающим автоматическое подтверждение потребления. По умолчанию он равен true, что означает, что как только я получу сообщение, разосланное потребителю в очереди, я автоматически верну подтверждение потребления. сообщение в очереди будет автоматически удалено после того, как очередь получит его.
- Но в этом есть очень важная проблема.Этот метод заключается в передаче риска потребителю.Например потребитель получает 10 сообщений которые ему нужно обработать,а просто потребляет 4.Потребитель вылетает и зависает.Следующие 6 сообщения теряются.
Если вы хотите модифицировать раздачу по способностям, есть два момента
-
Настройте канал так, чтобы он потреблял только одно сообщение за раз
-
Отключить автоматическое подтверждение сообщений и ручное подтверждение сообщений
4.2.2.2.3 Модель справедливого распределения — Код производителя
public class Producer {
@Test
public void sendMessage() throws IOException, TimeoutException {
// 通过工具类获取连接
Connection connection = RabbitMqUtil.getConnection();
// 获取连接通道
Channel channel = connection.createChannel();
// 一次只发送一条消息
channel.basicQos(1);
// 通道绑定消息队列
channel.queueDeclare("work", true, false, false, null);
for (int i = 1; i <= 20; i++) {
// 发布消息
channel.basicPublish("", "work", null, (i + "号消息").getBytes());
}
// 通过工具关闭channel和释放连接
RabbitMqUtil.close(channel, connection);
}
4.2.2.2.4 Модель справедливого распределения — Кодекс потребителей
- Потребитель 1
public class Consumer1 {
public static void main(String[] args) throws IOException {
// 通过工具类获取连接
Connection connection = RabbitMqUtil.getConnection();
// 获取连接通道
final Channel channel = connection.createChannel();
// 一次只接受一条未确认的消息
channel.basicQos(1);
// 通道绑定消息队列
channel.queueDeclare("work", true, false, false, null);
// 消费消息
channel.basicConsume("work", false, new DefaultConsumer(channel) {
@Override
public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
try {
Thread.sleep(2000);
} catch (InterruptedException e) {
e.printStackTrace();
}
System.out.println("消费者1号:消费-" + new String(body));
// 返回 deliveryTag 代表队列可以删除此消息了
channel.basicAck(envelope.getDeliveryTag(), false);
}
});
}
}
- Потребитель 2
public class Consumer2 {
public static void main(String[] args) throws IOException {
// 通过工具类获取连接
Connection connection = RabbitMqUtil.getConnection();
// 获取连接通道
final Channel channel = connection.createChannel();
//步骤一:一次只接受一条未确认的消息
channel.basicQos(1);
// 通道绑定消息队列
channel.queueDeclare("work", true, false, false, null);
// 消费消息
channel.basicConsume("work", false, new DefaultConsumer(channel) {
@Override
public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
System.out.println("消费者2号:消费-" + new String(body));
channel.basicAck(envelope.getDeliveryTag(), false);
}
});
}
результат операции:
- Consumer1
消费者1号:消费-1号消息
- Consumer2
消费者2号:消费-2号消息
消费者2号:消费-3号消息
消费者2号:消费-4号消息
消费者2号:消费-5号消息
消费者2号:消费-6号消息
消费者2号:消费-7号消息
消费者2号:消费-8号消息
消费者2号:消费-9号消息
消费者2号:消费-10号消息
消费者2号:消费-11号消息
消费者2号:消费-12号消息
消费者2号:消费-13号消息
消费者2号:消费-14号消息
消费者2号:消费-15号消息
消费者2号:消费-16号消息
消费者2号:消费-17号消息
消费者2号:消费-18号消息
消费者2号:消费-19号消息
消费者2号:消费-20号消息
4.2.3 Режим публикации и подписки (Fanout Broadcast)
-
Producer: производитель сообщения (программа, отправившая сообщение). -
Exchange: Exchange, отвечающий за отправку сообщений в указанную очередь. -
Queue: Очередь сообщений, понимаемая как контейнер, производитель отправляет в нее сообщения, он хранит сообщения и ждет, пока потребители их потреблят. -
Consumer: потребитель сообщения (программа, получающая сообщение).
4.2.3.1 Как понять
Fanout буквально переводится как "разветвление", но больше людей будут называть его широковещательным или публикацией и подпиской. Это режим без ключей маршрутизации. Производитель отправляет сообщение коммутатору, а коммутатор копирует и синхронизирует все сообщения для всех других пользователей. Он привязан к очереди, и в каждой очереди может быть только один потребитель для получения сообщения.Если в соединении с потребителем создано несколько каналов, это приведет к конкуренции за сообщение.
4.2.3.2 Реализация кода
Примечание: ниже описаны только те части, которые отличаются от вышеперечисленных.В простом режиме эти основные методы были введены
4.2.3.2.1 Код производителя
public class Producer {
@Test
public void sendMessage() throws IOException, TimeoutException {
// 通过工具类获取连接
Connection connection = RabbitMqUtil.getConnection();
// 获取连接通道
final Channel channel = connection.createChannel();
// 声明交换机
channel.exchangeDeclare("order", "fanout");
for (int i = 1; i <= 20; i++) {
// 发布消息
channel.basicPublish("order", "", null, "fanout!".getBytes());
}
// 通过工具关闭channel和释放连接
RabbitMqUtil.close(channel, connection);
}
}
-
объявить переключатель
- Объяснение метода exchangeDeclare
- Параметр 1: Exchange (Имя коммутатора) Если коммутатор не существует, автоматически создается
- Параметр 2: type (тип), здесь выбираем режим разветвления
- Объяснение метода exchangeDeclare
-
опубликовать новость: Введите имя коммутатора, определенное выше, в первый параметр метода basicPublish, второй параметр, ключ маршрутизации пуст.
- Цикл 20 для демонстрации потребителя
4.2.3.2.2 Код потребителя
- Потребитель 1
public class Consumer1 {
public static void main(String[] args) throws IOException {
// 通过工具类获取连接
Connection connection = RabbitMqUtil.getConnection();
Channel channel = connection.createChannel();
// 声明交换机
channel.exchangeDeclare("order", "fanout");
// 创建临时队列
String queue = channel.queueDeclare().getQueue();
// 绑定临时队列和交换机
channel.queueBind(queue, "order", "");
// 消费消息
channel.basicConsume(queue, true, new DefaultConsumer(channel) {
@Override
public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
System.out.println("消费者1号:消费-" + new String(body));
}
});
}
}
- объявить переключатель
- Создать временную очередь
-
Привязка временных очередей и обменов
- Объяснение метода queueBind
- Параметр 1: очередь (временная очередь)
- Параметр 2: обмен
- Параметр 3: routingKey (ключ маршрутизации)
- Объяснение метода queueBind
- Потребитель 2: демонстрирует ситуацию с несколькими каналами в одном соединении.
public class Consumer2 {
public static void main(String[] args) throws IOException {
// 通过工具类获取连接
Connection connection = RabbitMqUtil.getConnection();
// 获取连接通道
Channel channel = connection.createChannel();
Channel channel2 = connection.createChannel();
// 声明交换机
channel.exchangeDeclare("order", "fanout");
channel2.exchangeDeclare("order", "fanout");
// 创建临时队列
String queue = channel.queueDeclare().getQueue();
System.out.println(queue);
// 绑定临时队列和交换机
channel.queueBind(queue, "order", "");
channel2.queueBind(queue, "order", "");
// 消费消息
channel.basicConsume(queue, true, new DefaultConsumer(channel) {
@Override
public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
System.out.println("消费者2号:消费-" + new String(body));
}
});
// 消费消息
channel2.basicConsume(queue, true, new DefaultConsumer(channel2) {
@Override
public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
System.out.println("消费者2-2号:消费-" + new String(body));
}
});
}
}
результат операции:
消费者2号:消费-fanout!
消费者2号:消费-fanout!
消费者2-2号:消费-fanout!
消费者2号:消费-fanout!
消费者2号:消费-fanout!
消费者2号:消费-fanout!
消费者2号:消费-fanout!
消费者2号:消费-fanout!
消费者2号:消费-fanout!
消费者2号:消费-fanout!
消费者2号:消费-fanout!
消费者2-2号:消费-fanout!
消费者2-2号:消费-fanout!
消费者2-2号:消费-fanout!
消费者2-2号:消费-fanout!
消费者2-2号:消费-fanout!
消费者2-2号:消费-fanout!
消费者2-2号:消费-fanout!
消费者2-2号:消费-fanout!
消费者2-2号:消费-fanout!
4.2.3.2.3 Почему переключатель также объявлен в потребителе?
Как видно из приведенного выше кода, мы объявили коммутаторы в Producer и Conusmer соответственно, но из рисунка видно, что потребители не будут иметь прямого контакта со коммутаторами.Почему потребители также объявляют коммутаторы?
Это делается для того, чтобы при выполнении Producer или Producer никогда не возникало ошибки, поскольку переключатель не был объявлен.Например, если вы объявляете переключатель только в Producer, вы должны сначала запустить Producer.Если Вы напрямую запускаете Conusmer, переключателя там еще не будет, если он есть, то будет сообщено об ошибке. Написав все объявления, можно гарантировать, что независимо от того, кто начнет первым, он будет объявлен коммутатору.
4.2.4 Режим маршрутизации (Маршрутизация / Прямой)
-
Producer: производитель сообщения (программа, отправившая сообщение). -
Exchange: Exchange, отвечающий за отправку сообщений в указанную очередь. -
routingKey: Ключ маршрутизации, то есть ключ1, ключ2 и т. д. на приведенном выше рисунке, что эквивалентно добавлению еще одного уровня ограничений между коммутатором и очередью. -
Queue: Очередь сообщений, понимаемая как контейнер, производитель отправляет в нее сообщения, он хранит сообщения и ждет, пока потребители их потреблят. -
Consumer: потребитель сообщения (программа, получающая сообщение).
4.2.4.1 Как понять
Тип коммутатора в режиме маршрутизации прямой.По сравнению с режимом разветвления добавлена концепция ключа маршрутизации. Производитель отправляет сообщение, содержащее указанный routingKey (ключ маршрутизации), на биржу, а биржа берет routingKey, чтобы найти очередь, привязанную к routingKey, а затем отправляет ее в очередь.Очередь может быть привязана к нескольким routingKeys.
4.2.4.2 Реализация кода
4.2.4.2.1 Код производителя
public class Producer {
@Test
public void sendMessage() throws IOException, TimeoutException {
// 通过工具类获取连接
Connection connection = RabbitMqUtil.getConnection();
// 获取连接通道
Channel channel = connection.createChannel();
// 声明交换机
channel.exchangeDeclare("order_direct", "direct");
// 指定 routingKey
String key = "info";
// 发布消息
channel.basicPublish("order_direct", key, null, ("发送给指定路由" + key + "的消息").getBytes());
// 通过工具关闭channel和释放连接
RabbitMqUtil.close(channel, connection);
}
}
- Указать routingKey, то есть во втором параметре метода basicPublish указать значение ключа
4.2.4.2.2 Код потребителя
- Потребитель 1
public class Consumer1 {
public static void main(String[] args) throws IOException {
// 通过工具类获取连接
Connection connection = RabbitMqUtil.getConnection();
Channel channel = connection.createChannel();
// 声明交换机
channel.exchangeDeclare("order_direct", "direct");
// 获取临时队列
String queue = channel.queueDeclare().getQueue();
// 绑定临时队列和交换机
channel.queueBind(queue, "order_direct", "info");
channel.queueBind(queue, "order_direct", "error");
channel.queueBind(queue, "order_direct", "warn");
// 消费消息
channel.basicConsume(queue, true, new DefaultConsumer(channel) {
@Override
public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
System.out.println("消费者1:消费-" + new String(body));
}
});
}
}
- Просто добавьте значение ключа при привязке очереди и обмена
- Потребитель 2
public class Consumer2 {
public static void main(String[] args) throws IOException {
// 通过工具类获取连接
Connection connection = RabbitMqUtil.getConnection();
Channel channel = connection.createChannel();
// 声明交换机
channel.exchangeDeclare("order_direct", "direct");
// 获取临时队列
String queue = channel.queueDeclare().getQueue();
// 绑定临时队列和交换机
channel.queueBind(queue, "order_direct", "error");
// 消费消息
channel.basicConsume(queue, true, new DefaultConsumer(channel) {
@Override
public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
System.out.println("消费者2:消费-" + new String(body));
}
});
}
}
Результат выполнения: только потребитель 1 получил сообщение
消费者1:消费-发送给指定路由info的消息
4.2.5 Режим сопоставления подстановки (тема)
-
Producer: производитель сообщения (программа, отправившая сообщение). -
Exchange: Exchange, отвечающий за отправку сообщений в указанную очередь. -
routingKey: Ключ маршрутизации, то есть ключ1, ключ2 и т. д. на приведенном выше рисунке, что эквивалентно добавлению еще одного уровня ограничений между коммутатором и очередью.- Но ключ в теме находится в виде подстановочного знака, который может значительно повысить эффективность
-
Queue: Очередь сообщений, понимаемая как контейнер, производитель отправляет в нее сообщения, он хранит сообщения и ждет, пока потребители их потреблят. -
Consumer: потребитель сообщения (программа, получающая сообщение).
4.2.5.1 Как понять
Типом переключателя режима сопоставления подстановочных знаков является тема, поскольку он очень похож на прямой режим, поэтому иногда люди также объединяют прямой режим и тему в режиме маршрутизации.Разница между ними заключается в том, что routingKey режима прямого является заданным значением. routingKey режима Topic может использовать подстановочные знаки и обычно состоит из одного или нескольких слов, а несколько слов разделяются знаком «.», например: Ideal.insert.
-
*: соответствует ровно одному слову, например:order.*Может соответствовать заказу. вставка -
#: соответствует одному или нескольким словам, например:order.#Может совпадать с order.insert.common-
#как многослойное понятие, в то время как*просто однослойная концепция
-
4.2.5.2 Реализация кода
4.2.5.2.1 Код производителя
public class Producer {
@Test
public void sendMessage() throws IOException, TimeoutException {
// 通过工具类获取连接
Connection connection = RabbitMqUtil.getConnection();
// 获取连接通道
Channel channel = connection.createChannel();
channel.exchangeDeclare("order_topic", "topic");
// 声明交换机
String key = "user.query.all";
// 发布消息
channel.basicPublish("order_topic", key, null, ("发送给指定路由" + key + "的消息").getBytes());
RabbitMqUtil.close(channel, connection);
}
}
4.2.5.2.2 Код потребителя
- Потребитель 1
public class Consumer1 {
public static void main(String[] args) throws IOException {
// 通过工具类获取连接
Connection connection = RabbitMqUtil.getConnection();
// 获取连接通道
Channel channel = connection.createChannel();
// 声明交换机
channel.exchangeDeclare("order_topic", "topic");
// 获取临时队列
String queue = channel.queueDeclare().getQueue();
// 指定路由key
String key = "user.*";
channel.queueBind(queue, "order_topic", key);
// 发布消息
channel.basicConsume(queue, true, new DefaultConsumer(channel) {
@Override
public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
System.out.println("消费者1:消费-" + new String(body));
}
});
}
}
- Потребитель 2
public class Consumer2 {
public static void main(String[] args) throws IOException {
// 通过工具类获取连接
Connection connection = RabbitMqUtil.getConnection();
// 获取连接通道
Channel channel = connection.createChannel();
// 声明交换机
channel.exchangeDeclare("order_topic", "topic");
// 获取临时队列
String queue = channel.queueDeclare().getQueue();
// 指定路由key
String key = "user.#";
channel.queueBind(queue, "order_topic", key);
channel.basicConsume(queue, true, new DefaultConsumer(channel) {
@Override
public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
System.out.println("消费者2:消费-" + new String(body));
}
});
}
}
Результат выполнения: только потребитель 2 получил сообщение, поскольку сообщение представляет собой многоуровневую структуру, толькоuser.#Может соответствовать
消费者2:消费-发送给指定路由user.query.all的消息
5. Внедрение Springboot Rabbitmq
SpringBoot предоставляет стартовый пакет Spring For RabbitMQ, а также предоставляет серию аннотаций и шаблон RabbitTemplate, который может значительно упростить этапы разработки RabbitMQ. Следующее демонстрирует [5.1 На основе чистой аннотации] и [5.2 На основе аннотации + Класс конфигурации] Способ написания аналогичен, но позиция объявления и привязки обмена очередями и т. д. отличается. Обычно считается, что последний лучше подходит для обслуживания и управления, поэтому вы можете выбрать один из них.
Подготовка среды:
- Сначала создайте проект SprinBoot, затем выберите стартер RabbitMQ и базовые стартеры, такие как модульные тесты.
- Напишите файл конфигурации yml и запишите данные, необходимые для подключения к RabbitMQ.
Зависимости RabbitMQ
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-amqp</artifactId>
</dependency>
yml-файл конфигурации
spring:
rabbitmq:
host: 192.168.122.1 # 服务器地址
port: 5672 # tcp端口
username: admin # 用户名
password: admin # 用户密码
virtual-host: /rabbitmq_springboot_01 # 虚拟主机
5.1 На основе чистых аннотаций
Примечание. Этот метод не создает класс конфигурации для управления объявлением и привязкой очередей и переключателей, но все они записываются непосредственно в потребителях с помощью аннотаций.
5.1.1 Режим простой очереди
Весь код для создания сообщений, мы помещаем его в тест, чтобы сделать это
- режиссер
@SpringBootTest(classes = RabbitmqSpringbootApplication.class)
@RunWith(SpringRunner.class)
public class RabbitMqTest {
/**
* 注入 RabbitTemplate
*/
@Autowired
private RabbitTemplate rabbitTemplate;
@Test
public void testSimpleSendMessage() {
rabbitTemplate.convertAndSend("simple_queue", "This is a message !");
}
}
- Первый шаг — внедрить RabbitTemplate, предоставленный нам SpringBoot.
- Используется для отправки сообщений через конвертируемый метод раббитомата, у него есть множество способов перегруженных соответственно, будет использовать два и три аргумента сегодня
- Подробное объяснение метода convertAndSend (два параметра)
- Параметр 1: routingKey (ключ маршрутизации)
- Параметр 2: объект (тело отправленного сообщения)
- Подробное объяснение метода convertAndSend (три параметра)
- Параметр 1: обмен
- Параметр 2: routingKey (ключ маршрутизации)
- Параметр 3: объект (тело отправленного сообщения)
- Подробное объяснение метода convertAndSend (два параметра)
- потребитель
// 注入容器
@Component
// 监听 RabbitMQ
@RabbitListener(queuesToDeclare = @Queue(value = "simple_queue", durable = "true", exclusive = "false", autoDelete = "false"))
public class SimpleConsumer {
// 自动回调
@RabbitHandler
public void receiveMessage(String message) {
System.out.println("消费者:" + message);
}
}
-
контейнер для инъекций
-
Слушайте RabbitMQ, в аннотации @RabbitListener можно реализовать и объявление очереди, и привязку между коммутатором и очередью и т.д.
- @Queue может иметь четыре параметра, потому что каждый имеет значение по умолчанию, поэтому только при заданном значении значения он будет создан по умолчанию в виде сохраняемости, неисключительного, неавтоматического удаления.
- Параметр 1: значение (имя очереди)
- Параметр 2: прочный (постоянная очередь сообщений). После перезапуска RabbitMQ очередь все еще существует, по умолчанию установлено значение true.
- Параметр 3: Эксклюзив (исключая) указывает, что очередь сообщений вступает в силу только в текущем соединении, по умолчанию является ложным
- Параметр 4: auto-delete (автоматическое удаление) указывает, что очередь сообщений будет автоматически удаляться, когда она не используется, по умолчанию false
- @Queue может иметь четыре параметра, потому что каждый имеет значение по умолчанию, поэтому только при заданном значении значения он будет создан по умолчанию в виде сохраняемости, неисключительного, неавтоматического удаления.
-
Добавьте аннотацию @RabbitHandler к методу для реализации автоматических обратных вызовов, чтобы мы могли получить сообщение от производителя.
- Примечание. Тип параметра метода ReceiveMessage зависит от типа данных, которые вы отправили от производителя.
5.1.2 Режим рабочей очереди
5.1.2.1 Режим опроса
- Производитель: Нечего сказать, потому что режим работы имеет несколько потребителей, поэтому отправьте еще несколько сообщений.
@SpringBootTest(classes = RabbitmqSpringbootApplication.class)
@RunWith(SpringRunner.class)
public class RabbitMqTest {
/**
* 注入 RabbitTemplate
*/
@Autowired
@Test
public void testWorkSendMessage() {
for (int i = 0; i < 20; i++) {
rabbitTemplate.convertAndSend("work_queue", "This is a message !, 序号:" + i);
}
}
}
- потребитель
@Component
public class WorkConsumer {
// 监听 RabbitMQ
@RabbitListener(queuesToDeclare = @Queue("work_queue"))
// 消费者1
public void receiveMessage1(String message) {
System.out.println("消费者1:" + message);
// 监听 RabbitMQ
@RabbitListener(queuesToDeclare = @Queue("work_queue")
// 消费者2
public void receiveMessage2(String message) {
System.out.println("消费者2:" + message);
}
}
- Аннотация @RabbitListener может быть размещена либо в классе, либо в методе.Например, в приведенном выше коде мы помещаем ее в два метода для ссылки на разных потребителей.
- Однако, если вы добавите в класс аннотацию @RabbitListener, а в следующих двух методах добавление аннотации @RabbitHandler сообщит об ошибке, вам необходимо создать класс для каждого потребителя отдельно
5.1.2.2 Справедливая модель (распределение по мощности)
5.1.2.2.1 Как изменить файл конфигурации
-
производитель не изменился
-
Измените файл конфигурации yml/properties
spring:
rabbitmq:
host: 192.168.122.1 # 服务器地址
port: 5672 # tcp端口
username: admin # 用户名
password: admin # 用户密码
virtual-host: /rabbitmq_springboot_01 # 虚拟主机
# 新增部分
listener:
simple:
acknowledge-mode: manual # 开启 ack 手动应答
prefetch: 1 # 每次只能消费 1 条消息
- Введение в опцию режима подтверждения
- auto: автоматическое подтверждение, опция по умолчанию
- manual: ручное подтверждение (назначение по способностям должно быть установлено на ручное подтверждение)
- none: без подтверждения, автоматически отбрасывается после отправки
- потребитель
@Component
public class WorkConsumer {
// 监听 RabbitMQ
@RabbitListener(queuesToDeclare = @Queue("work_queue"))
// 消费者 1
public void receiveMessage(String body, Message message, Channel channel) throws IOException {
try {
// 打印输出消息主题
System.out.println("消费者1:" + body);
// 返回 deliveryTag 代表队列可以删除此消息了
channel.basicAck(message.getMessageProperties().getDeliveryTag(),false);
} catch (IOException e) {
e.printStackTrace();
// 消费者告诉队列信息消费失败
channel.basicNack(message.getMessageProperties().getDeliveryTag(), false, true);
}
}
// 监听 RabbitMQ
@RabbitListener(queuesToDeclare = @Queue("work_queue"))
// 消费者 2
public void receiveMessage2(String body, Message message, Channel channel) throws IOException{
try {
// 延迟 2s 代表处理业务慢
Thread.sleep(2000);
} catch (InterruptedException e) {
e.printStackTrace();
}
try {
// 打印输出消息主题
System.out.println("消费者2:" + body);
// 返回 deliveryTag 代表队列可以删除此消息了
channel.basicAck(message.getMessageProperties().getDeliveryTag(), false);
} catch (IOException e) {
e.printStackTrace();
// 消费者告诉队列信息消费失败
channel.basicNack(message.getMessageProperties().getDeliveryTag(), false, true);
}
}
}
-
Поскольку ручное подтверждение включено в конфигурации yml, необходимо возвращать подтверждающие сообщения после успеха и неудачи соответственно.
-
Объяснение метода basicAck
- Параметр 1: deliveryTag (тег доставки, то есть индекс сообщения), return означает, что сообщение получено, и очередь может удалить сообщение
- Параметр 2: множественный (независимо от того, пакетный или нет) выберите true, чтобы отклонить все сообщения меньше, чем deliveryTag за один раз.
-
Объяснение основного метода Nack
- Параметр 1 | Параметр 2 То же, что и выше
- Параметр 3: requeue (будут ли отклоненные повторно входить в очередь)
результат операции:
消费者1:This is a message !, 序号:2
消费者1:This is a message !, 序号:3
消费者1:This is a message !, 序号:4
消费者1:This is a message !, 序号:5
消费者1:This is a message !, 序号:6
消费者1:This is a message !, 序号:7
消费者1:This is a message !, 序号:8
消费者1:This is a message !, 序号:9
消费者1:This is a message !, 序号:10
消费者1:This is a message !, 序号:11
消费者1:This is a message !, 序号:12
消费者1:This is a message !, 序号:13
消费者1:This is a message !, 序号:14
消费者1:This is a message !, 序号:15
消费者1:This is a message !, 序号:16
消费者1:This is a message !, 序号:17
消费者1:This is a message !, 序号:18
消费者1:This is a message !, 序号:19
消费者1:This is a message !, 序号:20
消费者2:This is a message !, 序号:1
До сих пор был реализован метод изменения файла конфигурации для реализации распределения в соответствии с возможностями, и было добавлено несколько конфигураций.Мы использовали только часть вышеперечисленного, а остальные удобны для вашего ознакомления. Вы можете выбрать yml и свойства самостоятельно.
# 发送确认
spring.rabbitmq.publisher-confirm-type=correlated
# spring.rabbitmq.publisher-confirms=true(旧版)
# 发送回调
spring.rabbitmq.publisher-returns=true
# 消费手动确认
spring.rabbitmq.listener.direct.acknowledge-mode=manual
spring.rabbitmq.listener.simple.acknowledge-mode=manual
# 并发消费者初始化值
spring.rabbitmq.listener.simple.concurrency=1
# 并发消费者的最大值
spring.rabbitmq.listener.simple.max-concurrency=10
# 每个消费者每次监听时可拉取处理的消息数量
# 在单个请求中处理的消息个数,他应该大于等于事务数量(unack的最大数量)
spring.rabbitmq.listener.simple.prefetch=1
# 是否支持重试
spring.rabbitmq.listener.simple.retry.enabled=true
5.1.2.2.1 Как настроить фабрику
/**
* 设置消费者的确认机制,并达到能者多劳的效果
*
* @param connectionFactory 连接工厂
* @return
*/
@Bean("workListenerFactory")
public RabbitListenerContainerFactory myFactory(ConnectionFactory connectionFactory) {
SimpleRabbitListenerContainerFactory containerFactory =
new SimpleRabbitListenerContainerFactory();
containerFactory.setConnectionFactory(connectionFactory);
// 修改为手动确认
containerFactory.setAcknowledgeMode(AcknowledgeMode.MANUAL);
// 拒绝策略,true 回到队列 false丢弃,默认是true
containerFactory.setDefaultRequeueRejected(true);
// 默认的PrefetchCount是250 修改为 1
containerFactory.setPrefetchCount(1);
return containerFactory;
}
- Потребительская модификация
@RabbitListener(queuesToDeclare = @Queue("work_queue"))
// 将上面的监听,增加 containerFactory 属性,然后将配置好的工厂传入
@RabbitListener(queuesToDeclare = @Queue("work_queue"), containerFactory = "workListenerFactory")
5.1.3 Режим публикации и подписки
- режиссер
@SpringBootTest(classes = RabbitmqSpringbootApplication.class)
@RunWith(SpringRunner.class)
public class RabbitMqTest {
/**
* 注入 RabbitTemplate
*/
@Autowired
@Test
public void testFanoutSendMessage() {
rabbitTemplate.convertAndSend("order_exchange", "", "This is a message !");
}
}
- Т.к. из этого режима задействован переключатель, поэтому используется трехпараметрический метод.
- потребитель
@Component
public class FanoutConsumer {
// 绑定临时队列和交换机
@RabbitListener(bindings = {
@QueueBinding(
value = @Queue(), // 临时队列
exchange = @Exchange(name = "order_exchange", type = "fanout") // 交换机与类型
)
})
public void receiveMessage1(String message) {
System.out.println("消费者1:" + message);
}
// 绑定临时队列和交换机
@RabbitListener(bindings = {
@QueueBinding(
value = @Queue(), // 临时队列
exchange = @Exchange(name = "order_exchange", type = "fanout") // 交换机与类型
)
})
public void receiveMessage2(String message) {
System.out.println("消费者2:" + message);
}
}
5.1.4 Режим маршрутизации (прямой)
- режиссер
@SpringBootTest(classes = RabbitmqSpringbootApplication.class)
@RunWith(SpringRunner.class)
public class RabbitMqTest {
/**
* 注入 RabbitTemplate
*/
@Autowired
@Test
public void testDirectSendMessage() {
rabbitTemplate.convertAndSend("direct_exchange", "info", "This is a message !");
}
}
- потребитель
@Component
public class DirectConsumer {
// 绑定临时队列和交换机
@RabbitListener(bindings = {
@QueueBinding(
value = @Queue(), // 临时队列
exchange = @Exchange(name = "direct_exchange", type = "direct"), // 交换机和类型
key = {"info", "warn", "error"} // 路由key
)
})
public void receiveMessage1(String message) {
System.out.println("消费者1:" + message);
}
// 绑定临时队列和交换机
@RabbitListener(bindings = {
@QueueBinding(
value = @Queue(), // 临时队列
exchange = @Exchange(name = "direct_exchange", type = "direct"), // 交换机和类型
key = {"info", "warn", "error"} // 路由key
)
})
public void receiveMessage2(String message) {
System.out.println("消费者2:" + message);
}
}
5.1.5 Тематический режим
- режиссер
@SpringBootTest(classes = RabbitmqSpringbootApplication.class)
@RunWith(SpringRunner.class)
public class RabbitMqTest {
/**
* 注入 RabbitTemplate
*/
@Autowired
@Test
public void testTopicSendMessage() {
rabbitTemplate.convertAndSend("topic_exchange", "order.insert.common", "This is a message !");
}
}
- потребитель
@Component
public class TopicConsumer {
// 绑定临时队列和交换机
@RabbitListener(bindings = {
@QueueBinding(
value = @Queue(), // 临时队列
exchange = @Exchange(name = "topic_exchange", type = "topic"), // 交换机和类型
key = {"order.*"} // 通配符路由key
)
})
public void receiveMessage1(String message) {
System.out.println("消费者1:" + message);
}
// 绑定临时队列和交换机
@RabbitListener(bindings = {
@QueueBinding(
value = @Queue(), // 临时队列
exchange = @Exchange(name = "topic_exchange", type = "topic"), // 交换机和类型
key = {"order.*"} // 通配符路由key
)
})
public void receiveMessage2(String message) {
System.out.println("消费者2:" + message);
}
}
5.2 На основе аннотаций + класс конфигурации
Фактически таким образом осуществляется объявление и привязка коммутаторов и очередей в классе конфигурации.Одно - упрощение аннотаций в потребителях, а другое - унифицированное управление, которое более организовано, и производитель и потребитель ссылка Так же удобнее в момент модификации, и при модификации в дальнейшем не нужно будет модифицировать каждое место.
Из-за давности здесь самый сложный метод Топик, да и другие тоже верующие.
- класс конфигурации
@Configuration
public class RabbitMqConfiguration {
public static final String TOPIC_EXCHANGE = "topic_order_exchange";
public static final String TOPIC_QUEUE_NAME_1 = "test_topic_queue_1";
public static final String TOPIC_QUEUE_NAME_2 = "test_topic_queue_2";
public static final String TOPIC_ROUTINGKEY_1 = "test.*";
public static final String TOPIC_ROUTINGKEY_2 = "test.#";
@Bean
public TopicExchange topicExchange() {
return new TopicExchange(TOPIC_EXCHANGE);
}
@Bean
public Queue topicQueue1() {
return new Queue(TOPIC_QUEUE_NAME_1);
}
@Bean
public Queue topicQueue2() {
return new Queue(TOPIC_QUEUE_NAME_2);
}
@Bean
public Binding bindingTopic1(){
return BindingBuilder.bind(topicQueue1())
.to(topicExchange())
.with(TOPIC_ROUTINGKEY_1);
}
@Bean
public Binding bindingTopic2(){
return BindingBuilder.bind(topicQueue2())
.to(topicExchange())
.with(TOPIC_ROUTINGKEY_2);
}
}
-
Добавьте аннотацию @Configuration: указывает, что это класс конфигурации
-
определить константы: Имя коммутатора, имя очереди, ключ маршрутизации и т. д. могут быть созданы как константы, которые очень удобно вызывать, управлять и изменять.Также можно создать специальный класс констант RabbitMQ.
-
определить переключатель: Выберите тип TopicExchange, потому что этот пример — Topic
-
определить очередь: вы можете передать константу имени очереди, потому что есть значения по умолчанию для сохраняемости и т. д., вы также можете настроить сохраняемость, будь то эксклюзивная и другие параметры
-
Связывание обменов и очередей: Используйте метод привязки BindingBuilder для привязки очереди к указанному коммутатору с входящим ключом маршрутизации.
- режиссер
@SpringBootTest(classes = RabbitmqSpringbootApplication.class)
@RunWith(SpringRunner.class)
public class RabbitMqTest {
/**
* 注入 RabbitTemplate
*/
@Autowired
@Test
public void testTopicSendMessage() {
rabbitTemplate.convertAndSend(RabbitMqConfiguration.TOPIC_EXCHANGE, "test.order.insert", "This is a message !");
}
}
- потребитель
@Component
public class TopicConsumer {
// 绑定队列即可
@RabbitListener(queues = {RabbitMqConfiguration.TOPIC_QUEUE_NAME_1})
public void receiveMessage1(String message) {
System.out.println("消费者1:" + message);
}
// 绑定队列即可
@RabbitListener(queues = {RabbitMqConfiguration.TOPIC_QUEUE_NAME_2})
public void receiveMessage2(String message) {
System.out.println("消费者2:" + message);
}
}