Suporte do Spring Cloud Azure para o Spring Cloud Stream

O Spring Cloud Stream é uma estrutura para a criação de microsserviços altamente escaláveis orientados a eventos conectados com sistemas de mensagens compartilhados.

A estrutura fornece um modelo de programação flexível baseado em expressões idiomáticas e melhores práticas Spring já estabelecidas e familiares. Estas práticas recomendadas incluem suporte para semânticas pub/sub persistentes, grupos de consumidores e partições com estado.

As implementações atuais do ligador incluem:

Spring Cloud Stream Binder para Hubs de Eventos do Azure

Conceitos-chave

O Spring Cloud Stream Binder para Hubs de Eventos do Azure fornece a implementação de vinculação para a estrutura do Spring Cloud Stream. Esta implementação usa os adaptadores de canal do Spring Integration Event Hubs em sua base. Do ponto de vista do design, os Hubs de Eventos são semelhantes a Kafka. Além disso, os Hubs de Eventos podem ser acessados via Kafka API. Se o seu projeto tiver dependência total da API Kafka, você pode tentar Hub de Eventos com a API Kafka Sample

Grupo de consumidores

Os Hubs de Eventos fornecem suporte semelhante ao grupo de consumidores do Apache Kafka, mas com uma lógica ligeiramente diferente. Enquanto o Kafka armazena todos os offsets confirmados no broker, tem de armazenar manualmente os offsets das mensagens do Event Hubs que estão a ser processadas. O SDK do Event Hubs disponibiliza a funcionalidade de armazenar esses offsets no Armazenamento do Azure.

Suporte de particionamento

Os Hubs de Eventos fornecem um conceito de partição física semelhante ao Kafka. Mas, ao contrário do reequilíbrio automático do Kafka entre consumidores e partições, o Event Hubs disponibiliza um tipo de modo preemptivo. A conta de armazenamento atua como uma locação para determinar qual consumidor possui qual partição. Quando um novo consumidor começa, ele tenta roubar algumas partições dos consumidores mais carregados para alcançar o equilíbrio da carga de trabalho.

Para especificar a estratégia de balanceamento de carga, as propriedades de spring.cloud.stream.eventhubs.bindings.<binding-name>.consumer.load-balancing.* são fornecidas. Para obter mais informações, consulte a seção propriedades do consumidor .

Suporte ao consumidor em lote

O binder do Spring Cloud Azure Stream Event Hubs suporta a funcionalidade Spring Cloud Stream Batch Consumer.

Para trabalhar com o modo batch-consumer, defina a propriedade spring.cloud.stream.bindings.<binding-name>.consumer.batch-mode como true. Quando habilitada, uma mensagem com uma carga de uma lista de eventos em lote é recebida e passada para a função Consumer. Cada cabeçalho de mensagem também é convertido em uma lista, cujo conteúdo é o valor de cabeçalho associado analisado de cada evento. Os cabeçalhos comuns de ID de partição, ponteiro de verificação e últimas propriedades enfileiradas são apresentados como um único valor porque todo o lote de eventos compartilha o mesmo valor. Para obter mais informações, consulte a seção cabeçalhos de mensagem dos Hubs de Eventos de suporte do Spring Cloud Azure para o Spring Integration.

Observação

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

O ponto de verificação do consumidor em lote suporta dois modos: BATCH e MANUAL. O modo BATCH é um modo de criação automática de pontos de verificação para criar um ponto de verificação de todo o lote de eventos assim que o binder os recebe. MANUAL modo serve para registar pontos de controlo dos eventos dos utilizadores. Quando utilizado, o Checkpointer é incluído no cabeçalho da mensagem, e os utilizadores podem usá-lo para criar pontos de verificação.

Você pode especificar o tamanho do lote definindo as propriedades max-size e max-wait-time que têm um prefixo de spring.cloud.stream.eventhubs.bindings.<binding-name>.consumer.batch.. A propriedade max-size é necessária e a propriedade max-wait-time é opcional. Para obter mais informações, consulte a seção propriedades do consumidor .

Configuração de dependência

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

Como alternativa, você também pode usar o Spring Cloud Azure Stream Event Hubs Starter, conforme mostrado no exemplo a seguir para o Maven:

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

