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

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

Фреймворк предоставляет гибкую модель программирования, основанную на уже устоявшихся и знакомых идиомах Spring и лучших практиках. Эти рекомендации включают поддержку постоянной семантики Pub/Sub, групп потребителей и разделов с сохранением состояния.

К текущим реализациям привязки относятся:

Spring Cloud Stream Binder для Центров событий Azure

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

Связующее ПО Spring Cloud Stream для Центры событий Azure предоставляет реализацию механизма привязки для фреймворка Spring Cloud Stream. В основе этой реализации лежат адаптеры каналов Spring Integration Event Hubs. С точки зрения проектирования центры событий похожи на Kafka. Кроме того, к центрам событий можно получить доступ через API Kafka. Если проект имеет жесткую зависимость от API Kafka, вы можете попробовать Концентратор событий с помощью примера API Kafka

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

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

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

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

Чтобы указать стратегию балансировки нагрузки, предоставляются свойства spring.cloud.stream.eventhubs.bindings.<binding-name>.consumer.load-balancing.*. Дополнительные сведения см. в разделе Свойства потребителя.

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

Связующее ПО Spring Cloud Azure Stream для Event Hubs поддерживает функцию пакетного потребителя Spring Cloud Stream.

Чтобы работать с режимом пакетного потребителя, задайте для свойства spring.cloud.stream.bindings.<binding-name>.consumer.batch-mode значение true. При включении принимается сообщение с полезной нагрузкой в виде списка пакетных событий и передаётся в функцию Consumer. Каждый заголовок сообщения также преобразуется в список, из которого содержимое является соответствующим значением заголовка, проанализированным из каждого события. Общие заголовки идентификаторов секций, контрольных точек и последних вложенных свойств представлены в виде одного значения, так как весь пакет событий использует одно и то же значение. Дополнительные сведения см. в разделе Заголовки сообщений Event Hubs в Spring Cloud поддержка Azure for Spring Integration.

Заметка

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

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

Размер пакета можно указать, задав свойства max-size и max-wait-time с префиксом spring.cloud.stream.eventhubs.bindings.<binding-name>.consumer.batch.. Необходимо свойство max-size, а свойство max-wait-time является необязательным. Дополнительные сведения см. в разделе Свойства потребителя.

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

<dependency>
    <groupId>com.azure.spring</groupId>
    <artifactId>spring-cloud-azure-stream-binder-eventhubs</artifactId>
</dependency>

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

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

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

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

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

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

Заметка

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

Настраиваемые свойства spring-cloud-azure-stream-binder-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 Струна Адрес настраиваемой конечной точки.

Совет

Общие параметры конфигурации Azure Service SDK также можно настроить для binder Spring Cloud Azure Stream Event Hubs. Поддерживаемые параметры конфигурации представлены в конфигурации Spring Cloud Azureи могут быть настроены с помощью единого префикса spring.cloud.azure. или префикса spring.cloud.azure.eventhubs..

Биндер также по умолчанию поддерживает Spring Could Azure Resource Manager. Чтобы узнать, как получить строку подключения с субъектами безопасности, которым не назначены роли, связанные с Data, см. раздел Базовое использование в Spring Cloud Azure Resource Manager.

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

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

Заметка

Начиная с версии 4.0.0, если свойство spring.cloud.azure.eventhubs.processor.checkpoint-store.create-container-if-not-exists не включено вручную, контейнер Storage не будет создаваться автоматически с именем из spring.cloud.stream.bindings.binding-name.destination.

Создание контрольных точек для настраиваемых свойств spring-cloud-azure-stream-binder-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.

Свойства конфигурации привязки Центров событий Azure

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

Свойства потребителя

Эти свойства предоставляются через EventHubsConsumerProperties.

Заметка

Чтобы избежать повторения, начиная с версии 4.17.0 и 5.11.0, Центры событий Azure Stream Binder Spring Cloud поддерживают параметры всех каналов в формате spring.cloud.stream.eventhubs.default.consumer.<property>=<value>.

Настраиваемые spring-cloud-azure-stream-binder-eventhubsсвойства потребителя:

Свойство Тип Описание
spring.cloud.stream.eventhubs.bindings.binding-name.consumer.checkpoint.mode Режим контрольной точки Режим контрольной точки, используемый, когда потребитель решает, как отправлять сообщение о контрольной точке
spring.cloud.stream.eventhubs.bindings.binding-name.consumer.checkpoint.count Целое число Определяет количество сообщений для каждого раздела, необходимое для создания одной контрольной точки. Вступают в силу только в том случае, если используется режим контрольной точки PARTITION_COUNT.
spring.cloud.stream.eventhubs.bindings.binding-name.consumer.checkpoint.interval Длительность Определяет интервал времени для выполнения одной контрольной точки. Вступают в силу только в том случае, если используется режим контрольной точки TIME.
spring.cloud.stream.eventhubs.bindings.<binding-name>.consumer.batch.max-size Целое число Максимальное количество событий в пакете. Требуется для режима пакетной обработки.
spring.cloud.stream.eventhubs.bindings.binding-name.consumer.batch.max-wait-time Длительность Максимальная продолжительность пакетной обработки. Вступают в силу только в том случае, если включен режим пакетного потребителя и является необязательным.
spring.cloud.stream.eventhubs.bindings.binding-name.consumer.load-balancing.update-interval Длительность Длительность интервала обновления.
spring.cloud.stream.eventhubs.bindings.binding-name.consumer.load-balancing.strategy Стратегия балансировки нагрузки Стратегия балансировки нагрузки.
spring.cloud.stream.eventhubs.bindings.binding-name.consumer.load-balancing.partition-ownership-expiration-interval Длительность Период времени, по истечении которого право владения разделом прекращается.
spring.cloud.stream.eventhubs.bindings.binding-name.consumer.track-last-enqueued-event-properties Логический Должен ли обработчик событий запрашивать сведения о последнем заквеченном событии в связанной секции и отслеживать эти сведения по мере получения событий.
spring.cloud.stream.eventhubs.bindings.binding-name.consumer.prefetch-count Целое число Количество, используемое потребителем для управления числом событий, которые потребитель Event Hub будет активно получать и помещать в локальную очередь.
spring.cloud.stream.eventhubs.bindings.binding-name.consumer.initial-partition-event-position Сопоставление, где ключом является идентификатор раздела, а значениями — StartPositionProperties Карта, содержащая позицию события, используемую для каждой секции, если контрольная точка для секции не существует в хранилище контрольных точек. Эта карта использует идентификатор раздела в качестве ключа.

Заметка

Конфигурация initial-partition-event-position принимает map, чтобы указать начальную позицию для каждого концентратора событий. Таким образом, его ключом является идентификатор раздела, а значением — StartPositionProperties, который содержит свойства смещения, номера последовательности, времени постановки в очередь и признака включительности. Например, его можно задать как

spring:
  cloud:
    stream:
      eventhubs:
        bindings:
          <binding-name>:
            consumer:
              initial-partition-event-position:
                0:
                  offset: earliest
                1:
                  sequence-number: 100
                2:
                  enqueued-date-time: 2022-01-12T13:32:47.650005Z
                4:
                  inclusive: false
Расширенные пользовательские настройки

Приведённые выше параметры подключения, контрольной точки и общей конфигурации клиента Azure SDK поддерживают настройку для каждого потребителя binder, которую можно задать с префиксом spring.cloud.stream.eventhubs.bindings.<binding-name>.consumer..

Свойства производителя

Эти свойства предоставляются через EventHubsProducerProperties.

Заметка

Чтобы избежать повторения, начиная с версии 4.17.0 и 5.11.0, Центры событий Azure Stream Binder Spring Cloud поддерживают параметры всех каналов в формате spring.cloud.stream.eventhubs.default.producer.<property>=<value>.

Настраиваемые свойства производителя spring-cloud-azure-stream-binder-eventhubs:

Свойство Тип Описание
spring.cloud.stream.eventhubs.bindings.binding-name.producer.sync булев Флаг переключателя для синхронизации производителя. Если задано значение true, производитель ожидает ответа после операции отправки.
spring.cloud.stream.eventhubs.bindings.binding-name.producer.send-timeout длинный Время ожидания ответа после операции отправки. Вступают в силу только в том случае, если производитель синхронизации включен.
Расширенная конфигурация производителя

Приведенные выше параметры подключения и общей конфигурации клиента Azure SDK поддерживают настройку для каждого производителя binder, которую можно задать с помощью префикса spring.cloud.stream.eventhubs.bindings.<binding-name>.producer..

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

Отправка и получение сообщений из центров событий

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

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

      spring:
        cloud:
          azure:
            eventhubs:
              connection-string: ${EVENTHUB_NAMESPACE_CONNECTION_STRING}
              processor:
                checkpoint-store:
                  container-name: ${CHECKPOINT_CONTAINER}
                  account-name: ${CHECKPOINT_STORAGE_ACCOUNT}
                  account-key: ${CHECKPOINT_ACCESS_KEY}
          function:
            definition: consume;supply
          stream:
            bindings:
              consume-in-0:
                destination: ${EVENTHUB_NAME}
                group: ${CONSUMER_GROUP}
              supply-out-0:
                destination: ${THE_SAME_EVENTHUB_NAME_AS_ABOVE}
            eventhubs:
              bindings:
                consume-in-0:
                  consumer:
                    checkpoint:
                      mode: MANUAL
      

      Заметка

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

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

      spring:
        cloud:
          azure:
            credential:
              client-id: ${AZURE_CLIENT_ID}
              client-secret: ${AZURE_CLIENT_SECRET}
            profile:
              tenant-id: <tenant>
            eventhubs:
              namespace: ${EVENTHUB_NAMESPACE}
              processor:
                checkpoint-store:
                  container-name: ${CONTAINER_NAME}
                  account-name: ${ACCOUNT_NAME}
          function:
            definition: consume;supply
          stream:
            bindings:
              consume-in-0:
                destination: ${EVENTHUB_NAME}
                group: ${CONSUMER_GROUP}
              supply-out-0:
                destination: ${THE_SAME_EVENTHUB_NAME_AS_ABOVE}
            eventhubs:
              bindings:
                consume-in-0:
                  consumer:
                    checkpoint:
                      mode: MANUAL
      

