Поддержка Spring Cloud Azure для Spring Integration

Расширение Spring Integration для Azure предоставляет адаптеры Spring Integration для различных служб, предоставляемых пакетом SDK Azure для Java. Мы предоставляем поддержку Spring Integration для этих служб Azure: Центры событий, служебная шина, очередь хранилища. Ниже приведен список поддерживаемых адаптеров:

Интеграция Spring с Центрами событий Azure

Основные понятия

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

Spring Integration обеспечивает упрощенное обмен сообщениями в приложениях Spring и поддерживает интеграцию с внешними системами с помощью декларативных адаптеров. Эти адаптеры обеспечивают более высокий уровень абстракции по сравнению с поддержкой Spring для удаленного взаимодействия, обмена сообщениями и планирования. Проект расширения Spring Integration for Event Hubs предоставляет адаптеры и шлюзы для центров событий Azure для входящих и исходящих каналов.

Заметка

API-интерфейсы поддержки RxJava удаляются из версии 4.0.0. Дополнительные сведения см. в Javadoc.

Группа потребителей

Центры событий обеспечивают аналогичную поддержку группы потребителей, как Apache Kafka, но с небольшой другой логикой. В то время как Kafka сохраняет все зафиксированные смещения в брокере, вам нужно вручную хранить смещения обрабатываемых сообщений Event Hubs. Пакет SDK для Центров событий позволяет хранить такие смещения в служба хранилища Azure.

Поддержка секционирования

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

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

Поддержка пакетного потребителя

EventHubsInboundChannelAdapter поддерживает режим пакетного потребления. Чтобы включить его, пользователи могут указать режим прослушивателя как ListenerMode.BATCH при создании экземпляра EventHubsInboundChannelAdapter. Если этот параметр включен, будет получено сообщение, полезной нагрузкой которого является список пакетно обработанных событий, и передано в нисходящий канал. Каждый заголовок сообщения также преобразуется в список, из которого содержимое является соответствующим значением заголовка, проанализированным из каждого события. Для общих заголовков идентификатора секции, контрольных точек и последних вложенных свойств они отображаются в виде одного значения для всего пакета событий, совместного с одним и тем же. Дополнительные сведения см. в разделе Заголовки сообщений Центров событий.

Заметка

Заголовок контрольной точки присутствует только при использовании режима контрольных точек MANUAL.

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

Политику использования пакетной службы можно указать свойствами max-size и max-wait-time, где max-size является необходимым свойством, а max-wait-time является необязательным. Чтобы указать стратегию пакетной обработки, разработчики могут использовать EventHubsContainerProperties в конфигурации. См. следующий раздел, где приведён пример настройки .

Настройка зависимостей

<dependency>
    <groupId>com.azure.spring</groupId>
    <artifactId>spring-cloud-azure-starter-integration-eventhubs</artifactId>
</dependency>

Конфигурация

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

Свойства конфигурации подключения

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

Заметка

Если вы решили использовать субъект безопасности для проверки подлинности и авторизации через Microsoft Entra ID при доступе к ресурсу Azure, ознакомьтесь с разделом Авторизация доступа с помощью Microsoft Entra ID, чтобы убедиться, что субъекту безопасности предоставлены достаточные разрешения для доступа к ресурсу Azure.

Настраиваемые свойства spring-cloud-azure-starter-integration-eventhubsподключения:

Свойство Тип Описание
spring.cloud.azure.eventhubs.enabled булев Включена ли служба Центров событий Azure.
spring.cloud.azure.eventhubs.connection-string Струна Значение строки подключения для пространства имен Центров событий.
spring.cloud.azure.eventhubs.namespace Струна Значение пространства имен Event Hubs, которое является префиксом FQDN. Полное доменное имя должно состоять из NamespaceName.DomainName
spring.cloud.azure.eventhubs.domain-name Струна Значение доменного имени пространства имен Центры событий Azure.
spring.cloud.azure.eventhubs.custom-endpoint-address Струна Адрес настраиваемой конечной точки.
spring.cloud.azure.eventhubs.shared-connection логический Используют ли базовые EventProcessorClient и EventHubProducerAsyncClient одно и то же подключение. По умолчанию для каждого созданного клиента Центров событий создается и используется отдельное новое подключение.

Свойства конфигурации контрольной точки

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

Заметка

В версии 4.0.0, если свойство spring.cloud.azure.eventhubs.processor.checkpoint-store.create-container-if-not-exists не включено вручную, контейнер хранилища не будет создан автоматически.

Создание контрольных точек для настраиваемых свойств spring-cloud-azure-starter-integration-eventhubs:

Свойство Тип Описание
spring.cloud.azure.eventhubs.processor.checkpoint-store.create-container-if-not-exists логический Разрешить ли создание контейнеров, если они отсутствуют.
spring.cloud.azure.eventhubs.processor.checkpoint-store.account-name Струна Имя учетной записи хранения.
spring.cloud.azure.eventhubs.processor.checkpoint-store.account-key Струна Ключ доступа к учетной записи хранения.
spring.cloud.azure.eventhubs.processor.checkpoint-store.container-name Струна Имя контейнера хранилища.

