Suporte do Spring Cloud Azure para o Spring Integration

A Spring Integration Extension para Azure fornece adaptadores do Spring Integration para os vários serviços fornecidos pelo SDK do Azure para Java. Fornecemos suporte do Spring Integration para estes serviços do Azure: Event Hubs, Barramento de Serviço e Storage Queue. Veja a seguir uma lista de adaptadores com suporte:

Integração do Spring com os Hubs de Eventos do Azure

Principais conceitos

Os Hubs de Eventos do Azure são uma plataforma de streaming de Big Data e um serviço de ingestão de eventos. Ele pode receber e processar milhões de eventos por segundo. Os dados enviados para um hub de eventos podem ser transformados e armazenados usando qualquer provedor de análise em tempo real ou adaptadores de lote/armazenamento.

O Spring Integration permite mensagens leves em aplicativos baseados em Spring e dá suporte à integração com sistemas externos por meio de adaptadores declarativos. Esses adaptadores fornecem um nível mais alto de abstração sobre o suporte da Spring para comunicação remota, mensagens e agendamento. O projeto de extensão Spring Integration for Event Hubs fornece adaptadores e gateways de canais de entrada e saída para os Hubs de Eventos do Azure.

Nota

As APIs de suporte do RxJava são removidas da versão 4.0.0. Consulte Javadoc para obter detalhes.

Grupo de consumidores

Os Hubs de Eventos fornecem suporte semelhante ao grupo de consumidores como o Apache Kafka, mas com uma lógica ligeiramente diferente. Embora o Kafka armazene todos os offsets confirmados no broker, você precisa armazenar manualmente os offsets das mensagens do Event Hubs que estão sendo processadas. O SDK do Event Hubs oferece uma função para armazenar esses offsets no Armazenamento do Azure.

Suporte ao particionamento

Os Hubs de Eventos fornecem um conceito semelhante de partição física como Kafka. Mas, ao contrário do rebalanceamento automático do Kafka entre consumidores e partições, os Hubs de Eventos fornecem uma espécie de modo preemptivo. A conta de armazenamento atua como uma concessão para determinar qual partição pertence a qual consumidor. Quando um novo consumidor iniciar, ele tentará tomar algumas partições dos consumidores mais sobrecarregados para equilibrar a carga de trabalho.

Para especificar a estratégia de balanceamento de carga, os desenvolvedores podem usar EventHubsContainerProperties para a configuração. Consulte a seção a seguir para obter um exemplo de como configurar EventHubsContainerProperties.

Suporte a consumidor em lote

O EventHubsInboundChannelAdapter dá suporte ao modo de consumo em lote. Para habilitá-lo, os usuários podem especificar o modo de ouvinte como ListenerMode.BATCH ao construir uma instância de EventHubsInboundChannelAdapter. Quando habilitada, uma mensagem da qual o conteúdo é uma lista de eventos em lote será recebida e passada para o canal downstream. Cada cabeçalho da mensagem também é convertido em uma lista, cujo conteúdo é o valor de cabeçalho correspondente extraído de cada evento. Para os cabeçalhos comuns da ID da partição, do ponto de verificação e das últimas propriedades enfileiradas, eles são apresentados como um único valor porque todo o lote de eventos compartilha o mesmo. Para obter mais informações, consulte a seção cabeçalhos de mensagem dos Hubs de Eventos .

Nota

O cabeçalho de ponto de verificação só existe quando o modo de ponto de verificação MANUAL é usado.

O ponto de verificação do consumidor em lotes dá suporte a dois modos: BATCH e MANUAL. O modo BATCH é um modo do ponto de verificação automático para verificar todo o lote de eventos juntos depois que tiverem sido recebidos. O modo MANUAL é o ponto de verificação dos eventos para os usuários. Quando utilizado, o Checkpointer será incluído no cabeçalho da mensagem, e os usuários poderão usá-lo para realizar checkpointing.

A política de consumo em lote pode ser especificada por propriedades de max-size e max-wait-time, em que max-size é uma propriedade necessária enquanto max-wait-time é opcional. Para especificar a estratégia de consumo em lote, os desenvolvedores podem usar EventHubsContainerProperties para a configuração. Consulte a seção a seguir para obter um exemplo de como configurar EventHubsContainerProperties.

