Основы обучения очереди сообщений

задняя часть

Что такое МАМА

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

Идея MOM заключается в том, что два приложения A и B не отправляют сообщения напрямую. До того, как A и B отправили сообщения напрямую, было много проблем с эффективностью. Например, если B не принял сообщение вовремя после того, как его отправил A, тогда A продолжал бы блокировать там, и параллелизм был плохим. A должен ждать, пока B примет сообщение. получить результат, и тогда A может закончиться. И MOM должен решить такую ​​проблему.Он не позволяет взаимодействовать между A и B.Между A и B добавляется промежуточное программное обеспечение сообщения.A помещает сообщение в середину сообщения и может уйти и заняться другими делами.Когда B приходит к промежуточному программному обеспечению сообщений, чтобы получить сообщение, которое A не нужно знать или заботиться. Это повышает эффективность и обеспечивает параллелизм.После того, как B уходит, он может уведомить A через статус, уведомление, обратный вызов и т. д. На рынке существует множество технологий, реализующих эту идею, например IBM (MQSEVICES), Microsoft (MSMQ) и MessageMQ компании BEA. На этапе соперничества сотен школ каждая добивается своего, и единого стандарта реализации не существует. В это время для реализации единого стандарта появилась унифицированная спецификация реализации JMS. JMS в основном имеет две модели сообщений: «точка-точка» и «публикация-подписка».

Что такое очередь сообщений

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

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

В настоящее время в производственной среде наиболее часто используемыми очередями сообщений являются ActiveMQ, RabbitMQ, ZeroMQ, Kafka, MetaMQ, RocketMQ и т. д.

Преимущества очереди сообщений

  1. Передача данных из одного приложения в другое или из одного программного модуля в другой
  2. Отвечает за установление сетевых каналов связи для надежной передачи данных.
  3. Убедитесь, что данные не передаются повторно или не теряются
  4. Он может реализовать кроссплатформенную работу и может интегрировать службы передачи данных для программного обеспечения в разных операционных системах.

Сценарии применения очереди сообщений

Далее подробно описаны распространенные сценарии использования очередей сообщений в практических приложениях. Сцена делится наАсинхронная обработка, разделение приложений, отключение трафика и обмен сообщениямиЧетыре сцены.

Асинхронная обработка

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

последовательный режим

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


Параллельный путь

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

Из приведенного выше сравнения видно, что если предположить, что все три операции требуют времени выполнения 50 мс, исключая сетевые факторы, окончательное выполнение будет завершено.Последовательный режим занимает 150 мс, а параллельный режим - 100 мс.

Поскольку количество запросов, обрабатываемых ЦП в единицу времени, одинаково, при условии, что пропускная способность ЦП в секунду равна 100, количество запросов, которые могут быть выполнены за одну секунду в последовательном режиме, составляет 1000/150, меньше, чем 7 раз, в параллельном режиме количество запросов, которые могут быть выполнены за 1 секунду, составляет 1000/100, то есть 10 раз.

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


В соответствии с описанным выше процессом время ответа пользователя в основном эквивалентно времени записи данных в базу данных, отправке регистрационного электронного письма и отправке зарегистрированного SMS-сообщения после записи в очередь сообщений, и результат выполнения может быть возвращен. Время записи в очередь сообщений очень короткое, быстрое, почти незначительное, а также может увеличить пропускную способность системы до 20QPS, что почти в 3 раза выше, чем у последовательного метода и в 2 раза выше, чем у параллельного метода.

Разделение приложений

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

Традиционный подход таков: система заказов вызывает интерфейс системы инвентаризации. Как показано ниже:


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

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


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

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

Отключение трафика

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

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

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

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


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

обработка журнала

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

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


Архитектура очереди сообщений применительно к обработке журналов
  • Клиент сбора журналов: отвечает за сбор данных журнала и регулярную запись в очередь Kafka;
  • очередь сообщений Kafka: отвечает за получение, хранение и пересылку данных журнала;
  • Приложение для обработки журналов: подпишитесь и используйте данные журнала в очереди kafka;

Применение этой архитектуры в реальной разработке может относиться к случаю:Sina Technology Sharing: как мы проводим анализ и обработку 3,2 миллиарда журналов в режиме реального времени