Заметка

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

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

    spring:
      cloud:
        azure:
          credential:
            managed-identity-enabled: true
            client-id: ${AZURE_MANAGED_IDENTITY_CLIENT_ID} # Only needed when using a user-assigned managed identity
          eventhubs:
            namespace: ${EVENTHUB_NAMESPACE}
            processor:
              checkpoint-store:
                container-name: ${CONTAINER_NAME}
                account-name: ${ACCOUNT_NAME}
        function:
          definition: consume;supply
        stream:
          bindings:
            consume-in-0:
              destination: ${EVENTHUB_NAME}
              group: ${CONSUMER_GROUP}
            supply-out-0:
              destination: ${THE_SAME_EVENTHUB_NAME_AS_ABOVE}
    
          eventhubs:
            bindings:
              consume-in-0:
                consumer:
                  checkpoint:
                    mode: MANUAL
    
  1. Определение поставщика и потребителя.

    @Bean
    public Consumer<Message<String>> consume() {
        return message -> {
            Checkpointer checkpointer = (Checkpointer) message.getHeaders().get(CHECKPOINTER);
            LOGGER.info("New message received: '{}', partition key: {}, sequence number: {}, offset: {}, enqueued time: {}",
                    message.getPayload(),
                    message.getHeaders().get(EventHubsHeaders.PARTITION_KEY),
                    message.getHeaders().get(EventHubsHeaders.SEQUENCE_NUMBER),
                    message.getHeaders().get(EventHubsHeaders.OFFSET),
                    message.getHeaders().get(EventHubsHeaders.ENQUEUED_TIME)
            );
    
            checkpointer.success()
                    .doOnSuccess(success -> LOGGER.info("Message '{}' successfully checkpointed", message.getPayload()))
                    .doOnError(error -> LOGGER.error("Exception found", error))
                    .block();
        };
    }
    
    @Bean
    public Supplier<Message<String>> supply() {
        return () -> {
            LOGGER.info("Sending message, sequence " + i);
            return MessageBuilder.withPayload("Hello world, " + i++).build();
        };
    }
    

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

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

схема с блок-схемой процесса поддержки секционирования.

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

  1. Укажите параметры пакетной конфигурации, как показано в следующем примере:

    spring:
      cloud:
        function:
          definition: consume
        stream:
          bindings:
            consume-in-0:
              destination: ${AZURE_EVENTHUB_NAME}
              group: ${AZURE_EVENTHUB_CONSUMER_GROUP}
              consumer:
                batch-mode: true
          eventhubs:
            bindings:
              consume-in-0:
                consumer:
                  batch:
                    max-batch-size: 10 # Required for batch-consumer mode
                    max-wait-time: 1m # Optional, the default value is null
                  checkpoint:
                    mode: BATCH # or MANUAL as needed
    
  2. Определение поставщика и потребителя.

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

    @Bean
    public Consumer<Message<List<String>>> consume() {
        return message -> {
            for (int i = 0; i < message.getPayload().size(); i++) {
                LOGGER.info("New message received: '{}', partition key: {}, sequence number: {}, offset: {}, enqueued time: {}",
                        message.getPayload().get(i),
                        ((List<Object>) message.getHeaders().get(EventHubsHeaders.BATCH_CONVERTED_PARTITION_KEY)).get(i),
                        ((List<Object>) message.getHeaders().get(EventHubsHeaders.BATCH_CONVERTED_SEQUENCE_NUMBER)).get(i),
                        ((List<Object>) message.getHeaders().get(EventHubsHeaders.BATCH_CONVERTED_OFFSET)).get(i),
                        ((List<Object>) message.getHeaders().get(EventHubsHeaders.BATCH_CONVERTED_ENQUEUED_TIME)).get(i));
            }
        };
    }
    
    @Bean
    public Supplier<Message<String>> supply() {
        return () -> {
            LOGGER.info("Sending message, sequence " + i);
            return MessageBuilder.withPayload("\"test"+ i++ +"\"").build();
        };
    }
    

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

    @Bean
    public Consumer<Message<List<String>>> consume() {
        return message -> {
            for (int i = 0; i < message.getPayload().size(); i++) {
                LOGGER.info("New message received: '{}', partition key: {}, sequence number: {}, offset: {}, enqueued time: {}",
                    message.getPayload().get(i),
                    ((List<Object>) message.getHeaders().get(EventHubHeaders.BATCH_CONVERTED_PARTITION_KEY)).get(i),
                    ((List<Object>) message.getHeaders().get(EventHubHeaders.BATCH_CONVERTED_SEQUENCE_NUMBER)).get(i),
                    ((List<Object>) message.getHeaders().get(EventHubHeaders.BATCH_CONVERTED_OFFSET)).get(i),
                    ((List<Object>) message.getHeaders().get(EventHubHeaders.BATCH_CONVERTED_ENQUEUED_TIME)).get(i));
            }
    
            Checkpointer checkpointer = (Checkpointer) message.getHeaders().get(CHECKPOINTER);
            checkpointer.success()
                        .doOnSuccess(success -> LOGGER.info("Message '{}' successfully checkpointed", message.getPayload()))
                        .doOnError(error -> LOGGER.error("Exception found", error))
                        .block();
        };
    }
    
    @Bean
    public Supplier<Message<String>> supply() {
        return () -> {
            LOGGER.info("Sending message, sequence " + i);
            return MessageBuilder.withPayload("\"test"+ i++ +"\"").build();
        };
    }
    