Configuração de dependência

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

Configuração

Este início apresenta as 3 partes a seguir das opções de configuração:

Propriedades de configuração de conexão

Esta seção contém as opções de configuração usadas para se conectar aos Hubs de Eventos do Azure.

Nota

Se você optar por usar uma entidade de segurança para autenticar e autorizar com a ID do Microsoft Entra para acessar um recurso do Azure, consulte Autorizar o acesso com a ID do Microsoft Entra para garantir que a entidade de segurança tenha recebido a permissão suficiente para acessar o recurso do Azure.

Propriedades configuráveis de conexão de spring-cloud-azure-starter-integration-eventhubs:

Propriedade Tipo Descrição
spring.cloud.azure.eventhubs.enabled booleano Se um Hub de Eventos do Azure está habilitado.
spring.cloud.azure.eventhubs.connection-string Corda Valor da cadeia de caracteres de conexão do namespace do Event Hubs.
spring.cloud.azure.eventhubs.namespace Corda Valor do Namespace dos Hubs de Eventos, que é o prefixo do FQDN. Um FQDN deve ser composto por NamespaceName.DomainName
spring.cloud.azure.eventhubs.domain-name Corda Nome de domínio de um valor do namespace do Hubs de Eventos do Azure.
spring.cloud.azure.eventhubs.custom-endpoint-address Corda Endereço de endpoint personalizado.
spring.cloud.azure.eventhubs.shared-connection booleano Se o EventProcessorClient e o EventHubProducerAsyncClient subjacentes usam a mesma conexão. Por padrão, uma nova conexão é estabelecida e usada para cada cliente do Event Hub criado.

Propriedades de configuração de ponto de verificação

Esta seção contém as opções de configuração do serviço de Blobs de Armazenamento, que é usado para persistir as informações de propriedade da partição e de ponto de verificação.

Nota

A partir da versão 4.0.0, quando a propriedade de spring.cloud.azure.eventhubs.processor.checkpoint-store.create-container-if-not-exists não estiver habilitada manualmente, nenhum contêiner de armazenamento será criado automaticamente.

Propriedades configuráveis de ponto de verificação de spring-cloud-azure-starter-integration-eventhubs:

Propriedade Tipo Descrição
spring.cloud.azure.eventhubs.processor.checkpoint-store.create-container-if-not-exists booleano Permitir a criação de contêineres caso não existam.
spring.cloud.azure.eventhubs.processor.checkpoint-store.account-name Corda Nome da conta de armazenamento.
spring.cloud.azure.eventhubs.processor.checkpoint-store.account-key Corda Chave de acesso da conta de armazenamento.
spring.cloud.azure.eventhubs.processor.checkpoint-store.container-name Corda Nome do contêiner de armazenamento.

As opções comuns de configuração do SDK do Serviço do Azure também são configuráveis para o repositório de ponto de verificação de Blob de Armazenamento. As opções de configuração com suporte são apresentadas em Configuração do Spring Cloud Azure e podem ser configuradas com o prefixo unificado spring.cloud.azure. ou com o prefixo spring.cloud.azure.eventhubs.processor.checkpoint-store..

Propriedades de configuração do processador do Hub de Eventos

O EventHubsInboundChannelAdapter usa o EventProcessorClient para consumir mensagens de um hub de eventos, para configurar as propriedades gerais de um EventProcessorClient, os desenvolvedores podem usar EventHubsContainerProperties para a configuração. Consulte a seção a seguir sobre como trabalhar com EventHubsInboundChannelAdapter.

Uso básico