Общие параметры конфигурации пакета SDK служб Azure также можно настроить для хранилища контрольных точек для Blob Storage. Поддерживаемые параметры конфигурации представлены в конфигурации Spring Cloud Azureи могут быть настроены с помощью единого префикса spring.cloud.azure. или префикса spring.cloud.azure.eventhubs.processor.checkpoint-store..

Свойства конфигурации процессора Концентратора событий

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

Базовое использование

Отправка сообщений в Центры событий Azure

  1. Заполните параметры конфигурации учетных данных.

    • Для учетных данных в качестве строки подключения настройте следующие свойства в файле application.yml:

      spring:
        cloud:
          azure:
            eventhubs:
              connection-string: ${AZURE_EVENT_HUBS_CONNECTION_STRING}
              processor:
                checkpoint-store:
                  container-name: ${CHECKPOINT-CONTAINER}
                  account-name: ${CHECKPOINT-STORAGE-ACCOUNT}
                  account-key: ${CHECKPOINT-ACCESS-KEY}
      

      Заметка

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

    • Для учетных данных в качестве управляемых удостоверений настройте следующие свойства в файле application.yml:

      spring:
        cloud:
          azure:
            credential:
              managed-identity-enabled: true
              client-id: ${AZURE_CLIENT_ID}
            eventhubs:
              namespace: ${AZURE_EVENT_HUBS_NAMESPACE}
              processor:
                checkpoint-store:
                  container-name: ${CONTAINER_NAME}
                  account-name: ${ACCOUNT_NAME}
      
    • Для учетных данных в качестве субъекта-службы настройте следующие свойства в файле application.yml:

      spring:
        cloud:
          azure:
            credential:
              client-id: ${AZURE_CLIENT_ID}
              client-secret: ${AZURE_CLIENT_SECRET}
            profile:
              tenant-id: <tenant>
            eventhubs:
              namespace: ${AZURE_EVENT_HUBS_NAMESPACE}
              processor:
                checkpoint-store:
                  container-name: ${CONTAINER_NAME}
                  account-name: ${ACCOUNT_NAME}
      

Заметка

Значения, допустимые для tenant-id: common, organizations, consumersили идентификатор клиента. Дополнительные сведения об этих значениях см. в разделе Использована неправильная конечная точка (личные и организационные учетные записи) статьи Ошибка AADSTS50020: учетная запись пользователя от поставщика удостоверений не существует в арендаторе. Сведения о преобразовании однотенантного приложения см. в статье Преобразование однотенантного приложения в мультитенантное в Microsoft Entra ID.

  1. Создайте DefaultMessageHandler, используя компонент EventHubsTemplate, чтобы отправлять сообщения в Центры событий.

    class Demo {
        private static final String OUTPUT_CHANNEL = "output";
        private static final String EVENTHUB_NAME = "eh1";
    
        @Bean
        @ServiceActivator(inputChannel = OUTPUT_CHANNEL)
        public MessageHandler messageSender(EventHubsTemplate eventHubsTemplate) {
            DefaultMessageHandler handler = new DefaultMessageHandler(EVENTHUB_NAME, eventHubsTemplate);
            handler.setSendCallback(new ListenableFutureCallback<Void>() {
                @Override
                public void onSuccess(Void result) {
                    LOGGER.info("Message was sent successfully.");
                }
                @Override
                public void onFailure(Throwable ex) {
                    LOGGER.error("There was an error sending the message.", ex);
                }
            });
            return handler;
        }
    }
    
  2. Создайте привязку шлюза сообщений с указанным выше обработчиком сообщений через канал сообщений.

    class Demo {
        @Autowired
        EventHubOutboundGateway messagingGateway;
    
        @MessagingGateway(defaultRequestChannel = OUTPUT_CHANNEL)
        public interface EventHubOutboundGateway {
            void send(String text);
        }
    }
    
  3. Отправка сообщений с помощью шлюза.

    class Demo {
        public void demo() {
            this.messagingGateway.send(message);
        }
    }
    