Configuração

O fichário fornece as seguintes três partes 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.

Observação

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 acesso com o Microsoft Entra ID para verificar se a entidade de segurança recebeu a permissão suficiente para acessar o recurso do Azure.

Propriedades configuráveis de ligação de spring-cloud-azure-stream-binder-eventhubs:

Propriedade Tipo Descrição
spring.cloud.azure.eventhubs.enabled Booleano Se um Hubs de Eventos do Azure está habilitado.
spring.cloud.azure.eventhubs.connection-string Cordão Valor da cadeia de ligação do Namespace do Event Hubs.
spring.cloud.azure.eventhubs.namespace Cordão Valor do espaço de nomes do Event Hubs, que é o prefixo do FQDN. Um FQDN deve ser composto por NamespaceName.DomainName
spring.cloud.azure.eventhubs.domain-name Cordão Valor do nome de domínio de um espaço de nomes do Hubs de Eventos do Azure.
spring.cloud.azure.eventhubs.custom-endpoint-address Cordão Endereço de ponto final personalizado.

Dica

As opções comuns de configuração do SDK do Serviço do Azure também são configuráveis para o fichário dos Hubs de Eventos do Azure Stream do Spring Cloud. As opções de configuração suportadas são apresentadas em Spring Cloud Azure configuration e podem ser configuradas com o prefixo unificado spring.cloud.azure. ou com o prefixo de spring.cloud.azure.eventhubs..

O binder também suporta Spring Could Azure Resource Manager por predefinição. Para saber como recuperar a cadeia de ligação com principais de segurança aos quais não foram atribuídas as funções relacionadas com Data, consulte a secção Utilização básica de Spring Could Azure Resource Manager.

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

Esta secção contém as opções de configuração do serviço Storage Blobs, utilizado para persistir as informações sobre a propriedade das partições e os pontos de controlo.

Observação

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 está ativada manualmente, nenhum contentor de armazenamento será criado automaticamente com o nome de spring.cloud.stream.bindings.binding-name.destination.

Pontos de verificação de propriedades configuráveis de spring-cloud-azure-stream-binder-eventhubs:

Propriedade Tipo Descrição
spring.cloud.azure.eventhubs.processor.checkpoint-store.create-container-if-not-exists booleano Se deve ser permitida a criação de contentores caso não existam.
spring.cloud.azure.eventhubs.processor.checkpoint-store.account-name Cordão Nome da conta de armazenamento.
spring.cloud.azure.eventhubs.processor.checkpoint-store.account-key Cordão Chave de acesso da conta de armazenamento.
spring.cloud.azure.eventhubs.processor.checkpoint-store.container-name Cordão Nome do recipiente de armazenamento.

Dica

As opções comuns de configuração do SDK de serviço do Azure também podem ser configuradas para o arquivo de pontos de verificação do Armazenamento de Blobs. As opções de configuração suportadas são apresentadas em Spring Cloud Azure configuration e podem ser configuradas com o prefixo unificado spring.cloud.azure. ou com o prefixo de spring.cloud.azure.eventhubs.processor.checkpoint-store.

Propriedades de configuração de vinculação dos Hubs de Eventos do Azure

As opções a seguir estão divididas em quatro seções: Propriedades do consumidor, Configurações avançadas do consumidor, Propriedades do produtor e Configurações avançadas do produtor.

Propriedades do consumidor

Estas propriedades são expostas através de EventHubsConsumerProperties.

Observação

Para evitar repetições, desde as versões 4.17.0 e 5.11.0, o Spring Cloud Azure Stream Binder Event Hubs oferece suporte à configuração de valores para todos os canais, no formato de spring.cloud.stream.eventhubs.default.consumer.<property>=<value>.

Propriedades configuráveis pelo consumidor de spring-cloud-azure-stream-binder-eventhubs:

Propriedade Tipo Descrição
spring.cloud.stream.eventhubs.bindings.binding-name.consumer.checkpoint.mode CheckpointMode Modo de ponto de verificação usado quando o consumidor decide como enviar a mensagem de ponto de verificação
spring.cloud.stream.eventhubs.bindings.binding-name.consumer.checkpoint.count Inteiro Decide a quantidade de mensagem para cada partição para fazer um ponto de verificação. Entrará em vigor somente quando PARTITION_COUNT modo de ponto de verificação for usado.
spring.cloud.stream.eventhubs.bindings.binding-name.consumer.checkpoint.interval Duração Decide o intervalo de tempo para fazer um ponto de verificação. Entrará em vigor somente quando TIME modo de ponto de verificação for usado.
spring.cloud.stream.eventhubs.bindings.<binding-name>.consumer.batch.max-size Inteiro O número máximo de eventos em um lote. Necessário para o modo batch-consumer.
spring.cloud.stream.eventhubs.bindings.binding-name.consumer.batch.max-wait-time Duração A duração máxima do consumo em lote. Entrará em vigor somente quando o modo batch-consumer estiver ativado e for opcional.
spring.cloud.stream.eventhubs.bindings.binding-name.consumer.load-balancing.update-interval Duração A duração do intervalo de tempo para atualização.
spring.cloud.stream.eventhubs.bindings.binding-name.consumer.load-balancing.strategy LoadBalancingStrategy A estratégia de balanceamento de carga.
spring.cloud.stream.eventhubs.bindings.binding-name.consumer.load-balancing.partition-ownership-expiration-interval Duração O período de tempo ao fim do qual a propriedade da partição expira.
spring.cloud.stream.eventhubs.bindings.binding-name.consumer.track-last-enqueued-event-properties booleano Se o processador de eventos deve solicitar informações sobre o último evento colocado na fila na partição associada e acompanhar essas informações à medida que os eventos são recebidos.
spring.cloud.stream.eventhubs.bindings.binding-name.consumer.prefetch-count Inteiro O número utilizado pelo consumidor para controlar o número de eventos que o consumidor do Event Hub receberá ativamente e colocará em fila localmente.
spring.cloud.stream.eventhubs.bindings.binding-name.consumer.initial-partition-event-position Mapa com a chave como o ID da partição e os valores de StartPositionProperties O mapa que contém a posição do evento a ser usada para cada partição se um ponto de verificação para a partição não existir no armazenamento de pontos de verificação. Este mapa baseia-se no ID da partição.

Observação

A configuração initial-partition-event-position aceita um map para especificar a posição inicial para cada hub de eventos. Assim, sua chave é o ID da partição, e o valor é de StartPositionProperties, que inclui propriedades de deslocamento, número de sequência, data e hora enfileirada e se inclusive. Por exemplo, você pode defini-lo como

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
Configuração avançada do cliente

A configuração de ligação, ponto de verificação e cliente comum do SDK do Azure acima suporta personalização para cada consumidor de associador, que pode configurar-se com o prefixo .

Propriedades do produtor

Estas propriedades são expostas através de EventHubsProducerProperties.

Observação

Para evitar repetições, desde as versões 4.17.0 e 5.11.0, o Spring Cloud Azure Stream Binder Event Hubs oferece suporte à configuração de valores para todos os canais, no formato de spring.cloud.stream.eventhubs.default.producer.<property>=<value>.

Propriedades configuráveis pelo produtor de spring-cloud-azure-stream-binder-eventhubs:

Propriedade Tipo Descrição
spring.cloud.stream.eventhubs.bindings.binding-name.producer.sync Booleano O indicador de comutação para sincronização do produtor. Se verdadeiro, o produtor aguardará uma resposta após uma operação de envio.
spring.cloud.stream.eventhubs.bindings.binding-name.producer.send-timeout longo A quantidade de tempo para aguardar uma resposta após uma operação de envio. Terá efeito somente quando um produtor de sincronização estiver habilitado.
Configuração avançada do produtor

A ligação acima e a configuração do cliente comum do SDK do Azure suportam a personalização para cada produtor de binder, que pode configurar com o prefixo .

Utilização básica

Enviar e receber mensagens de/para Hubs de Eventos

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

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

      Observação

      A Microsoft recomenda o uso do fluxo de autenticação mais seguro disponível. O fluxo de autenticação descrito neste procedimento, como para bancos de dados, caches, mensagens ou serviços de IA, requer 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áquina local, prefira identidades de usuário para conexões sem senha ou sem chave.

    • Para credenciais de principal de serviço, configure as seguintes propriedades no ficheiro 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
      