Enviar mensagens aos Hubs de Eventos do Azure

  1. Preencha as opções de configuração de credencial.

    • Para credenciais como cadeia de conexão, configure as seguintes propriedades no arquivo 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}
      

      Nota

      A Microsoft recomenda usar o fluxo de autenticação mais seguro disponível. O fluxo de autenticação descrito nesse procedimento, como para bancos de dados, caches, mensagens ou serviços de IA, exige um grau muito alto de confiança no aplicativo e traz riscos não presentes em outros fluxos. Use esse fluxo somente quando opções mais seguras, como identidades gerenciadas para conexões sem senha ou sem chave, não forem viáveis. Para operações de máquinas locais, prefira identidades de usuário para conexões sem senha ou sem chave.

    • Para credenciais como identidades gerenciadas, configure as seguintes propriedades em seu arquivo 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}
      
    • Para credenciais de service principal, configure as seguintes propriedades no arquivo 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}
      

Nota

Os valores permitidos para tenant-id são: common, organizations, consumersou a ID do locatário. Para obter mais informações sobre esses valores, consulte a seção Uso do ponto de extremidade incorreto (contas pessoais e de organização) do Error AADSTS50020 – A conta de usuário do provedor de identidade não existe no locatário. Para obter informações sobre como converter seu aplicativo de locatário único, consulte Converter aplicativo de locatário único em multilocatário no Microsoft Entra ID.

  1. Crie DefaultMessageHandler com o bean EventHubsTemplate para enviar mensagens aos Hubs de Eventos.

    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. Crie um vínculo de gateway de mensagens com o manipulador de mensagens acima por meio de um canal de mensagens.

    class Demo {
        @Autowired
        EventHubOutboundGateway messagingGateway;
    
        @MessagingGateway(defaultRequestChannel = OUTPUT_CHANNEL)
        public interface EventHubOutboundGateway {
            void send(String text);
        }
    }
    
  3. Enviar mensagens usando o gateway.

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

Receber mensagens dos Hubs de Eventos do Azure

  1. Preencha as opções de configuração de credencial.

  2. Crie um bean de canal de mensagens que seja o canal de entrada.

    @Configuration
    class Demo {
        @Bean
        public MessageChannel input() {
            return new DirectChannel();
        }
    }
    
  3. Crie EventHubsInboundChannelAdapter com o bean EventHubsMessageListenerContainer para receber mensagens dos Hubs de Eventos.

    @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. Crie uma associação de receptor de mensagem com EventHubsInboundChannelAdapter por meio do canal de mensagem criado antes.

    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();
        }
    }
    

Configurar EventHubsMessageConverter para personalizar objectMapper

EventHubsMessageConverter é feito como um bean configurável para permitir que os usuários personalizem ObjectMapper.

Suporte a consumidor em lote

Consumir mensagens do Event Hubs em lotes é semelhante ao exemplo acima, mas os usuários também devem definir as opções de configuração relacionadas ao consumo em lotes para EventHubsInboundChannelAdapter.

Ao criar EventHubsInboundChannelAdapter, o modo de ouvinte deve ser definido como BATCH. Ao criar o bean de EventHubsMessageListenerContainer, defina o modo de ponto de verificação como MANUAL ou BATCH, e as opções de lote podem ser configuradas conforme necessário.

@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);
    }
}

Cabeçalhos de mensagem dos Hubs de Eventos

A tabela a seguir ilustra como as propriedades de mensagem dos Hubs de Eventos são mapeadas para cabeçalhos de mensagem spring. Para os Hubs de Eventos do Azure, a mensagem é chamada como event.

Mapeamento entre as propriedades de mensagem e de evento do Event Hubs e os cabeçalhos de mensagem do Spring no modo de ouvinte de registro:

Propriedades do Evento nos Hubs de Eventos Constantes de Cabeçalho de Mensagem do Spring Tipo Descrição
Tempo enfileirado EventHubsHeaders#ENQUEUED_TIME Instante O instante, em UTC, em que o evento foi enfileirado na partição do Event Hub.
Offset EventHubsHeaders#OFFSET long O deslocamento do evento quando ele foi recebido da partição do Hub de Eventos associada.
Chave de partição AzureHeaders#PARTITION_KEY Corda A chave de hash da partição, caso tenha sido definida quando o evento foi publicado originalmente.
ID da partição AzureHeaders#RAW_PARTITION_ID Corda A ID de partição do Hub de Eventos.
Número da sequência EventHubsHeaders#SEQUENCE_NUMBER Longo O número de sequência atribuído ao evento quando ele foi enfileirado na partição do Hub de Eventos associada.
Propriedades do último evento enfileirado EventHubsHeaders#LAST_ENQUEUED_EVENT_PROPERTIES LastEnqueuedEventProperties As propriedades do último evento enfileirado nesta partição.
NA AzureHeaders#CHECKPOINTER Ponto de verificação O cabeçalho para marcar a mensagem específica.