Получение сообщений из Центров событий Azure

  1. Заполните параметры конфигурации учетных данных.

  2. Создайте компонент (bean) канала сообщений для использования в качестве входного канала.

    @Configuration
    class Demo {
        @Bean
        public MessageChannel input() {
            return new DirectChannel();
        }
    }
    
  3. Создайте EventHubsInboundChannelAdapter с помощью EventHubsMessageListenerContainer bean для получения сообщений из Центров событий.

    @Configuration
    class Demo {
        private static final String INPUT_CHANNEL = "input";
        private static final String EVENTHUB_NAME = "eh1";
        private static final String CONSUMER_GROUP = "$Default";
    
        @Bean
        public EventHubsInboundChannelAdapter messageChannelAdapter(
                @Qualifier(INPUT_CHANNEL) MessageChannel inputChannel,
                EventHubsMessageListenerContainer listenerContainer) {
            EventHubsInboundChannelAdapter adapter = new EventHubsInboundChannelAdapter(listenerContainer);
            adapter.setOutputChannel(inputChannel);
            return adapter;
        }
    
        @Bean
        public EventHubsMessageListenerContainer messageListenerContainer(EventHubsProcessorFactory processorFactory) {
            EventHubsContainerProperties containerProperties = new EventHubsContainerProperties();
            containerProperties.setEventHubName(EVENTHUB_NAME);
            containerProperties.setConsumerGroup(CONSUMER_GROUP);
            containerProperties.setCheckpointConfig(new CheckpointConfig(CheckpointMode.MANUAL));
            return new EventHubsMessageListenerContainer(processorFactory, containerProperties);
        }
    }
    
  4. Создайте привязку приемника сообщений с помощью EventHubsInboundChannelAdapter через канал сообщений, созданный ранее.

    class Demo {
        @ServiceActivator(inputChannel = INPUT_CHANNEL)
        public void messageReceiver(byte[] payload, @Header(AzureHeaders.CHECKPOINTER) Checkpointer checkpointer) {
            String message = new String(payload);
            LOGGER.info("New message received: '{}'", message);
            checkpointer.success()
                    .doOnSuccess(s -> LOGGER.info("Message '{}' successfully checkpointed", message))
                    .doOnError(e -> LOGGER.error("Error found", e))
                    .block();
        }
    }
    

Настройка EventHubsMessageConverter для настройки objectMapper

EventHubsMessageConverter реализован как настраиваемый компонент, чтобы пользователи могли настроить ObjectMapper.

Поддержка пакетного потребителя

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

При создании EventHubsInboundChannelAdapterрежим прослушивателя должен быть задан как BATCH. При создании компонента EventHubsMessageListenerContainer установите режим создания контрольных точек как MANUAL либо BATCH; параметры пакетной обработки можно настроить по мере необходимости.

@Configuration
class Demo {
    private static final String INPUT_CHANNEL = "input";
    private static final String EVENTHUB_NAME = "eh1";
    private static final String CONSUMER_GROUP = "$Default";

    @Bean
    public EventHubsInboundChannelAdapter messageChannelAdapter(
            @Qualifier(INPUT_CHANNEL) MessageChannel inputChannel,
            EventHubsMessageListenerContainer listenerContainer) {
        EventHubsInboundChannelAdapter adapter = new EventHubsInboundChannelAdapter(listenerContainer, ListenerMode.BATCH);
        adapter.setOutputChannel(inputChannel);
        return adapter;
    }

    @Bean
    public EventHubsMessageListenerContainer messageListenerContainer(EventHubsProcessorFactory processorFactory) {
        EventHubsContainerProperties containerProperties = new EventHubsContainerProperties();
        containerProperties.setEventHubName(EVENTHUB_NAME);
        containerProperties.setConsumerGroup(CONSUMER_GROUP);
        containerProperties.getBatch().setMaxSize(100);
        containerProperties.setCheckpointConfig(new CheckpointConfig(CheckpointMode.MANUAL));
        return new EventHubsMessageListenerContainer(processorFactory, containerProperties);
    }
}

Заголовки сообщений Центров событий

В следующей таблице показано, как свойства сообщений Центров событий сопоставляются с заголовками сообщений Spring. Для Центров событий Azure сообщение вызывается как event.

Сопоставление между свойствами сообщений и событий Event Hubs и заголовками сообщений Spring в режиме прослушивателя записей:

Свойства событий в Центрах событий Константы заголовков сообщений Spring Тип Описание
Время постановки в очередь EventHubsHeaders#ENQUEUED_TIME Мгновение Момент времени по UTC, когда событие было поставлено в очередь в разделе Event Hub.
Смещение EventHubsHeaders#OFFSET Длинный Смещение события в момент его получения из связанного раздела Event Hub.
Ключ раздела AzureHeaders#PARTITION_KEY Струна Ключ хэширования раздела, если он был задан при первоначальной публикации события.
Идентификатор раздела AzureHeaders#RAW_PARTITION_ID Струна Идентификатор раздела Центра событий.
Порядковый номер EventHubsHeaders#SEQUENCE_NUMBER Длинный Порядковый номер, присвоенный событию, когда оно было поставлено в очередь в соответствующем разделе Центра событий.
Свойства последнего события, помещённого в очередь EventHubsHeaders#LAST_ENQUEUED_EVENT_PROPERTIES Свойства последнего поставленного в очередь события Свойства последнего события, добавленного в очередь в этом разделе.
NA AzureHeaders#CHECKPOINTER Средство создания контрольных точек Заголовок для контрольной точки определенного сообщения.

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