Observação

Os valores permitidos para tenant-id são: common, organizations, consumersou o ID do locatário. Para obter mais informações sobre estes valores, consulte a secção Utilizou o ponto final errado (contas pessoais e organizacionais) de Erro AADSTS50020 - A conta de utilizador do fornecedor de identidades não existe no inquilino. 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 como identidades gerenciadas, configure as seguintes propriedades em seu arquivo 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. Definir fornecedor e consumidor.

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

Suporte de particionamento

Um PartitionSupplier com informações de partição fornecidas pelo usuário é criado para configurar as informações de partição sobre a mensagem a ser enviada. O fluxograma a seguir mostra o processo de obtenção de diferentes prioridades para o ID e a chave da partição:

Diagrama mostrando um fluxograma do processo de suporte ao particionamento.

Suporte ao consumidor em lote

  1. Forneça as opções de configuração de lote, conforme mostrado no exemplo a seguir:

    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. Definir fornecedor e consumidor.

    Para o modo de ponto de verificação como BATCH, você pode usar o código a seguir para enviar mensagens e consumir em lotes.

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

    Para o modo de checkpointing definido como MANUAL, pode usar o código seguinte para enviar mensagens e consumir/fazer checkpoint em lotes.

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

Observação

No modo de consumo em lote, o tipo de conteúdo padrão do fichário do Spring Cloud Stream é application/json, portanto, verifique se a carga útil da mensagem está alinhada com o tipo de conteúdo. Por exemplo, ao usar o tipo de conteúdo predefinido de application/json para receber mensagens com carga útil String, a carga útil deve estar em JSON String, entre aspas duplas para o texto String original. Já no caso do tipo de conteúdo text/plain, pode tratar-se diretamente de um objeto String. Para obter mais informações, consulte Spring Cloud Stream Content Type Negotiation.

Tratar mensagens de erro

  • Tratar mensagens de erro de associação de saída

    Por padrão, o Spring Integration cria um canal de erro global chamado errorChannel. Configure o seguinte ponto final de mensagem para processar mensagens de erro de associação de saída.

    @ServiceActivator(inputChannel = IntegrationContextUtils.ERROR_CHANNEL_BEAN_NAME)
    public void handleError(ErrorMessage message) {
        LOGGER.error("Handling outbound binding error: " + message);
    }
    
  • Processar mensagens de erro de associação de entrada

    O Spring Cloud Stream Event Hubs Binder suporta uma solução para lidar com erros para as ligações de mensagens de entrada: manipuladores de erros.

    Manipulador de erros:

    O Spring Cloud Stream expõe um mecanismo para você fornecer um manipulador de erros personalizado adicionando um Consumer que aceita instâncias ErrorMessage. Para obter mais informações, consulte Manipular mensagens de erro na documentação do Spring Cloud Stream.

    • Manipulador de erros predefinido de associação

      Configure um bean Consumer único para consumir todas as mensagens de erro de associação recebidas. A seguinte função padrão se inscreve em cada canal de erro de vinculação de entrada:

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

      Você também precisa definir a propriedade spring.cloud.stream.default.error-handler-definition para o nome da função.

    • Manipulador de erros específico da vinculação

      Configure um Consumer bean para consumir as mensagens de erro de vinculação de entrada específicas. A função seguinte subscreve o canal de erros específico da associação de entrada e tem uma prioridade mais elevada do que o processador de erros predefinido da associação:

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

      Você também precisa definir a propriedade spring.cloud.stream.bindings.<input-binding-name>.error-handler-definition para o nome da função.

Cabeçalhos de mensagem dos Hubs de Eventos

Para obter os cabeçalhos de mensagem básicos suportados, consulte a seção cabeçalhos de mensagem dos Hubs de Eventos de suporte do Spring Cloud Azure para o Spring Integration.

Suporte para vários fichários