Os usuários podem analisar os cabeçalhos da mensagem para obter as informações relacionadas de cada evento. Para definir um cabeçalho de mensagem para o evento, todos os cabeçalhos personalizados serão colocados como uma propriedade de aplicativo de um evento, em que o cabeçalho é definido como a chave de propriedade. Quando os eventos forem recebidos dos Hubs de Eventos, todas as propriedades do aplicativo serão convertidas no cabeçalho da mensagem.

Nota

Os cabeçalhos de mensagem da chave de partição, tempo enfileirado, deslocamento e número de sequência não são compatíveis para serem definidos manualmente.

Quando o modo de consumidor em lote está habilitado, os cabeçalhos específicos das mensagens em lote são listados a seguir e contêm uma lista de valores de cada evento individual do Event Hubs.

Mapeamento entre as propriedades de mensagens/eventos dos Hubs de Eventos e os cabeçalhos de mensagens do Spring em modo de ouvinte em lote:

Propriedades do Evento nos Hubs de Eventos Constantes de cabeçalho da mensagem do Spring Batch Tipo Descrição
Tempo enfileirado EventHubsHeaders#ENQUEUED_TIME Lista de instantâneos Lista dos instantes, em UTC, em que cada evento foi enfileirado na partição do Event Hub.
Offset EventHubsHeaders#OFFSET Lista de longos Lista do deslocamento de cada evento quando ele foi recebido da partição do Hub de Eventos associada.
Chave de partição AzureHeaders#PARTITION_KEY Lista de strings Lista das chaves de hash da partição, caso tenham sido definidas quando cada evento foi publicado originalmente.
Número da sequência EventHubsHeaders#SEQUENCE_NUMBER Lista de longos Lista dos números de sequência atribuídos a cada evento quando ele foi enfileirado na partição associada do Hub de Eventos.
Propriedades do sistema EventHubsHeaders#BATCH_CONVERTED_SYSTEM_PROPERTIES Lista de Mapas Lista das propriedades do sistema de cada evento.
Propriedades do aplicativo EventHubsHeaders#BATCH_CONVERTED_APPLICATION_PROPERTIES Lista de Mapas Lista das propriedades da aplicação de cada evento, onde são colocados todos os cabeçalhos de mensagem personalizados ou as propriedades do evento.

Nota

Ao publicar mensagens, todos os cabeçalhos de lote acima serão removidos das mensagens, se existirem.

Amostras

Para obter mais informações, consulte o azure-spring-boot-samples repositório no GitHub.

Integração do Spring com Barramento de Serviço do Azure

Principais conceitos

O Spring Integration permite mensagens leves em aplicativos baseados em Spring e dá suporte à integração com sistemas externos por meio de adaptadores declarativos.

O projeto de extensão Spring Integration para Barramento de Serviço do Azure fornece adaptadores de canal de entrada e de saída para o Barramento de Serviço do Azure.

Nota

As APIs de suporte do CompletableFuture foram preteridas da versão 2.10.0 e substituídas pelo Reactor Core da versão 4.0.0. Consulte Javadoc para obter detalhes.

Configuração de dependência

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

Configuração

Este início apresenta as 2 partes a seguir das opções de configuração:

Propriedades de configuração de conexão

Esta seção contém as opções de configuração usadas para conectar-se ao Barramento de Serviço do Azure.

Nota

Se você optar por usar uma entidade de segurança para autenticar e autorizar com a ID do Microsoft Entra para acessar um recurso do Azure, consulte Autorizar o acesso com a ID do Microsoft Entra para garantir que a entidade de segurança tenha recebido a permissão suficiente para acessar o recurso do Azure.