Заметка

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

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

Сопоставление между сообщением / свойствами события Event Hubs и заголовками сообщений Spring в режиме пакетного прослушивателя:

Свойства событий в Центрах событий Константы заголовка сообщений Spring Batch Тип Описание
Время постановки в очередь EventHubsHeaders#ENQUEUED_TIME Список Instant Список моментов времени (UTC), когда каждое событие было поставлено в очередь в разделе концентратора событий.
Смещение EventHubsHeaders#OFFSET Список длинных Список смещений каждого события на момент его получения из связанного раздела Центра событий.
Ключ раздела AzureHeaders#PARTITION_KEY Список строк Список ключей хеширования раздела, если такой ключ был задан при первоначальной публикации каждого события.
Порядковый номер EventHubsHeaders#SEQUENCE_NUMBER Список длинных Список номеров последовательности, назначенных каждому событию при его постановке в очередь в соответствующем разделе Event Hub.
Системные свойства EventHubsHeaders#BATCH_CONVERTED_SYSTEM_PROPERTIES Список карт Список системных свойств каждого события.
Свойства приложения EventHubsHeaders#BATCH_CONVERTED_APPLICATION_PROPERTIES Список карт Список свойств приложения каждого события, где размещаются все настраиваемые заголовки сообщений или свойства события.

Заметка

При публикации сообщений все указанные выше заголовки пакета будут удалены из сообщений, если они существуют.

Образцы

Дополнительные сведения см. в репозитории azure-spring-boot-samples на сайте GitHub.

Интеграция Spring с служебной шиной Azure

Основные понятия

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

Проект расширения Spring Integration для Служебная шина Azure предоставляет адаптеры входящих и исходящих каналов для Служебная шина Azure.

Заметка

API поддержки CompletableFuture объявлены устаревшими, начиная с версии 2.10.0, и заменены на Reactor Core, начиная с версии 4.0.0. Дополнительные сведения см. в Javadoc.

Настройка зависимостей

<dependency>
    <groupId>com.azure.spring</groupId>
    <artifactId>spring-cloud-azure-starter-integration-servicebus</artifactId>
</dependency>

Конфигурация

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

Свойства конфигурации подключения

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

Заметка

Если вы решили использовать субъект безопасности для проверки подлинности и авторизации через Microsoft Entra ID при доступе к ресурсу Azure, ознакомьтесь с разделом Авторизация доступа с помощью Microsoft Entra ID, чтобы убедиться, что субъекту безопасности предоставлены достаточные разрешения для доступа к ресурсу Azure.

Настраиваемые свойства spring-cloud-azure-starter-integration-servicebusподключения:

Свойство Тип Описание
spring.cloud.azure.servicebus.enabled булев Включена ли служебная шина Azure.
spring.cloud.azure.servicebus.connection-string Струна Значение строки подключения пространства имен служебная шина.
spring.cloud.azure.servicebus.custom-endpoint-address Струна Адрес настраиваемой конечной точки, используемый при подключении к служебная шина.
spring.cloud.azure.servicebus.namespace Струна Значение пространства имен служебная шина, которое является префиксом полного доменного имени. Полное доменное имя должно состоять из NamespaceName.DomainName
spring.cloud.azure.servicebus.domain-name Струна Доменное имя значения пространства имен служебной шины Azure.

Свойства конфигурации процессора служебной шины

ServiceBusInboundChannelAdapter использует ServiceBusProcessorClient для получения сообщений, а для настройки общих параметров ServiceBusProcessorClient разработчики могут использовать ServiceBusContainerProperties. См. следующем разделе о работе с ServiceBusInboundChannelAdapter.

Базовое использование

Отправка сообщений в служебную шину Azure

  1. Заполните параметры конфигурации учетных данных.

    • Для учетных данных в качестве строки подключения настройте следующие свойства в файле application.yml:

      spring:
        cloud:
          azure:
            servicebus:
              connection-string: ${AZURE_SERVICE_BUS_CONNECTION_STRING}
      

      Заметка

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

    • Для учетных данных в качестве управляемых удостоверений настройте следующие свойства в файле application.yml:

      spring:
        cloud:
          azure:
            credential:
              managed-identity-enabled: true
              client-id: ${AZURE_CLIENT_ID}
            profile:
              tenant-id: <tenant>
            servicebus:
              namespace: ${AZURE_SERVICE_BUS_NAMESPACE}
      

Заметка

Значения, допустимые для tenant-id: common, organizations, consumersили идентификатор клиента. Дополнительные сведения об этих значениях см. в разделе Используется неправильная конечная точка (личные учетные записи и учетные записи организации) статьи Ошибка AADSTS50020 — учетная запись пользователя от поставщика удостоверений не существует в клиенте. Сведения о преобразовании однотенантного приложения см. в статье Преобразование однотенантного приложения в мультитенантное в Microsoft Entra ID.

  • Для учетных данных в качестве субъекта-службы настройте следующие свойства в файле application.yml:

    spring:
      cloud:
        azure:
          credential:
            client-id: ${AZURE_CLIENT_ID}
            client-secret: ${AZURE_CLIENT_SECRET}
          profile:
            tenant-id: <tenant>
          servicebus:
            namespace: ${AZURE_SERVICE_BUS_NAMESPACE}
    