Дизайн технической архитектуры сервиса

  1. Kafka: очередь сообщений, которая получает пользовательские журналы.
  2. Logstash: анализируйте журнал, объединяйте его в JSON и выводите в Elasticsearch.
  3. Elasticsearch: основная технология службы анализа журналов в реальном времени, служба хранения данных без схемы в реальном времени, организует данные с помощью индекса и имеет мощные функции поиска и статистики.
  4. Kibana: Компонент визуализации данных, основанный на Elasticsearch, возможность супервизуализации данных является важной причиной, по которой многие компании выбирают стек ELK.

Новостная рассылка

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

одноранговое общение

Архитектура одноранговой связи

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

общение в чате

Архитектура общения в чате

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

Служба сообщений JMS

Говоря об очередях сообщений, нельзя не упомянуть JMS. JMS (служба сообщений Java, служба сообщений Java) JMS называется службой сообщений Java (служба сообщений Java) и представляет собой техническую спецификацию для MOM на платформе Java. обработка сообщений.Разработка приложений, аналогичная абстракции JDBC и взаимодействию с реляционными базами данных.

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

В архитектуре EJB есть bean-компоненты сообщений, которые можно легко интегрировать со службой сообщений JM. В шаблоне архитектуры J2EE есть шаблон сервера сообщений, который используется для реализации прямого разделения сообщений и приложений.


Общие понятия

  • Поставщик: реализация интерфейса JMS, написанная на чистом языке Java (например, ActiveMQ).
  • Домены: методы обмена сообщениями, включая одноранговые (P2P) и публикации/подписки (Pub/Sub).
  • Фабрика соединений: клиент использует фабрику соединений для создания соединения с провайдером JMS.
  • Назначение: объект, которому адресовано, отправлено и получено сообщение.

модель сообщения

В стандарте JMS есть две модели сообщений: P2P (Point-to-Point), Publish/Subscribe (Pub/Sub)

P2P-режим


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

Домен сообщений P2P использует очередь в качестве пункта назначения.Сообщения могут быть отправлены и получены синхронно или асинхронно, и каждое сообщение будет отправлено только одному Потребителю один раз. Потребители могут получать сообщения синхронно, используя MessageConsumer.receive(), или асинхронно, регистрируя MessageListener, используя MessageConsumer.setMessageListener().

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


Особенности P2P

  1. У каждого сообщения есть только один потребитель (Consumer) (то есть после потребления сообщение больше не находится в очереди сообщений, и другие потребители не могут получить это сообщение).
  2. Проверки качества отправителя и получателя не зависят от времени, то есть когда отправитель отправляет сообщение, независимо от того, запущен ли получатель, это не повлияет на сообщение, отправляемое в очередь.
  3. Потребитель должен подтвердить получение сообщения

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

  4. Непостоянные сообщения отправляются не более одного раза

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

    1. Поставщик услуг JMS не работает, что приводит к потере непостоянной информации.

    2. Сообщение в очереди просрочено и не получено

  5. Постоянные сообщения отправляются строго один раз

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

Режим p2p необходим, если вы хотите, чтобы каждое отправленное сообщение обрабатывалось успешно

Режим паб/саб



Содержит три роли: топик (Topic), издатель (Publisher), подписчик (Subscriber). Несколько издателей отправляют сообщения в темы, и система доставляет эти сообщения нескольким подписчикам.

Домен сообщений Pub/Sub (публикация/подписка, публикация/подписка) использует тему в качестве адресата, издатели отправляют сообщения в тему, а подписчики регистрируются для получения сообщений из темы. Любое сообщение, отправленное в тему, будет автоматически доставлено всем подписчикам. Способы получения (синхронный и асинхронный) такие же, как и в домене P2P.

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

Особенности Pub/Sub

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

  • Между издателями и подписчиками существует временная зависимость. Подписчики могут потреблять сообщения, опубликованные только после того, как они подпишутся.Для подписчика темы (Topic) он должен создать подписчика перед использованием сообщений от издателя. Это требует, чтобы подписчики сначала подписывались, а производители публиковали. То есть сначала должен запуститься подписчик, а затем дождаться запуска производителя, который отличается от типа точка-точка.
  • Чтобы получать сообщения, подписчики должны оставаться в рабочем состоянии. То есть подписчик должен оставаться активным в ожидании сообщения, опубликованного издателем, и если подписчик запускается после публикации сообщения издателем, он не может получить сообщение, опубликованное предыдущим издателем.