Propriedades configuráveis de conexão de spring-cloud-azure-starter-integration-servicebus:

Propriedade Tipo Descrição
spring.cloud.azure.servicebus.enabled booleano Indica se o Barramento de Serviço do Azure está habilitado.
spring.cloud.azure.servicebus.connection-string Corda Valor da string de conexão do namespace do Barramento de Serviço.
spring.cloud.azure.servicebus.custom-endpoint-address Corda O endereço de endpoint personalizado usado ao se conectar ao Barramento de Serviço.
spring.cloud.azure.servicebus.namespace Corda Valor do namespace do Barramento de Serviço, que é o prefixo do FQDN. Um FQDN deve ser composto por NamespaceName.DomainName
spring.cloud.azure.servicebus.domain-name Corda Valor do nome de domínio de um namespace do Barramento de Serviço do Azure.

Propriedades de configuração do processador do Barramento de Serviço

O ServiceBusInboundChannelAdapter usa o ServiceBusProcessorClient para consumir mensagens, para configurar as propriedades gerais de um ServiceBusProcessorClient, os desenvolvedores podem usar ServiceBusContainerProperties para a configuração. Consulte a seção a seguir sobre como trabalhar com ServiceBusInboundChannelAdapter.

Uso básico

Enviar mensagens para o Barramento de Serviço do Azure

  1. Preencha as opções de configuração de credencial.

    • Para credenciais como cadeia de conexão, configure as seguintes propriedades no arquivo application.yml:

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

      Nota

      A Microsoft recomenda usar o fluxo de autenticação mais seguro disponível. O fluxo de autenticação descrito nesse procedimento, como para bancos de dados, caches, mensagens ou serviços de IA, exige um grau muito alto de confiança no aplicativo e traz riscos não presentes em outros fluxos. Use esse fluxo somente quando opções mais seguras, como identidades gerenciadas para conexões sem senha ou sem chave, não forem viáveis. Para operações de máquinas locais, prefira identidades de usuário para conexões sem senha ou sem chave.

    • Para credenciais como identidades gerenciadas, configure as seguintes propriedades em seu arquivo 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}
      

Nota

Os valores permitidos para tenant-id são: common, organizations, consumersou a ID do locatário. Para obter mais informações sobre esses valores, consulte a seção Uso do ponto de extremidade incorreto (contas pessoais e de organização) do Error AADSTS50020 – A conta de usuário do provedor de identidade não existe no locatário. Para obter informações sobre como converter seu aplicativo de locatário único, consulte Converter aplicativo de locatário único em multilocatário no Microsoft Entra ID.

  • Para credenciais de service principal, configure as seguintes propriedades no arquivo 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}
    

Nota

Os valores permitidos para tenant-id são: common, organizations, consumersou a ID do locatário. Para obter mais informações sobre esses valores, consulte a seção Uso do ponto de extremidade incorreto (contas pessoais e de organização) do Error AADSTS50020 – A conta de usuário do provedor de identidade não existe no locatário. Para obter informações sobre como converter seu aplicativo de locatário único, consulte Converter aplicativo de locatário único em multilocatário no Microsoft Entra ID.

  1. Crie DefaultMessageHandler usando o bean ServiceBusTemplate para enviar mensagens ao Barramento de Serviço e defina o tipo de entidade no ServiceBusTemplate. Este exemplo usa a Fila do Barramento de Serviço.

    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. Crie um vínculo de gateway de mensagens com o manipulador de mensagens acima por meio de um canal de mensagens.

    class Demo {
        @Autowired
        QueueOutboundGateway messagingGateway;
    
        @MessagingGateway(defaultRequestChannel = OUTPUT_CHANNEL)
        public interface QueueOutboundGateway {
            void send(String text);
        }
    }
    
  3. Enviar mensagens usando o gateway.

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