Заметка

Значения, допустимые для tenant-id: common, organizations, consumersили идентификатор клиента. Дополнительные сведения об этих значениях см. в разделе Использована неправильная конечная точка (личные и организационные учетные записи) статьи Ошибка AADSTS50020: учетная запись пользователя от поставщика удостоверений не существует в арендаторе. Сведения о преобразовании однотенантного приложения см. в статье Преобразование однотенантного приложения в мультитенантное в Microsoft Entra ID.

  1. Создайте DefaultMessageHandler с компонентом ServiceBusTemplate для отправки сообщений в служебная шина; задайте для ServiceBusTemplate тип сущности. В этом примере рассматривается очередь служебная шина.

    class Demo {
        private static final String OUTPUT_CHANNEL = "queue.output";
    
        @Bean
        @ServiceActivator(inputChannel = OUTPUT_CHANNEL)
        public MessageHandler queueMessageSender(ServiceBusTemplate serviceBusTemplate) {
            serviceBusTemplate.setDefaultEntityType(ServiceBusEntityType.QUEUE);
            DefaultMessageHandler handler = new DefaultMessageHandler(QUEUE_NAME, serviceBusTemplate);
            handler.setSendCallback(new ListenableFutureCallback<Void>() {
                @Override
                public void onSuccess(Void result) {
                    LOGGER.info("Message was sent successfully.");
                }
    
                @Override
                public void onFailure(Throwable ex) {
                    LOGGER.error("There was an error sending the message.", ex);
                }
            });
    
            return handler;
        }
    }
    
  2. Создайте привязку шлюза сообщений с указанным выше обработчиком сообщений через канал сообщений.

    class Demo {
        @Autowired
        QueueOutboundGateway messagingGateway;
    
        @MessagingGateway(defaultRequestChannel = OUTPUT_CHANNEL)
        public interface QueueOutboundGateway {
            void send(String text);
        }
    }
    
  3. Отправка сообщений с помощью шлюза.

    class Demo {
        public void demo() {
            this.messagingGateway.send(message);
        }
    }
    

Получение сообщений из служебной шины Azure

  1. Заполните параметры конфигурации учетных данных.

  2. Создайте компонент (bean) канала сообщений для использования в качестве входного канала.

    @Configuration
    class Demo {
        private static final String INPUT_CHANNEL = "input";
    
        @Bean
        public MessageChannel input() {
            return new DirectChannel();
        }
    }
    
  3. Создайте ServiceBusInboundChannelAdapter с компонентом ServiceBusMessageListenerContainer для отправки сообщений в служебная шина. В этом примере рассматривается очередь служебная шина.

    @Configuration
    class Demo {
        private static final String QUEUE_NAME = "queue1";
    
        @Bean
        public ServiceBusMessageListenerContainer messageListenerContainer(ServiceBusProcessorFactory processorFactory) {
            ServiceBusContainerProperties containerProperties = new ServiceBusContainerProperties();
            containerProperties.setEntityName(QUEUE_NAME);
            containerProperties.setAutoComplete(false);
            return new ServiceBusMessageListenerContainer(processorFactory, containerProperties);
        }
    
        @Bean
        public ServiceBusInboundChannelAdapter queueMessageChannelAdapter(
            @Qualifier(INPUT_CHANNEL) MessageChannel inputChannel,
            ServiceBusMessageListenerContainer listenerContainer) {
            ServiceBusInboundChannelAdapter adapter = new ServiceBusInboundChannelAdapter(listenerContainer);
            adapter.setOutputChannel(inputChannel);
            return adapter;
        }
    }
    
  4. Создайте привязку приемника сообщений с ServiceBusInboundChannelAdapter через канал сообщений, который мы создали ранее.

    class Demo {
        @ServiceActivator(inputChannel = INPUT_CHANNEL)
        public void messageReceiver(byte[] payload, @Header(AzureHeaders.CHECKPOINTER) Checkpointer checkpointer) {
            String message = new String(payload);
            LOGGER.info("New message received: '{}'", message);
            checkpointer.success()
                    .doOnSuccess(s -> LOGGER.info("Message '{}' successfully checkpointed", message))
                    .doOnError(e -> LOGGER.error("Error found", e))
                    .block();
        }
    }
    

Настройка ServiceBusMessageConverter для настройки objectMapper

ServiceBusMessageConverter реализован как настраиваемый bean-компонент, чтобы пользователи могли настраивать ObjectMapper.

Заголовки сообщений служебная шина

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