Чтобы смягчить такие строгие временные зависимости, JMS позволяет подписчикам создавать долговременную подписку. Таким образом, даже если подписчик не активирован (работает), он может получать сообщения от издателя.
Если сообщение, которое вы хотите отправить, может быть обработано без какой-либо обработки, или может быть обработано только одним отправителем сообщений, или может быть обработано несколькими потребителями, тогда можно использовать модель Pub/Sub.

потребление сообщений

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

  1. Синхронизировать
    Подписчики или получатели принимают сообщения через метод получения, и получение будет заблокировано до тех пор, пока сообщение не будет получено (или до истечения времени ожидания).
  2. асинхронный
    Подписчики или получатели также могут зарегистрировать прослушиватель сообщений. Когда сообщение приходит, система автоматически вызывает метод слушателя onMessage.
JDNI: Интерфейс именования и каталогов Java — это стандартный интерфейс системы именования Java. Услуги можно найти и получить к ним доступ в Интернете. При указании имени ресурса это имя соответствует записи в базе данных или службе имен и возвращает информацию, необходимую для установления соединения с ресурсом.

JNDI играет роль в поиске второго адресата доступа или источника сообщений в JMS.

Программирование JMS

Общие шаги JMS

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

Модель программирования JMS

1.ConnectionFactory

Фабрики, которые создают объекты Connection, — это QueueConnectionFactory и TopicConnectionFactory для двух разных моделей сообщений JMS. Объект ConnectionFactory можно найти через JNDI.

2.Destination

Пункт назначения означает пункт назначения отправки сообщения производителем сообщения или источник сообщения потребителя сообщения. Для производителей сообщений. Его Destination — это очередь (queue) или топик (Topic); для потребителей сообщений его Destination также является очередью или топиком (то есть источником сообщения).

Таким образом, Destination на самом деле представляет собой два типа объектов: Queue, Topic может найти Destination через JNDI.

3.Connection

Соединение представляет собой соединение, установленное между клиентом и системой JMS (обертывание сокета TCP/IP). Соединение может генерировать один или несколько сеансов. Как и ConnectionFactory, Connection также имеет два типа: QueueConnection и TopicConnection.

4.Session

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

5. Производитель сообщения

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

6. Потребители сообщений

Потребители сообщений создаются сеансом для получения сообщений, отправленных в пункт назначения. Два типа: QueueReceiver и TopicSubscriber. Его можно создать с помощью createReceiver(Queue) или createSubscriber(Topic) сеанса соответственно. Конечно, метод сеанса creatDurableSubscriber также можно использовать для создания постоянных подписчиков.

7. MessageListener

прослушиватель сообщений. Если прослушиватель сообщений зарегистрирован, после прибытия сообщения метод прослушивателя onMessage будет вызываться автоматически. MDB (Message-Driven Bean) в EJB является разновидностью MessageListener.

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

Практика JMS-программирования

возьми это здесьПример ActiveMQ

public class JMSDemo {
        ConnectionFactory connectionFactory;
        Connection connection;
        Session session;
        Destination destination;
        MessageProducer producer;
        MessageConsumer consumer;
        Message message;
        boolean useTransaction = false;
        try {
                Context ctx = new InitialContext();
                connectionFactory = (ConnectionFactory) ctx.lookup("ConnectionFactoryName");
                //使用ActiveMQ时:connectionFactory = new ActiveMQConnectionFactory(user, password, getOptimizeBrokerUrl(broker));
                connection = connectionFactory.createConnection();
                connection.start();
                session = connection.createSession(useTransaction, Session.AUTO_ACKNOWLEDGE);
                destination = session.createQueue("TEST.QUEUE");
                //生产者发送消息
                producer = session.createProducer(destination);
                message = session.createTextMessage("this is a test");

                //消费者同步接收
                consumer = session.createConsumer(destination);
                message = (TextMessage) consumer.receive(1000);
                System.out.println("Received message: " + message);
                //消费者异步接收
                consumer.setMessageListener(new MessageListener() {
                        @Override
                        public void onMessage(Message message) {
                                if (message != null) {
                                        doMessageEvent(message);
                                }
                        }
                });
        } catch (JMSException e) {
                ...
        } finally {
                producer.close();
                session.close();
                connection.close();
        }
}