Receber mensagens do Barramento de Serviço do Azure

  1. Preencha as opções de configuração de credencial.

  2. Crie um bean de canal de mensagens que seja o canal de entrada.

    @Configuration
    class Demo {
        private static final String INPUT_CHANNEL = "input";
    
        @Bean
        public MessageChannel input() {
            return new DirectChannel();
        }
    }
    
  3. Crie ServiceBusInboundChannelAdapter com o bean ServiceBusMessageListenerContainer para receber mensagens do Barramento de Serviço. Este exemplo usa a fila do Barramento de Serviço.

    @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. Crie uma associação de receptor de mensagem com ServiceBusInboundChannelAdapter por meio do canal de mensagem que criamos antes.

    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();
        }
    }
    

Configurar ServiceBusMessageConverter para personalizar objectMapper

ServiceBusMessageConverter é feito como um bean configurável para permitir que os usuários personalizem ObjectMapper.

Cabeçalhos de mensagem do Barramento de Serviço

Para alguns cabeçalhos do Barramento de Serviço que podem ser mapeados para mais de uma constante de cabeçalho do Spring, a prioridade dos diferentes cabeçalhos do Spring é indicada.

Mapeamento entre cabeçalhos do Barramento de Serviço e cabeçalhos do Spring:

Cabeçalhos e propriedades da mensagem do Barramento de Serviço Constantes de cabeçalho de mensagem do Spring Tipo Configurável Descrição
Tipo de conteúdo MessageHeaders#CONTENT_TYPE Corda Sim O descritor Content-Type RFC2045 da mensagem.
ID de correlação ServiceBusMessageHeaders#CORRELATION_ID Corda Sim A ID de correlação da mensagem
ID da mensagem ServiceBusMessageHeaders#MESSAGE_ID Corda Sim A ID da mensagem, esse cabeçalho tem prioridade maior do que MessageHeaders#ID.
ID da mensagem MessageHeaders#ID Identificador Único Universal (UUID) Sim O ID da mensagem; este cabeçalho tem prioridade inferior a ServiceBusMessageHeaders#MESSAGE_ID.
Chave de partição ServiceBusMessageHeaders#PARTITION_KEY Corda Sim A chave de partição para enviar a mensagem para uma entidade particionada.
Responder a MessageHeaders#REPLY_CHANNEL Corda Sim O endereço de uma entidade para a qual enviar respostas.
Responder à ID da sessão ServiceBusMessageHeaders#REPLY_TO_SESSION_ID Corda Sim O valor da propriedade ReplyToGroupId da mensagem.
Hora de enfileiramento agendada utc ServiceBusMessageHeaders#SCHEDULED_ENQUEUE_TIME OffsetDateTime Sim A data e hora em que a mensagem deve ser enfileirada no Barramento de Serviço; este cabeçalho tem prioridade sobre AzureHeaders#SCHEDULED_ENQUEUE_MESSAGE.
Hora de enfileiramento agendada utc AzureHeaders#SCHEDULED_ENQUEUE_MESSAGE Inteiro Sim A data e hora na qual a mensagem deve ser enfileirada no Barramento de Serviço; este cabeçalho tem prioridade inferior à de ServiceBusMessageHeaders#SCHEDULED_ENQUEUE_TIME.
ID da sessão ServiceBusMessageHeaders#SESSION_ID Corda Sim O identificador da sessão de uma entidade ciente da sessão.
Vida útil ServiceBusMessageHeaders#TIME_TO_LIVE Duração Sim A duração do tempo antes que essa mensagem expire.
Para ServiceBusMessageHeaders#TO Corda Sim O endereço de destino ("to") da mensagem, reservado para uso futuro em cenários de roteamento e atualmente ignorado pelo próprio broker.
Assunto ServiceBusMessageHeaders#SUBJECT Corda Sim O assunto da mensagem.
Descrição de erro da mensagem morta ServiceBusMessageHeaders#DEAD_LETTER_ERROR_DESCRIPTION Corda Não A descrição de uma mensagem que foi enviada para a fila de mensagens mortas.
Motivo da mensagem morta ServiceBusMessageHeaders#DEAD_LETTER_REASON Corda Não O motivo pelo qual uma mensagem foi encaminhada para a fila de mensagens mortas.
Fonte da fila de mensagens mortas ServiceBusMessageHeaders#DEAD_LETTER_SOURCE Corda Não A entidade na qual a mensagem foi enviada como morta.
Número de entregas ServiceBusMessageHeaders#DELIVERY_COUNT long Não O número de vezes que essa mensagem foi entregue aos clientes.
Número da sequência enfileirada ServiceBusMessageHeaders#ENQUEUED_SEQUENCE_NUMBER longo Não O número de sequência atribuído a uma mensagem quando ela é enfileirada pelo Barramento de Serviço.
Tempo enfileirado ServiceBusMessageHeaders#ENQUEUED_TIME OffsetDateTime Não A data e hora em que esta mensagem foi enfileirada no Barramento de Serviço.
Expira em ServiceBusMessageHeaders#EXPIRES_AT OffsetDateTime Não A data e hora em que esta mensagem expirará.
Token de bloqueio ServiceBusMessageHeaders#LOCK_TOKEN Corda Não O token de bloqueio da mensagem atual.
Bloqueado até ServiceBusMessageHeaders#LOCKED_UNTIL OffsetDateTime Não A data e hora em que o bloqueio desta mensagem expira.
Número da sequência ServiceBusMessageHeaders#SEQUENCE_NUMBER Longo Não O número exclusivo atribuído a uma mensagem pelo Barramento de Serviço.
Estado ServiceBusMessageHeaders#STATE ServiceBusMessageState Não O estado da mensagem, que pode ser Ativa, Adiada ou Agendada.