Заметка

В режиме пакетного потребления тип содержимого биндера Spring Cloud Stream по умолчанию равен application/json, поэтому убедитесь, что полезная нагрузка сообщения соответствует этому типу содержимого. Например, при использовании типа контента по умолчанию application/json для получения сообщений с полезными данными String полезные данные должны быть JSON String, окруженные двойными кавычками для исходного текста String. Хотя для типа контента text/plain это может быть объект String напрямую. Дополнительные сведения см. в разделе Согласование типов контента Spring Cloud Stream.

Обработка сообщений об ошибках

  • Обработка сообщений об ошибках исходящей привязки

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

    @ServiceActivator(inputChannel = IntegrationContextUtils.ERROR_CHANNEL_BEAN_NAME)
    public void handleError(ErrorMessage message) {
        LOGGER.error("Handling outbound binding error: " + message);
    }
    
  • Обработка сообщений об ошибках входящей привязки

    Spring Cloud Stream Event Hubs Binder поддерживает одно решение для обработки ошибок для привязок входящих сообщений: обработчики ошибок.

    обработчик ошибок:

    Spring Cloud Stream предоставляет механизм для предоставления пользовательского обработчика ошибок путем добавления Consumer, который принимает экземпляры ErrorMessage. Дополнительные сведения см. в разделе Обработка сообщений об ошибках в документации Spring Cloud Stream.

    • Стандартный обработчик ошибок привязки

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

      @Bean
      public Consumer<ErrorMessage> myDefaultHandler() {
          return message -> {
              // consume the error message
          };
      }
      

      Кроме того, необходимо задать для свойства spring.cloud.stream.default.error-handler-definition имя функции.

    • Обработчик ошибок, зависящих от привязки

      Настройте компонент Consumer bean для обработки определённых сообщений об ошибках входящей привязки. Следующая функция подписывается на определённый канал ошибок входящей привязки и имеет более высокий приоритет, чем обработчик ошибок привязки по умолчанию:

      @Bean
      public Consumer<ErrorMessage> myErrorHandler() {
          return message -> {
              // consume the error message
          };
      }
      

      Кроме того, необходимо задать для свойства spring.cloud.stream.bindings.<input-binding-name>.error-handler-definition имя функции.

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

Сведения об основных поддерживаемых заголовках сообщений см. в разделе Заголовки сообщений Event Hubs в документации Spring Cloud Azure для Spring Integration.

Поддержка нескольких биндеров

Подключение к нескольким пространствам имён Event Hubs также поддерживается с помощью нескольких биндеров. В этом примере в качестве примера используется строка подключения. Также поддерживаются учетные данные субъектов-служб и управляемых удостоверений. Связанные свойства можно задать в параметрах среды привязки.

  1. Чтобы использовать несколько привязок с Центрами событий, настройте следующие свойства в файле application.yml:

    spring:
      cloud:
        function:
          definition: consume1;supply1;consume2;supply2
        stream:
          bindings:
            consume1-in-0:
              destination: ${EVENTHUB_NAME_01}
              group: ${CONSUMER_GROUP_01}
            supply1-out-0:
              destination: ${THE_SAME_EVENTHUB_NAME_01_AS_ABOVE}
            consume2-in-0:
              binder: eventhub-2
              destination: ${EVENTHUB_NAME_02}
              group: ${CONSUMER_GROUP_02}
            supply2-out-0:
              binder: eventhub-2
              destination: ${THE_SAME_EVENTHUB_NAME_02_AS_ABOVE}
          binders:
            eventhub-1:
              type: eventhubs
              default-candidate: true
              environment:
                spring:
                  cloud:
                    azure:
                      eventhubs:
                        connection-string: ${EVENTHUB_NAMESPACE_01_CONNECTION_STRING}
                        processor:
                          checkpoint-store:
                            container-name: ${CHECKPOINT_CONTAINER_01}
                            account-name: ${CHECKPOINT_STORAGE_ACCOUNT}
                            account-key: ${CHECKPOINT_ACCESS_KEY}
            eventhub-2:
              type: eventhubs
              default-candidate: false
              environment:
                spring:
                  cloud:
                    azure:
                      eventhubs:
                        connection-string: ${EVENTHUB_NAMESPACE_02_CONNECTION_STRING}
                        processor:
                          checkpoint-store:
                            container-name: ${CHECKPOINT_CONTAINER_02}
                            account-name: ${CHECKPOINT_STORAGE_ACCOUNT}
                            account-key: ${CHECKPOINT_ACCESS_KEY}
          eventhubs:
            bindings:
              consume1-in-0:
                consumer:
                  checkpoint:
                    mode: MANUAL
              consume2-in-0:
                consumer:
                  checkpoint:
                    mode: MANUAL
          poller:
            initial-delay: 0
            fixed-delay: 1000
    

    Заметка

    В предыдущем файле приложения показано, как настроить один опрашиватель по умолчанию для приложения для всех привязок. Если вы хотите настроить механизм опроса для конкретной привязки, можно использовать конфигурацию, подобную spring.cloud.stream.bindings.<binding-name>.producer.poller.fixed-delay=3000.

    Заметка

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

  2. Нам нужно определить двух поставщиков и двух потребителей:

    @Bean
    public Supplier<Message<String>> supply1() {
        return () -> {
            LOGGER.info("Sending message1, sequence1 " + i);
            return MessageBuilder.withPayload("Hello world1, " + i++).build();
        };
    }
    
    @Bean
    public Supplier<Message<String>> supply2() {
        return () -> {
            LOGGER.info("Sending message2, sequence2 " + j);
            return MessageBuilder.withPayload("Hello world2, " + j++).build();
        };
    }
    
    @Bean
    public Consumer<Message<String>> consume1() {
        return message -> {
            Checkpointer checkpointer = (Checkpointer) message.getHeaders().get(CHECKPOINTER);
            LOGGER.info("New message1 received: '{}'", message);
            checkpointer.success()
                    .doOnSuccess(success -> LOGGER.info("Message1 '{}' successfully checkpointed", message))
                    .doOnError(error -> LOGGER.error("Exception found", error))
                    .block();
        };
    }
    
    @Bean
    public Consumer<Message<String>> consume2() {
        return message -> {
            Checkpointer checkpointer = (Checkpointer) message.getHeaders().get(CHECKPOINTER);
            LOGGER.info("New message2 received: '{}'", message);
            checkpointer.success()
                    .doOnSuccess(success -> LOGGER.info("Message2 '{}' successfully checkpointed", message))
                    .doOnError(error -> LOGGER.error("Exception found", error))
                    .block();
        };
    }
    