A ligação a vários namespaces do Event Hubs também é suportada através da utilização de vários associadores. Este exemplo usa uma cadeia de conexão como exemplo. As credenciais dos principais de serviço e das identidades geridas também são suportadas. Você pode definir propriedades relacionadas nas configurações de ambiente de cada fichário.

  1. Para usar vários binders com o Event Hubs, configure as seguintes propriedades no ficheiro 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
    

    Observação

    O ficheiro da aplicação anterior mostra como configurar um único sondador predefinido para ser aplicado a todas as ligações. Se quiser configurar o poller para uma ligação específica, você pode usar uma configuração como spring.cloud.stream.bindings.<binding-name>.producer.poller.fixed-delay=3000.

    Observação

    A Microsoft recomenda o uso do fluxo de autenticação mais seguro disponível. O fluxo de autenticação descrito neste procedimento, como para bancos de dados, caches, mensagens ou serviços de IA, requer 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áquina local, prefira identidades de usuário para conexões sem senha ou sem chave.

  2. Precisamos definir dois fornecedores e dois consumidores:

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

Provisionamento de recursos

O fichário de Hubs de Eventos oferece suporte ao provisionamento de hub de eventos e grupo de consumidores, os usuários podem usar as seguintes propriedades para habilitar o provisionamento.

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

Observação

Os valores permitidos para tenant-id são: common, organizations, consumersou o ID do locatário. Para obter mais informações sobre estes valores, consulte a secção Utilizou o ponto final errado (contas pessoais e organizacionais) de Erro AADSTS50020 - A conta de utilizador do fornecedor de identidades não existe no inquilino. 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.

Amostras

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

Spring Cloud Stream Binder para o Azure Service Bus

Conceitos-chave

O Spring Cloud Stream Binder for Azure Service Bus fornece a implementação de vinculação para o Spring Cloud Stream Framework. Esta implementação usa adaptadores de canal do Spring Integration Service Bus em sua base.

Mensagem agendada

Este binder permite submeter mensagens a um tópico para processamento diferido. Os usuários podem enviar mensagens agendadas com cabeçalho x-delay expressando em milissegundos um tempo de atraso para a mensagem. A mensagem será entregue aos respetivos tópicos após x-delay milissegundos.

Grupo de consumidores

O Service Bus Topic fornece suporte semelhante ao grupo de consumidores do Apache Kafka, mas com uma lógica ligeiramente diferente. Este aglutinante baseia-se na Subscription de um tópico para atuar como um grupo de consumidores.

Configuração de dependência

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

Como alternativa, você também pode usar o Spring Cloud Azure Stream Service Bus Starter, conforme mostrado no exemplo a seguir para o Maven:

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

Configuração

O fichário fornece as seguintes duas partes das opções de configuração:

Propriedades de configuração de conexão

Esta secção contém as opções de configuração utilizadas para ligar ao Azure Service Bus.

Observação

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 acesso com o Microsoft Entra ID para verificar se a entidade de segurança recebeu a permissão suficiente para acessar o recurso do Azure.

Propriedades configuráveis de ligação de spring-cloud-azure-stream-binder-servicebus:

Propriedade Tipo Descrição
spring.cloud.azure.servicebus.enabled Booleano Se o Azure Service Bus está ativado.
spring.cloud.azure.servicebus.connection-string Cordão Valor da cadeia de ligação do namespace do Service Bus.
spring.cloud.azure.servicebus.custom-endpoint-address Cordão O endereço do ponto final personalizado a utilizar ao estabelecer ligação ao Service Bus.
spring.cloud.azure.servicebus.namespace Cordão Valor do espaço de nomes do Service Bus, que é o prefixo do FQDN. Um FQDN deve ser composto por NamespaceName.DomainName
spring.cloud.azure.servicebus.domain-name Cordão Valor do nome de domínio de um Namespace do Azure Service Bus.

Observação

As opções comuns de configuração do SDK de Serviço do Azure também podem ser configuradas para o binder do Spring Cloud Azure Stream Service Bus. As opções de configuração suportadas são apresentadas em Spring Cloud Azure configuration e podem ser configuradas com o prefixo unificado spring.cloud.azure. ou com o prefixo de spring.cloud.azure.servicebus..

O binder também suporta Spring Could Azure Resource Manager por predefinição. Para saber como recuperar a cadeia de ligação com principais de segurança aos quais não foram atribuídas as funções relacionadas com Data, consulte a secção Utilização básica de Spring Could Azure Resource Manager.