Сопоставление между заголовками служебная шина и заголовками Spring:

Заголовки и свойства сообщений служебная шина Константы заголовка сообщения Spring Тип Конфигурируемый Описание
Тип контента MessageHeaders#CONTENT_TYPE Струна Да Дескриптор типа контента RFC2045 сообщения.
Идентификатор корреляции ServiceBusMessageHeaders#CORRELATION_ID Струна Да Идентификатор корреляции сообщения
Идентификатор сообщения ServiceBusMessageHeaders#MESSAGE_ID Струна Да Идентификатор сообщения, этот заголовок имеет более высокий приоритет, чем MessageHeaders#ID.
Идентификатор сообщения MessageHeaders#ID UUID (Универсальный уникальный идентификатор) Да Идентификатор сообщения, этот заголовок имеет более низкий приоритет, чем ServiceBusMessageHeaders#MESSAGE_ID.
Ключ раздела ServiceBusMessageHeaders#PARTITION_KEY Струна Да Ключ раздела для отправки сообщения в секционированный объект.
Ответить на MessageHeaders#REPLY_CHANNEL Струна Да Адрес сущности, на который следует отправлять ответы.
Ответ на идентификатор сеанса ServiceBusMessageHeaders#REPLY_TO_SESSION_ID Струна Да Значение свойства ReplyToGroupId сообщения.
Запланированное время постановки в очередь (UTC) ServiceBusMessageHeaders#SCHEDULED_ENQUEUE_TIME СмещениеДатыВремени Да Дата и время, в которые сообщение должно быть помещено в очередь в служебная шина; этот заголовок имеет более высокий приоритет, чем AzureHeaders#SCHEDULED_ENQUEUE_MESSAGE.
Запланированное время постановки в очередь (UTC) AzureHeaders#SCHEDULED_ENQUEUE_MESSAGE Целое число Да Дата и время, в которые сообщение должно быть поставлено в очередь в служебная шина; приоритет этого заголовка ниже, чем у ServiceBusMessageHeaders#SCHEDULED_ENQUEUE_TIME.
Идентификатор сеанса ServiceBusMessageHeaders#SESSION_ID Струна Да Идентификатор сеанса для сущности, поддерживающей сеанс.
Время жить ServiceBusMessageHeaders#TIME_TO_LIVE Длительность Да Длительность времени до истечения срока действия этого сообщения.
Кому ServiceBusMessageHeaders#TO Струна Да Адрес сообщения "to", зарезервированный для дальнейшего использования в сценариях маршрутизации и в настоящее время игнорируемый самим брокером.
Тема ServiceBusMessageHeaders#SUBJECT Струна Да Тема сообщения.
Описание ошибки недоставленной буквы ServiceBusMessageHeaders#DEAD_LETTER_ERROR_DESCRIPTION Струна Нет Описание сообщения, помещенного в очередь недоставленных сообщений.
Причина недоставленных писем ServiceBusMessageHeaders#DEAD_LETTER_REASON Струна Нет Причина, по которой сообщение было перемещено в очередь недоставленных сообщений.
Источник недоставленных писем ServiceBusMessageHeaders#DEAD_LETTER_SOURCE Струна Нет Сущность, в которую было помещено сообщение в очередь недоставленных сообщений.
Количество доставки ServiceBusMessageHeaders#DELIVERY_COUNT длинный Нет Количество случаев доставки этого сообщения клиентам.
Порядковый номер в очереди ServiceBusMessageHeaders#ENQUEUED_SEQUENCE_NUMBER длинный Нет Номер последовательности, назначенный сообщению при постановке в очередь служебная шина.
Время постановки в очередь ServiceBusMessageHeaders#ENQUEUED_TIME СмещениеДатыВремени Нет Дата и время, когда это сообщение было поставлено в очередь в служебная шина.
Истекает по адресу ServiceBusMessageHeaders#EXPIRES_AT СмещениеДатыВремени Нет Дата окончания срока действия этого сообщения.
Токен блокировки ServiceBusMessageHeaders#LOCK_TOKEN Струна Нет Маркер блокировки для текущего сообщения.
Заблокировано до ServiceBusMessageHeaders#LOCKED_UNTIL СмещениеДатыВремени Нет Дата, в течение которого истекает блокировка этого сообщения.
Порядковый номер ServiceBusMessageHeaders#SEQUENCE_NUMBER длинный Нет Уникальный номер, присвоенный сообщению служебная шина.
Государство ServiceBusMessageHeaders#STATE ServiceBusMessageState Нет Состояние сообщения, которое может быть активным, отложенным или запланированным.

Поддержка ключа раздела

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

Рекомендуется: использовать ServiceBusMessageHeaders.PARTITION_KEY в качестве ключа заголовка.