Suporte à chave de partição

Este starter oferece suporte ao particionamento do Barramento de Serviço, permitindo definir a chave de partição e o ID de sessão no cabeçalho da mensagem. Esta seção apresenta como definir a chave de partição para mensagens.

Recomendado: use ServiceBusMessageHeaders.PARTITION_KEY como a chave do cabeçalho.

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!");
    }
}

Não recomendado, mas com suporte no momento: AzureHeaders.PARTITION_KEY como a chave do cabeçalho.

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!");
    }
}

Nota

Quando ServiceBusMessageHeaders.PARTITION_KEY e AzureHeaders.PARTITION_KEY são definidos nos cabeçalhos da mensagem, ServiceBusMessageHeaders.PARTITION_KEY é preferencial.

Suporte à sessão

Este exemplo demonstra como definir manualmente a ID da sessão de uma mensagem no aplicativo.

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!");
    }
}

Nota

Quando o ServiceBusMessageHeaders.SESSION_ID é definido nos cabeçalhos da mensagem e um cabeçalho ServiceBusMessageHeaders.PARTITION_KEY diferente também é definido, o valor da ID da sessão será eventualmente usado para substituir o valor da chave de partição.

Personalizar propriedades do cliente do Barramento de Serviço

Os desenvolvedores podem usar AzureServiceClientBuilderCustomizer para personalizar as propriedades do cliente do Barramento de Serviço. O exemplo a seguir personaliza a propriedade sessionIdleTimeout em ServiceBusClientBuilder:

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

Amostras

Para obter mais informações, consulte o azure-spring-boot-samples repositório no GitHub.

Integração do Spring com a Fila de Armazenamento do Azure

Principais conceitos

O Armazenamento de Filas do Azure é um serviço para armazenar um grande número de mensagens. Você acessa mensagens de qualquer lugar do mundo por meio de chamadas autenticadas usando HTTP ou HTTPS. Uma mensagem de fila pode ter até 64 KB de tamanho. Uma fila pode conter milhões de mensagens, até o limite total de capacidade de uma conta de armazenamento. As filas são comumente usadas para criar uma fila de trabalho pendente para processamento assíncrono.

Configuração de dependência

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

Configuração

Este modelo inicial fornece as seguintes opções de configuração:

Propriedades de configuração de conexão

Esta seção contém as opções de configuração usadas para se conectar à Fila de Armazenamento do Azure.

Nota

Se você optar por usar uma entidade de segurança para autenticar e autorizar com a ID do Microsoft Entra para acessar um recurso do Azure, consulte Autorizar o acesso com a ID do Microsoft Entra para garantir que a entidade de segurança tenha recebido a permissão suficiente para acessar o recurso do Azure.

Propriedades configuráveis de conexão de spring-cloud-azure-starter-integration-storage-queue:

Propriedade Tipo Descrição
spring.cloud.azure.storage.queue.enabled booleano Indica se uma fila de armazenamento do Azure está habilitada.
spring.cloud.azure.storage.queue.connection-string Corda Valor da string de conexão do namespace da fila de armazenamento.
spring.cloud.azure.storage.queue.accountName Corda Nome da conta da Fila de Armazenamento.
spring.cloud.azure.storage.queue.accountKey Corda Chave da conta da Fila de Armazenamento.
spring.cloud.azure.storage.queue.endpoint Corda Ponto de extremidade de serviço Fila de Armazenamento.
spring.cloud.azure.storage.queue.sasToken Corda Credencial de token SAS
spring.cloud.azure.storage.queue.serviceVersion QueueServiceVersion QueueServiceVersion que é usado ao fazer solicitações de API.
spring.cloud.azure.storage.queue.messageEncoding Corda Codificação de mensagens de fila.

Uso básico

Enviar mensagens para a Fila de Armazenamento do Azure

  1. Preencha as opções de configuração de credencial.

    • Para credenciais como cadeia de conexão, configure as seguintes propriedades no arquivo application.yml:

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

      Nota

      A Microsoft recomenda usar o fluxo de autenticação mais seguro disponível. O fluxo de autenticação descrito nesse procedimento, como para bancos de dados, caches, mensagens ou serviços de IA, exige um grau muito alto de confiança no aplicativo e traz riscos não presentes em outros fluxos. Use esse fluxo somente quando opções mais seguras, como identidades gerenciadas para conexões sem senha ou sem chave, não forem viáveis. Para operações de máquinas locais, prefira identidades de usuário para conexões sem senha ou sem chave.

    • Para credenciais como identidades gerenciadas, configure as seguintes propriedades em seu arquivo 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}
      

Nota

Os valores permitidos para tenant-id são: common, organizations, consumersou a ID do locatário. Para obter mais informações sobre esses valores, consulte a seção Uso do ponto de extremidade incorreto (contas pessoais e de organização) do Error AADSTS50020 – A conta de usuário do provedor de identidade não existe no locatário. Para obter informações sobre como converter seu aplicativo de locatário único, consulte Converter aplicativo de locatário único em multilocatário no Microsoft Entra ID.

  • Para credenciais de service principal, configure as seguintes propriedades no arquivo 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}
    

Nota

Os valores permitidos para tenant-id são: common, organizations, consumersou a ID do locatário. Para obter mais informações sobre esses valores, consulte a seção Uso do ponto de extremidade incorreto (contas pessoais e de organização) do Error AADSTS50020 – A conta de usuário do provedor de identidade não existe no locatário. Para obter informações sobre como converter seu aplicativo de locatário único, consulte Converter aplicativo de locatário único em multilocatário no Microsoft Entra ID.

  1. Crie DefaultMessageHandler com o bean StorageQueueTemplate para enviar mensagens para o 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. Crie uma vinculação de gateway de mensagens com o manipulador de mensagens acima por meio de um canal de mensagens.

    class Demo {
        @Autowired
        StorageQueueOutboundGateway storageQueueOutboundGateway;
    
        @MessagingGateway(defaultRequestChannel = OUTPUT_CHANNEL)
        public interface StorageQueueOutboundGateway {
            void send(String text);
        }
    }
    
  3. Enviar mensagens usando o gateway.

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

Receber mensagens da Fila de Armazenamento do Azure

  1. Preencha as opções de configuração de credencial.

  2. Crie um bean de canal de mensagens que seja o canal de entrada.

    class Demo {
        private static final String INPUT_CHANNEL = "input";
    
        @Bean
        public MessageChannel input() {
            return new DirectChannel();
        }
    }
    
  3. Crie StorageQueueMessageSource usando o bean StorageQueueTemplate para receber mensagens da 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. Crie uma associação de receptor de mensagem com StorageQueueMessageSource criado na última etapa por meio do canal de mensagem que criamos antes.

    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();
        }
    }
    

Amostras

Para obter mais informações, consulte o azure-spring-boot-samples repositório no GitHub.