Propriedades de configuração da associação do Azure Service Bus

As opções a seguir estão divididas em quatro seções: Propriedades do consumidor, Configurações avançadas do consumidor, Propriedades do produtor e Configurações avançadas do produtor.

Propriedades do consumidor

Estas propriedades são expostas através de ServiceBusConsumerProperties.

Observação

Para evitar repetições, desde as versões 4.17.0 e 5.11.0, o Spring Cloud Azure Stream Binder Service Bus suporta a definição de valores para todos os canais, no formato de spring.cloud.stream.servicebus.default.consumer.<property>=<value>.

Propriedades configuráveis pelo consumidor de spring-cloud-azure-stream-binder-servicebus:

Propriedade Tipo Predefinido Descrição
spring.cloud.stream.servicebus.bindings.binding-name.consumer.requeue-rejected Booleano falso Se as mensagens com falha forem roteadas para o DLQ.
spring.cloud.stream.servicebus.bindings.binding-name.consumer.max-concurrent-calls Inteiro 1 Número máximo de mensagens em simultâneo que o cliente de processador do Service Bus deve processar. Quando a sessão está ativada, aplica-se a cada sessão.
spring.cloud.stream.servicebus.bindings.binding-name.consumer.max-concurrent-sessions Inteiro null Número máximo de sessões simultâneas a serem processadas a qualquer momento.
spring.cloud.stream.servicebus.bindings.binding-name.consumer.session-enabled booleano null Se a sessão está habilitada.
spring.cloud.stream.servicebus.bindings.binding-name.consumer.session-idle-timeout Duração null Define a quantidade máxima de tempo (Duração) para aguardar o recebimento de uma mensagem para a sessão ativa no momento.
spring.cloud.stream.servicebus.bindings.binding-name.consumer.prefetch-count Inteiro 0 A contagem de pré-carregamento do cliente processador do Service Bus.
spring.cloud.stream.servicebus.bindings.binding-name.consumer.sub-queue Subfila de espera nenhum O tipo de subfila à qual estabelecer ligação.
spring.cloud.stream.servicebus.bindings.binding-name.consumer.max-auto-lock-renew-duration Duração 5 metros O período de tempo durante o qual o bloqueio continuará a ser renovado automaticamente.
spring.cloud.stream.servicebus.bindings.binding-name.consumer.receive-mode ServiceBusReceiveMode peek_lock O modo de receção do cliente processador do Service Bus.
spring.cloud.stream.servicebus.bindings.binding-name.consumer.auto-complete booleano verdadeiro Se as mensagens devem ser confirmadas automaticamente. Se definido como false, um cabeçalho de mensagem de Checkpointer será adicionado para permitir que os desenvolvedores liquidem as mensagens manualmente.
spring.cloud.stream.servicebus.bindings.binding-name.consumer.max-size-in-megabytes Longo 1024 O tamanho máximo da fila/tópico em megabytes, que é o tamanho da memória alocada para a fila/tópico.
spring.cloud.stream.servicebus.bindings.binding-name.consumer.default-message-time-to-live Duração P10675199DT2H48M5.4775807S. (10675199 dias, 2 horas, 48 minutos, 5 segundos e 477 milissegundos) A duração após a qual a mensagem expira, a partir de quando a mensagem é enviada para o Service Bus.

Importante

Ao usar o Azure Resource Manager (ARM), você deve configurar a propriedade spring.cloud.stream.servicebus.bindings.<binding-name>.consume.entity-type. Para mais informações, consulte o servicebus-queue-binder-arm exemplo no GitHub.

Configuração avançada do utilizador

A configuração de ligação acima e a configuração comum do cliente do SDK do Azure suportam a personalização para cada consumidor do binder, que pode configurar com o prefixo .

Propriedades do produtor

Estas propriedades são expostas através de ServiceBusProducerProperties.

Observação

Para evitar repetições, desde as versões 4.17.0 e 5.11.0, o Spring Cloud Azure Stream Binder Service Bus suporta a definição de valores para todos os canais, no formato de spring.cloud.stream.servicebus.default.producer.<property>=<value>.