public class SampleController {
    @PostMapping("/messages")
    public ResponseEntity<String> sendMessage(@RequestParam String message) {
        LOGGER.info("Going to add message {} to Sinks.Many.", message);
        many.emitNext(MessageBuilder.withPayload(message)
                                    .setHeader(ServiceBusMessageHeaders.PARTITION_KEY, "Customize partition key")
                                    .build(), Sinks.EmitFailureHandler.FAIL_FAST);
        return ResponseEntity.ok("Sent!");
    }
}

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

public class SampleController {
    @PostMapping("/messages")
    public ResponseEntity<String> sendMessage(@RequestParam String message) {
        LOGGER.info("Going to add message {} to Sinks.Many.", message);
        many.emitNext(MessageBuilder.withPayload(message)
                                    .setHeader(AzureHeaders.PARTITION_KEY, "Customize partition key")
                                    .build(), Sinks.EmitFailureHandler.FAIL_FAST);
        return ResponseEntity.ok("Sent!");
    }
}

Заметка

Если в заголовках сообщений заданы ServiceBusMessageHeaders.PARTITION_KEY и AzureHeaders.PARTITION_KEY, ServiceBusMessageHeaders.PARTITION_KEY предпочтительнее.

Поддержка сеансов

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

public class SampleController {
    @PostMapping("/messages")
    public ResponseEntity<String> sendMessage(@RequestParam String message) {
        LOGGER.info("Going to add message {} to Sinks.Many.", message);
        many.emitNext(MessageBuilder.withPayload(message)
                                    .setHeader(ServiceBusMessageHeaders.SESSION_ID, "Customize session ID")
                                    .build(), Sinks.EmitFailureHandler.FAIL_FAST);
        return ResponseEntity.ok("Sent!");
    }
}

Заметка

Если в заголовках сообщения задан ServiceBusMessageHeaders.SESSION_ID, а также задан другой заголовок ServiceBusMessageHeaders.PARTITION_KEY, значение идентификатора сеанса в конечном итоге заменит значение ключа раздела.

Настройка свойств клиента служебная шина

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

@Bean
public AzureServiceClientBuilderCustomizer<ServiceBusClientBuilder.ServiceBusSessionProcessorClientBuilder> customizeBuilder() {
    return builder -> builder.sessionIdleTimeout(Duration.ofSeconds(10));
}

Образцы

Дополнительные сведения см. в репозитории azure-spring-boot-samples на сайте GitHub.

Интеграция Spring с очередью хранилища Azure

Основные понятия

Хранилище очередей Azure — это служба для хранения большого количества сообщений. Вы обращаетесь к сообщениям из любой точки мира через прошедшие проверку подлинности вызовы с помощью HTTP или HTTPS. Сообщение очереди может быть размером до 64 КБ. Очередь может содержать миллионы сообщений до общего ограничения емкости учетной записи хранения. Очереди обычно используются для создания очереди задач, которые обрабатываются асинхронно.

Настройка зависимостей

<dependency>
    <groupId>com.azure.spring</groupId>
    <artifactId>spring-cloud-azure-starter-integration-storage-queue</artifactId>
</dependency>

Конфигурация

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

Свойства конфигурации подключения

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

Заметка

Если вы решили использовать субъект безопасности для проверки подлинности и авторизации через Microsoft Entra ID при доступе к ресурсу Azure, ознакомьтесь с разделом Авторизация доступа с помощью Microsoft Entra ID, чтобы убедиться, что субъекту безопасности предоставлены достаточные разрешения для доступа к ресурсу Azure.

Настраиваемые свойства spring-cloud-azure-starter-integration-storage-queueподключения:

Свойство Тип Описание
spring.cloud.azure.storage.queue.enabled булев Включена ли очередь службы хранилища Azure.
spring.cloud.azure.storage.queue.connection-string Струна Значение строки подключения пространства имен очереди хранилища.
spring.cloud.azure.storage.queue.accountName Струна Имя учетной записи очереди хранения.
spring.cloud.azure.storage.queue.accountKey Струна Ключ учетной записи очереди хранения.
spring.cloud.azure.storage.queue.endpoint Струна Конечная точка службы очередей хранилища
spring.cloud.azure.storage.queue.sasToken Струна Учетные данные токена SAS
spring.cloud.azure.storage.queue.serviceVersion QueueServiceVersion QueueServiceVersion, используемый при выполнении API-запросов.
spring.cloud.azure.storage.queue.messageEncoding Струна Кодировка сообщений очереди.

Базовое использование

Отправка сообщений в очередь службы хранилища Azure

  1. Заполните параметры конфигурации учетных данных.

    • Для учетных данных в качестве строки подключения настройте следующие свойства в файле application.yml:

      spring:
        cloud:
          azure:
            storage:
              queue:
                connection-string: ${AZURE_STORAGE_QUEUE_CONNECTION_STRING}
      

      Заметка

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

    • Для учетных данных в качестве управляемых удостоверений настройте следующие свойства в файле application.yml:

      spring:
        cloud:
          azure:
            credential:
              managed-identity-enabled: true
              client-id: ${AZURE_CLIENT_ID}
            profile:
              tenant-id: <tenant>
            storage:
              queue:
                account-name: ${AZURE_STORAGE_QUEUE_ACCOUNT_NAME}
      