Подготовка ресурсов

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

spring:
  cloud:
    azure:
      credential:
        tenant-id: <tenant>
      profile:
        subscription-id: ${AZURE_SUBSCRIPTION_ID}
      eventhubs:
        resource:
          resource-group: ${AZURE_EVENTHUBS_RESOURCE_GROUP}

Заметка

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

Образцы

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

Spring Cloud Stream Binder для служебной шины Azure

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

Spring Cloud Stream Binder для служебной шины Azure предоставляет реализацию привязки для Spring Cloud Stream Framework. Эта реализация использует адаптеры канала служебной шины Spring Integration в своей основе.

Запланированное сообщение

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

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

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

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

<dependency>
    <groupId>com.azure.spring</groupId>
    <artifactId>spring-cloud-azure-stream-binder-servicebus</artifactId>
</dependency>

В качестве альтернативы можно также использовать Spring Cloud Azure Stream служебная шина Starter, как показано в следующем примере Maven:

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

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

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

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

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

Заметка

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

Настраиваемые свойства spring-cloud-azure-stream-binder-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.

Заметка

Общие параметры конфигурации Azure Service SDK также можно настраивать для биндера Spring Cloud Azure Stream служебная шина. Поддерживаемые параметры конфигурации представлены в конфигурации Spring Cloud Azureи могут быть настроены с помощью единого префикса spring.cloud.azure. или префикса spring.cloud.azure.servicebus..

Биндер также по умолчанию поддерживает Spring Could Azure Resource Manager. Чтобы узнать, как получить строку подключения с субъектами безопасности, которым не назначены роли, связанные с Data, см. раздел Базовое использование в Spring Cloud Azure Resource Manager.

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

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

Свойства потребителя

Эти свойства предоставляются через ServiceBusConsumerProperties.

Заметка

Чтобы избежать повторения, начиная с версии 4.17.0 и 5.11.0, служебная шина Azure Stream Binder Spring Cloud поддерживает значения для всех каналов в формате spring.cloud.stream.servicebus.default.consumer.<property>=<value>.

Настраиваемые spring-cloud-azure-stream-binder-servicebusсвойства потребителя:

Свойство Тип По умолчанию Описание
spring.cloud.stream.servicebus.bindings.binding-name.consumer.requeue-rejected булев ложь Если неудачные сообщения перенаправляются в DLQ.
spring.cloud.stream.servicebus.bindings.binding-name.consumer.max-concurrent-calls Целое число 1 Максимальное число одновременных сообщений, которые должен обрабатывать клиент обработчика служебной шины. Если сеанс включен, он применяется к каждому сеансу.
spring.cloud.stream.servicebus.bindings.binding-name.consumer.max-concurrent-sessions Целое число null Максимальное количество одновременных сеансов для обработки в любое время.
spring.cloud.stream.servicebus.bindings.binding-name.consumer.session-enabled Логический null Включен ли сеанс.
spring.cloud.stream.servicebus.bindings.binding-name.consumer.session-idle-timeout Длительность null Задает максимальное время (длительность) ожидания получения сообщения для текущего активного сеанса.
spring.cloud.stream.servicebus.bindings.binding-name.consumer.prefetch-count Целое число 0 Количество сообщений для предварительной выборки клиента-процессора служебная шина.
spring.cloud.stream.servicebus.bindings.binding-name.consumer.sub-queue Подочередь никакой Тип подочереди, к которой выполняется подключение.
spring.cloud.stream.servicebus.bindings.binding-name.consumer.max-auto-lock-renew-duration Длительность 5 м Время для продолжения автоматического продления блокировки.
spring.cloud.stream.servicebus.bindings.binding-name.consumer.receive-mode ServiceBusReceiveMode peek_lock Режим получения клиента-процессора служебная шина.
spring.cloud.stream.servicebus.bindings.binding-name.consumer.auto-complete логический истина Следует ли автоматически урегулировать сообщения. Если установлено значение false, в сообщение будет добавлен заголовок Checkpointer, чтобы разработчики могли вручную подтверждать сообщения.
spring.cloud.stream.servicebus.bindings.binding-name.consumer.max-size-in-megabytes Длинный 1024 Максимальный размер очереди или раздела в мегабайтах, который является размером памяти, выделенной для очереди или раздела.
spring.cloud.stream.servicebus.bindings.binding-name.consumer.default-message-time-to-live Длительность P10675199DT2H48M5.4775807S. (10675199 дней, 2 часа, 48 минут, 5 секунд и 477 миллисекунд) Период времени, по истечении которого срок действия сообщения истекает, начиная с момента отправки сообщения в служебная шина.

Важно

При использовании Azure Resource Manager (ARM) необходимо настроить свойство spring.cloud.stream.servicebus.bindings.<binding-name>.consume.entity-type. Дополнительные сведения см. в servicebus-queue-binder-arm примере на сайте GitHub.

Расширенные пользовательские настройки

Приведенные выше настройки подключения и общей конфигурации клиента Azure SDK поддерживают настройку для каждого потребителя привязки, которую можно задать с префиксом spring.cloud.stream.servicebus.bindings.<binding-name>.consumer..

Свойства производителя

Эти свойства предоставляются через ServiceBusProducerProperties.

Заметка

Чтобы избежать повторения, начиная с версии 4.17.0 и 5.11.0, служебная шина Azure Stream Binder Spring Cloud поддерживает значения для всех каналов в формате spring.cloud.stream.servicebus.default.producer.<property>=<value>.

Настраиваемые свойства производителя spring-cloud-azure-stream-binder-servicebus:

Свойство Тип По умолчанию Описание
spring.cloud.stream.servicebus.bindings.binding-name.producer.sync булев ложь Переключение флага для синхронизации производителя.
spring.cloud.stream.servicebus.bindings.binding-name.producer.send-timeout длинный 10 000 Значение тайм-аута отправки производителя сообщений.
spring.cloud.stream.servicebus.bindings.binding-name.producer.entity-type ServiceBusEntityType null Тип сущности служебной шины производителя, необходимый для производителя привязки.
spring.cloud.stream.servicebus.bindings.binding-name.producer.max-size-in-megabytes Длинный 1024 Максимальный размер очереди или раздела в мегабайтах, который является размером памяти, выделенной для очереди или раздела.
spring.cloud.stream.servicebus.bindings.binding-name.producer.default-message-time-to-live Длительность P10675199DT2H48M5.4775807S. (10675199 дней, 2 часа, 48 минут, 5 секунд и 477 миллисекунд) Период, по истечении которого срок действия сообщения истекает, начиная с момента отправки сообщения в служебная шина.

Важно

При использовании производителя привязки необходимо настроить свойство spring.cloud.stream.servicebus.bindings.<binding-name>.producer.entity-type.

Расширенная конфигурация производителя

Приведенные выше параметры подключения и общей конфигурации клиента Azure SDK поддерживают настройку для каждого производителя binder, которую можно задать с помощью префикса spring.cloud.stream.servicebus.bindings.<binding-name>.producer..

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

Отправка сообщений в служебная шина и получение из неё

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

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

      spring:
         cloud:
           azure:
             servicebus:
               connection-string: ${SERVICEBUS_NAMESPACE_CONNECTION_STRING}
           function:
             definition: consume;supply
           stream:
             bindings:
               consume-in-0:
                 destination: ${SERVICEBUS_ENTITY_NAME}
                 # If you use Service Bus Topic, add the following configuration
                 # group: ${SUBSCRIPTION_NAME}
               supply-out-0:
                 destination: ${SERVICEBUS_ENTITY_NAME_SAME_AS_ABOVE}
             servicebus:
               bindings:
                 consume-in-0:
                   consumer:
                     auto-complete: false
                 supply-out-0:
                   producer:
                     entity-type: queue # set as "topic" if you use Service Bus Topic
      

      Заметка

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

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

      spring:
         cloud:
           azure:
             credential:
               client-id: ${AZURE_CLIENT_ID}
               client-secret: ${AZURE_CLIENT_SECRET}
             profile:
               tenant-id: <tenant>
             servicebus:
               namespace: ${SERVICEBUS_NAMESPACE}
           function:
             definition: consume;supply
           stream:
             bindings:
               consume-in-0:
                 destination: ${SERVICEBUS_ENTITY_NAME}
                 # If you use Service Bus Topic, add the following configuration
                 # group: ${SUBSCRIPTION_NAME}
               supply-out-0:
                 destination: ${SERVICEBUS_ENTITY_NAME_SAME_AS_ABOVE}
             servicebus:
               bindings:
                 consume-in-0:
                   consumer:
                     auto-complete: false
                 supply-out-0:
                   producer:
                     entity-type: queue # set as "topic" if you use Service Bus Topic
      