Propriedades configuráveis pelo produtor de spring-cloud-azure-stream-binder-servicebus:

Propriedade Tipo Predefinido Descrição
spring.cloud.stream.servicebus.bindings.binding-name.producer.sync Booleano falso Sinalizador de switch para sincronização do produtor.
spring.cloud.stream.servicebus.bindings.binding-name.producer.send-timeout longo 10 000 Valor do tempo limite para envio pelo produtor.
spring.cloud.stream.servicebus.bindings.binding-name.producer.entity-type ServiceBusEntityType null Tipo de entidade do Service Bus do produtor, necessário para o produtor de ligação.
spring.cloud.stream.servicebus.bindings.binding-name.producer.max-size-in-megabytes Longo 1024 O tamanho máximo da fila/tópico em megabytes, que é o tamanho da memória alocada para a fila/tópico.
spring.cloud.stream.servicebus.bindings.binding-name.producer.default-message-time-to-live Duração P10675199DT2H48M5.4775807S. (10675199 dias, 2 horas, 48 minutos, 5 segundos e 477 milissegundos) A duração após a qual a mensagem expira, a partir de quando a mensagem é enviada para o Service Bus.

Importante

Ao usar o produtor de ligação, a propriedade de spring.cloud.stream.servicebus.bindings.<binding-name>.producer.entity-type deve ser configurada.

Configuração avançada do produtor

A ligação acima e a configuração do cliente comum do SDK do Azure suportam a personalização para cada produtor de binder, que pode configurar com o prefixo .

Utilização básica

Enviar e receber mensagens de/para o Service Bus

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

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

      Observação

      A Microsoft recomenda o uso do fluxo de autenticação mais seguro disponível. O fluxo de autenticação descrito neste procedimento, como para bancos de dados, caches, mensagens ou serviços de IA, requer 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áquina local, prefira identidades de usuário para conexões sem senha ou sem chave.

    • Para credenciais de principal de serviço, configure as seguintes propriedades no ficheiro 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
      

Observação

Os valores permitidos para tenant-id são: common, organizations, consumersou o ID do locatário. Para obter mais informações sobre estes valores, consulte a secção Utilizou o ponto final errado (contas pessoais e organizacionais) de Erro AADSTS50020 - A conta de utilizador do fornecedor de identidades não existe no inquilino. 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 como identidades gerenciadas, configure as seguintes propriedades em seu arquivo 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. Definir fornecedor e consumidor.

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

Suporte de chave de partição

O binder suporta o particionamento do Service Bus ao permitir 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.

O Spring Cloud Stream fornece uma propriedade de expressão SpEL para a chave de partição spring.cloud.stream.bindings.<binding-name>.producer.partition-key-expression. Por exemplo, definir essa propriedade como "'partitionKey-' + headers[<message-header-key>]" e adicionar um cabeçalho chamado message-header-key. O Spring Cloud Stream usa o valor desse cabeçalho ao avaliar a expressão para atribuir uma chave de partição. O código a seguir fornece um exemplo de produtor:

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

Suporte de sessão

O binder suporta as sessões de mensagens do Service Bus. O ID de sessão de uma mensagem pode ser definido através do cabeçalho da mensagem.

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

Observação

De acordo com o particionamento do Service Bus, o ID da sessão tem maior prioridade do que a chave de partição. Assim, quando os cabeçalhos ServiceBusMessageHeaders#SESSION_ID e ServiceBusMessageHeaders#PARTITION_KEY são definidos, o valor do ID da sessão é eventualmente usado para substituir o valor da chave de partição.