Заметка

Значения, допустимые для tenant-id: common, organizations, consumersили идентификатор клиента. Дополнительные сведения об этих значениях см. в разделе Использована неправильная конечная точка (личные и организационные учетные записи) статьи Ошибка AADSTS50020: учетная запись пользователя от поставщика удостоверений не существует в арендаторе. Сведения о преобразовании однотенантного приложения см. в статье Преобразование однотенантного приложения в мультитенантное в Microsoft Entra ID.

  • Для учетных данных в качестве субъекта-службы настройте следующие свойства в файле application.yml:

    spring:
      cloud:
        azure:
          credential:
            client-id: ${AZURE_CLIENT_ID}
            client-secret: ${AZURE_CLIENT_SECRET}
          profile:
            tenant-id: <tenant>
          storage:
            queue:
              account-name: ${AZURE_STORAGE_QUEUE_ACCOUNT_NAME}
    

Заметка

Значения, допустимые для tenant-id: common, organizations, consumersили идентификатор клиента. Дополнительные сведения об этих значениях см. в разделе Использована неправильная конечная точка (личные и организационные учетные записи) статьи Ошибка AADSTS50020: учетная запись пользователя от поставщика удостоверений не существует в арендаторе. Сведения о преобразовании однотенантного приложения см. в статье Преобразование однотенантного приложения в мультитенантное в Microsoft Entra ID.

  1. Создайте DefaultMessageHandler с использованием bean-компонента StorageQueueTemplate для отправки сообщений в Storage Queue.

    class Demo {
        private static final String STORAGE_QUEUE_NAME = "example";
        private static final String OUTPUT_CHANNEL = "output";
    
        @Bean
        @ServiceActivator(inputChannel = OUTPUT_CHANNEL)
        public MessageHandler messageSender(StorageQueueTemplate storageQueueTemplate) {
            DefaultMessageHandler handler = new DefaultMessageHandler(STORAGE_QUEUE_NAME, storageQueueTemplate);
            handler.setSendCallback(new ListenableFutureCallback<Void>() {
                @Override
                public void onSuccess(Void result) {
                    LOGGER.info("Message was sent successfully.");
                }
    
                @Override
                public void onFailure(Throwable ex) {
                    LOGGER.error("There was an error sending the message.", ex);
                }
            });
            return handler;
        }
    }
    
  2. Создайте привязку шлюза сообщений с указанным выше обработчиком сообщений через канал сообщений.

    class Demo {
        @Autowired
        StorageQueueOutboundGateway storageQueueOutboundGateway;
    
        @MessagingGateway(defaultRequestChannel = OUTPUT_CHANNEL)
        public interface StorageQueueOutboundGateway {
            void send(String text);
        }
    }
    
  3. Отправка сообщений с помощью шлюза.

    class Demo {
        public void demo() {
            this.storageQueueOutboundGateway.send(message);
        }
    }
    

Получение сообщений из очереди службы хранилища Azure

  1. Заполните параметры конфигурации учетных данных.

  2. Создайте компонент (bean) канала сообщений для использования в качестве входного канала.

    class Demo {
        private static final String INPUT_CHANNEL = "input";
    
        @Bean
        public MessageChannel input() {
            return new DirectChannel();
        }
    }
    
  3. Создайте StorageQueueMessageSource с использованием bean-компонента StorageQueueTemplate для получения сообщений из Storage Queue.

    class Demo {
        private static final String STORAGE_QUEUE_NAME = "example";
    
        @Bean
        @InboundChannelAdapter(channel = INPUT_CHANNEL, poller = @Poller(fixedDelay = "1000"))
        public StorageQueueMessageSource storageQueueMessageSource(StorageQueueTemplate storageQueueTemplate) {
            return new StorageQueueMessageSource(STORAGE_QUEUE_NAME, storageQueueTemplate);
        }
    }
    
  4. Создайте привязку приемника сообщений с помощью StorageQueueMessageSource, созданной на последнем шаге через канал сообщений, который мы создали раньше.

    class Demo {
        @ServiceActivator(inputChannel = INPUT_CHANNEL)
        public void messageReceiver(byte[] payload, @Header(AzureHeaders.CHECKPOINTER) Checkpointer checkpointer) {
            String message = new String(payload);
            LOGGER.info("New message received: '{}'", message);
            checkpointer.success()
                .doOnError(Throwable::printStackTrace)
                .doOnSuccess(t -> LOGGER.info("Message '{}' successfully checkpointed", message))
                .block();
        }
    }
    

Образцы

Дополнительные сведения см. в репозитории azure-spring-boot-samples на сайте GitHub.