Заметка

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

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

    spring:
      cloud:
        azure:
          credential:
            managed-identity-enabled: true
            client-id: ${MANAGED_IDENTITY_CLIENT_ID} # Only needed when using a user-assigned managed identity
          servicebus:
            namespace: ${SERVICEBUS_NAMESPACE}
        function:
          definition: consume;supply
        stream:
          bindings:
            consume-in-0:
              destination: ${SERVICEBUS_ENTITY_NAME}
              # If you use Service Bus Topic, add the following configuration
              # group: ${SUBSCRIPTION_NAME}
            supply-out-0:
              destination: ${SERVICEBUS_ENTITY_NAME_SAME_AS_ABOVE}
          servicebus:
            bindings:
              consume-in-0:
                consumer:
                  auto-complete: false
              supply-out-0:
                producer:
                  entity-type: queue # set as "topic" if you use Service Bus Topic
    
  1. Определение поставщика и потребителя.

    @Bean
    public Consumer<Message<String>> consume() {
        return message -> {
            Checkpointer checkpointer = (Checkpointer) message.getHeaders().get(CHECKPOINTER);
            LOGGER.info("New message received: '{}'", message.getPayload());
    
            checkpointer.success()
                    .doOnSuccess(success -> LOGGER.info("Message '{}' successfully checkpointed", message.getPayload()))
                    .doOnError(error -> LOGGER.error("Exception found", error))
                    .block();
        };
    }
    
    @Bean
    public Supplier<Message<String>> supply() {
        return () -> {
            LOGGER.info("Sending message, sequence " + i);
            return MessageBuilder.withPayload("Hello world, " + i++).build();
        };
    }
    

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

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

Spring Cloud Stream предоставляет свойство выражения spEL ключа раздела spring.cloud.stream.bindings.<binding-name>.producer.partition-key-expression. Например, присвойите этому свойству значение "'partitionKey-' + headers[<message-header-key>]" и добавьте заголовок с именем message-header-key. Spring Cloud Stream использует значение этого заголовка при оценке выражения для назначения ключа секции. В следующем коде представлен пример производителя:

@Bean
public Supplier<Message<String>> generate() {
    return () -> {
        String value = "random payload";
        return MessageBuilder.withPayload(value)
            .setHeader("<message-header-key>", value.length() % 4)
            .build();
    };
}

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

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

@Bean
public Supplier<Message<String>> generate() {
    return () -> {
        String value = "random payload";
        return MessageBuilder.withPayload(value)
            .setHeader(ServiceBusMessageHeaders.SESSION_ID, "Customize session ID")
            .build();
    };
}

Заметка

Согласно секционированию служебная шина, идентификатор сеанса имеет более высокий приоритет, чем ключ секционирования. Поэтому, если заданы оба заголовка ServiceBusMessageHeaders#SESSION_ID и ServiceBusMessageHeaders#PARTITION_KEY, значение идентификатора сеанса в итоге используется для перезаписи значения ключа раздела.

Обработка сообщений об ошибках

  • Обработка сообщений об ошибках исходящей привязки

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

    @ServiceActivator(inputChannel = IntegrationContextUtils.ERROR_CHANNEL_BEAN_NAME)
    public void handleError(ErrorMessage message) {
        LOGGER.error("Handling outbound binding error: " + message);
    }
    
  • Обработка сообщений об ошибках входящей привязки

    Биндер служебная шина для Spring Cloud Stream поддерживает два способа обработки ошибок для входящих привязок сообщений: обработчик ошибок биндинга и обработчики.

    Обработчик ошибок Binder:

    Стандартный обработчик ошибок компонента привязки обрабатывает входящую привязку. Этот обработчик используется для отправки неудачных сообщений в очередь недоставленных писем при включении spring.cloud.stream.servicebus.bindings.<binding-name>.consumer.requeue-rejected. В противном случае сообщения, обработка которых завершилась сбоем, переводятся в состояние abandoned. Обработчик ошибок привязки взаимоисключаем с другими указанными обработчиками ошибок.

    обработчик ошибок :

    Spring Cloud Stream предоставляет механизм для предоставления пользовательского обработчика ошибок путем добавления Consumer, который принимает экземпляры ErrorMessage. Дополнительные сведения см. в разделе Обработка сообщений об ошибках в документации Spring Cloud Stream.

    • Стандартный обработчик ошибок привязки

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

      @Bean
      public Consumer<ErrorMessage> myDefaultHandler() {
          return message -> {
              // consume the error message
          };
      }
      

      Кроме того, необходимо задать для свойства spring.cloud.stream.default.error-handler-definition имя функции.

    • Обработчик ошибок, зависящих от привязки

      Настройте компонент Consumer bean для обработки определённых сообщений об ошибках входящей привязки. Следующая функция подписывается на конкретный канал ошибок входящей привязки с более высоким приоритетом, чем обработчик ошибок привязки по умолчанию.

      @Bean
      public Consumer<ErrorMessage> myDefaultHandler() {
          return message -> {
              // consume the error message
          };
      }
      

      Кроме того, необходимо задать для свойства spring.cloud.stream.bindings.<input-binding-name>.error-handler-definition имя функции.

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

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