Tratar mensagens de erro

  • Tratar mensagens de erro de associação de saída

    Por padrão, o Spring Integration cria um canal de erro global chamado errorChannel. Configure o ponto de extremidade da seguinte mensagem para manipular a mensagem de erro de vinculação de saída.

    @ServiceActivator(inputChannel = IntegrationContextUtils.ERROR_CHANNEL_BEAN_NAME)
    public void handleError(ErrorMessage message) {
        LOGGER.error("Handling outbound binding error: " + message);
    }
    
  • Processar mensagens de erro de associação de entrada

    O Spring Cloud Stream Service Bus Binder suporta duas soluções para processar erros nas bindings de mensagens de entrada: o processador de erros do binder e os processadores.

    Processador de erros do Binder:

    O processador de erros predefinido do binder processa a ligação de entrada. Use esse manipulador para enviar mensagens com falha para a fila de mensagens mortas quando spring.cloud.stream.servicebus.bindings.<binding-name>.consumer.requeue-rejected estiver habilitado. Caso contrário, as mensagens com falha serão abandonadas. O manipulador de erros do binder é mutuamente exclusivo de outros manipuladores de erros disponibilizados.

    Manipulador de erros:

    O Spring Cloud Stream expõe um mecanismo para você fornecer um manipulador de erros personalizado adicionando um Consumer que aceita instâncias ErrorMessage. Para obter mais informações, consulte Manipular mensagens de erro na documentação do Spring Cloud Stream.

    • Manipulador de erros de associação predefinido

      Configure um único bean de Consumer para consumir todas as mensagens de erro de vinculação de entrada. A seguinte função padrão se inscreve em cada canal de erro de vinculação de entrada:

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

      Você também precisa definir a propriedade spring.cloud.stream.default.error-handler-definition para o nome da função.

    • Manipulador de erros específico da vinculação

      Configure um Consumer bean para consumir as mensagens de erro de vinculação de entrada específicas. A função a seguir se inscreve no canal de erro de vinculação de entrada específico com uma prioridade maior do que o manipulador de erro de vinculação padrão.

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

      Você também precisa definir a propriedade spring.cloud.stream.bindings.<input-binding-name>.error-handler-definition para o nome da função.

Cabeçalhos de mensagens do Service Bus

Para os cabeçalhos básicos de mensagem suportados, consulte a secção cabeçalhos de mensagem do Service Bus de suporte do Spring Cloud Azure para o Spring Integration.

Observação

Ao definir a chave de partição, a prioridade do cabeçalho da mensagem é maior do que a propriedade do Spring Cloud Stream. Portanto, spring.cloud.stream.bindings.<binding-name>.producer.partition-key-expression só entra em vigor quando nenhum dos cabeçalhos ServiceBusMessageHeaders#SESSION_ID e ServiceBusMessageHeaders#PARTITION_KEY estiver configurado.

Suporte para vários fichários

A ligação a vários espaços de nomes do Service Bus também é suportada através da utilização de vários binders. Este exemplo usa a cadeia de conexão como exemplo. Credenciais de entidades de serviço e identidades gerenciadas também são suportadas, os usuários podem definir propriedades relacionadas nas configurações de ambiente de cada fichário.

  1. Para usar vários associadores do ServiceBus, configure as seguintes propriedades no ficheiro 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
    

    Observação

    O ficheiro da aplicação anterior mostra como configurar um único sondador predefinido para ser aplicado a todas as ligações. Se quiser configurar o poller para uma ligação específica, você pode usar uma configuração como spring.cloud.stream.bindings.<binding-name>.producer.poller.fixed-delay=3000.

    Observação

    A Microsoft recomenda o uso do fluxo de autenticação mais seguro disponível. O fluxo de autenticação descrito neste procedimento, como para bancos de dados, caches, mensagens ou serviços de IA, requer 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áquina local, prefira identidades de usuário para conexões sem senha ou sem chave.

  2. precisamos definir dois fornecedores e dois consumidores

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

Provisionamento de recursos

O binder do Service Bus suporta o provisionamento de filas, tópicos e subscrições; os utilizadores podem usar as seguintes propriedades para ativar o provisionamento.

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}

Observação

Os valores permitidos para tenant-id são: common, organizations, consumersou o ID do locatário. Para obter mais informações sobre estes valores, consulte a secção Utilizou o ponto final errado (contas pessoais e organizacionais) de Erro AADSTS50020 - A conta de utilizador do fornecedor de identidades não existe no inquilino. 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.

Personalizar as propriedades do cliente do Service Bus

Os desenvolvedores podem usar AzureServiceClientBuilderCustomizer para personalizar as propriedades do Service Bus Client. O exemplo a seguir personaliza a propriedade sessionIdleTimeout no ServiceBusClientBuilder:

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

Amostras

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