Заметка

При настройке ключа секции приоритет заголовка сообщения выше свойства Spring Cloud Stream. Поэтому spring.cloud.stream.bindings.<binding-name>.producer.partition-key-expression действует только в том случае, если ни один из ServiceBusMessageHeaders#SESSION_ID и ServiceBusMessageHeaders#PARTITION_KEY заголовков не настроен.

Поддержка нескольких биндеров

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

  1. Чтобы использовать несколько связующих компонентов ServiceBus, настройте следующие свойства в файле application.yml:

    spring:
      cloud:
        function:
          definition: consume1;supply1;consume2;supply2
        stream:
          bindings:
            consume1-in-0:
              destination: ${SERVICEBUS_TOPIC_NAME}
              group: ${SUBSCRIPTION_NAME}
            supply1-out-0:
              destination: ${SERVICEBUS_TOPIC_NAME_SAME_AS_ABOVE}
            consume2-in-0:
              binder: servicebus-2
              destination: ${SERVICEBUS_QUEUE_NAME}
            supply2-out-0:
              binder: servicebus-2
              destination: ${SERVICEBUS_QUEUE_NAME_SAME_AS_ABOVE}
          binders:
            servicebus-1:
              type: servicebus
              default-candidate: true
              environment:
                spring:
                  cloud:
                    azure:
                      servicebus:
                        connection-string: ${SERVICEBUS_NAMESPACE_01_CONNECTION_STRING}
            servicebus-2:
              type: servicebus
              default-candidate: false
              environment:
                spring:
                  cloud:
                    azure:
                      servicebus:
                        connection-string: ${SERVICEBUS_NAMESPACE_02_CONNECTION_STRING}
          servicebus:
            bindings:
              consume1-in-0:
                consumer:
                  auto-complete: false
              supply1-out-0:
                producer:
                  entity-type: topic
              consume2-in-0:
                consumer:
                  auto-complete: false
              supply2-out-0:
                producer:
                  entity-type: queue
          poller:
            initial-delay: 0
            fixed-delay: 1000
    

    Заметка

    В предыдущем файле приложения показано, как настроить один опрашиватель по умолчанию для приложения для всех привязок. Если вы хотите настроить механизм опроса для конкретной привязки, можно использовать конфигурацию, подобную spring.cloud.stream.bindings.<binding-name>.producer.poller.fixed-delay=3000.

    Заметка

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

  2. нам нужно определить двух поставщиков и двух потребителей

    @Bean
    public Supplier<Message<String>> supply1() {
        return () -> {
            LOGGER.info("Sending message1, sequence1 " + i);
            return MessageBuilder.withPayload("Hello world1, " + i++).build();
        };
    }
    
    @Bean
    public Supplier<Message<String>> supply2() {
        return () -> {
            LOGGER.info("Sending message2, sequence2 " + j);
            return MessageBuilder.withPayload("Hello world2, " + j++).build();
        };
    }
    
    @Bean
    public Consumer<Message<String>> consume1() {
        return message -> {
            Checkpointer checkpointer = (Checkpointer) message.getHeaders().get(CHECKPOINTER);
            LOGGER.info("New message1 received: '{}'", message);
            checkpointer.success()
                    .doOnSuccess(s -> LOGGER.info("Message '{}' successfully checkpointed", message.getPayload()))
                    .doOnError(e -> LOGGER.error("Error found", e))
                    .block();
        };
    }
    
    @Bean
    public Consumer<Message<String>> consume2() {
        return message -> {
            Checkpointer checkpointer = (Checkpointer) message.getHeaders().get(CHECKPOINTER);
            LOGGER.info("New message2 received: '{}'", message);
            checkpointer.success()
                    .doOnSuccess(s -> LOGGER.info("Message '{}' successfully checkpointed", message.getPayload()))
                    .doOnError(e -> LOGGER.error("Error found", e))
                    .block();
        };
    
    }
    

Подготовка ресурсов

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

spring:
  cloud:
    azure:
      credential:
        tenant-id: <tenant>
      profile:
        subscription-id: ${AZURE_SUBSCRIPTION_ID}
      servicebus:
        resource:
          resource-group: ${AZURE_SERVICEBUS_RESOURCE_GROUP}
    stream:
      servicebus:
        bindings:
          <binding-name>:
            consumer:
              entity-type: ${SERVICEBUS_CONSUMER_ENTITY_TYPE}

Заметка

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

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

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

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

Образцы

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