Obsługa Spring Cloud Stream w Spring Cloud Azure

Spring Cloud Stream to struktura umożliwiająca tworzenie wysoce skalowalnych mikrousług opartych na zdarzeniach połączonych z udostępnionymi systemami obsługi komunikatów.

Platforma udostępnia elastyczny model programowania oparty na już ustalonych i znanych idiomach Spring i najlepszych rozwiązaniach. Te najlepsze rozwiązania obejmują obsługę trwałych semantyki pub/podsieci, grup odbiorców i partycji stanowych.

Obecne implementacje mechanizmu wiązania obejmują:

Spring Cloud Stream Binder dla Azure Event Hubs

Kluczowe pojęcia

Binder Spring Cloud Stream dla Azure Event Hubs udostępnia implementację mechanizmu wiązania dla frameworka Spring Cloud Stream. Ta implementacja opiera się na adapterach kanałów Spring Integration Event Hubs. Z perspektywy projektu usługa Event Hubs jest podobna do platformy Kafka. Ponadto dostęp do usługi Event Hubs można uzyskać za pośrednictwem interfejsu API platformy Kafka. Jeśli Twój projekt jest ściśle zależny od interfejsu API Kafka, możesz wypróbować Events Hub with Kafka API Sample

Grupa odbiorców

Usługa Event Hubs zapewnia podobną obsługę grupy odbiorców jako platformę Apache Kafka, ale z niewielką inną logiką. Podczas gdy platforma Kafka przechowuje wszystkie zatwierdzone przesunięcia w brokerze, musisz przechowywać przesunięcia komunikatów usługi Event Hubs przetwarzanych ręcznie. Zestaw SDK usługi Event Hubs udostępnia funkcję do przechowywania takich przesunięć w usłudze Azure Storage.

Obsługa partycjonowania

Usługa Event Hubs oferuje podobną koncepcję fizycznej partycji jak Kafka. Jednak w przeciwieństwie do automatycznego równoważenia platformy Kafka między użytkownikami i partycjami usługa Event Hubs zapewnia rodzaj trybu wyprzedzania. Konto magazynu danych pełni funkcję dzierżawy służącej do określenia, który konsument jest właścicielem której partycji. Po uruchomieniu nowego użytkownika próbuje ukraść niektóre partycje od najbardziej obciążonych użytkowników w celu osiągnięcia równowagi obciążenia.

Aby określić strategię równoważenia obciążenia, udostępniane są właściwości spring.cloud.stream.eventhubs.bindings.<binding-name>.consumer.load-balancing.*. Aby uzyskać więcej informacji, zobacz sekcję Właściwości konsumenta.

Obsługa konsumentów usługi Batch

Binder Spring Cloud Azure Stream Event Hubs obsługuje funkcję Spring Cloud Stream Batch Consumer.

Aby pracować z trybem batch-consumer, ustaw właściwość spring.cloud.stream.bindings.<binding-name>.consumer.batch-mode na wartość true. Po włączeniu odbierany jest komunikat z ładunkiem zawierającym listę zdarzeń wsadowych i przekazywany do funkcji Consumer. Każdy nagłówek wiadomości jest również konwertowany na listę, której zawartość stanowi odpowiadająca mu wartość nagłówka wyodrębniana z każdego zdarzenia. Wspólne pola nagłówka — identyfikator partycji, checkpointer i właściwości ostatniego zdarzenia umieszczonego w kolejce — są przedstawiane jako pojedyncza wartość, ponieważ cała partia zdarzeń ma te same wartości. Aby uzyskać więcej informacji, zobacz sekcję Nagłówki komunikatów usługi Event Hubs w Spring Cloud pomoc techniczna platformy Azure for Spring Integration.

Uwaga

Nagłówek punktu kontrolnego występuje tylko wtedy, gdy używany jest tryb punktu kontrolnego MANUAL.

Punkt kontrolny odbiorcy wsadowego obsługuje dwa tryby: BATCH i MANUAL. Tryb BATCH to tryb automatycznego tworzenia punktów kontrolnych, w którym punkt kontrolny dla całej partii zdarzeń jest tworzony naraz po ich odebraniu przez powiązanie. MANUAL tryb polega na określeniu punktów kontrolnych zdarzeń przez użytkowników. Gdy jest używany, element Checkpointer jest przekazywany do nagłówka wiadomości, a użytkownicy mogą go używać do tworzenia punktów kontrolnych.

Rozmiar partii można określić, ustawiając właściwości max-size i max-wait-time, które mają prefiks spring.cloud.stream.eventhubs.bindings.<binding-name>.consumer.batch.. Właściwość max-size jest wymagana, a właściwość max-wait-time jest opcjonalna. Aby uzyskać więcej informacji, zobacz sekcję Właściwości konsumenta.

Konfiguracja zależności

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

Alternatywnie możesz również użyć szablonu startowego Spring Cloud Azure Stream Event Hubs, jak pokazano w poniższym przykładzie dla narzędzia Maven:

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

Konfiguracja

Powiążnik udostępnia następujące trzy części opcji konfiguracji:

Właściwości konfiguracji połączenia

Ta sekcja zawiera opcje konfiguracji używane do nawiązywania połączenia z usługą Azure Event Hubs.

Uwaga

Jeśli zdecydujesz się użyć podmiotu zabezpieczeń do uwierzytelniania i autoryzacji za pomocą identyfikatora Entra firmy Microsoft na potrzeby uzyskiwania dostępu do zasobu platformy Azure, zobacz Autoryzuj dostęp za pomocą identyfikatora Entra firmy Microsoft, aby upewnić się, że podmiot zabezpieczeń otrzymał wystarczające uprawnienia dostępu do zasobu platformy Azure.

Konfigurowalne właściwości połączenia spring-cloud-azure-stream-binder-eventhubs:

Własność Typ Opis
spring.cloud.azure.eventhubs.enabled logiczny Określa, czy usługa Azure Event Hubs jest włączona.
spring.cloud.azure.eventhubs.connection-string Struna Wartość parametrów połączenia dla przestrzeni nazw usługi Event Hubs.
spring.cloud.azure.eventhubs.namespace Struna Wartość przestrzeni nazw usługi Event Hubs, stanowiąca prefiks nazwy FQDN. FQDN powinien mieć postać NamespaceName.DomainName
spring.cloud.azure.eventhubs.domain-name Struna Nazwa domeny dla wartości przestrzeni nazw usługi Azure Event Hubs.
spring.cloud.azure.eventhubs.custom-endpoint-address Struna Niestandardowy adres punktu końcowego.

Wskazówka

Typowe opcje konfiguracji pakietu SDK usług Azure można również skonfigurować dla bindera Spring Cloud Azure Stream Event Hubs. Obsługiwane opcje konfiguracji są wprowadzane w konfiguracji platformy Azure Spring Cloudi można je skonfigurować za pomocą ujednoliconego prefiksu spring.cloud.azure. lub prefiksu spring.cloud.azure.eventhubs..

Binder obsługuje również Spring Could Azure Resource Manager domyślnie. Aby dowiedzieć się, jak pobrać ciąg połączenia przy użyciu jednostek zabezpieczeń, którym nie przypisano ról powiązanych z Data, zobacz sekcję Basic usage w Spring Could Azure Resource Manager.

Właściwości konfiguracji punktu kontrolnego

Ta sekcja zawiera opcje konfiguracji usługi Storage Blobs Service, która jest używana do utrwalania własności partycji i informacji o punkcie kontrolnym.

Uwaga

Od wersji 4.0.0, jeśli właściwość spring.cloud.azure.eventhubs.processor.checkpoint-store.create-container-if-not-exists nie zostanie włączona ręcznie, kontener usługi Storage o nazwie z spring.cloud.stream.bindings.binding-name.destination nie będzie tworzony automatycznie.

Tworzenie punktów kontrolnych dla konfigurowalnych właściwości elementu spring-cloud-azure-stream-binder-eventhubs:

Własność Typ Opis
spring.cloud.azure.eventhubs.processor.checkpoint-store.create-container-if-not-exists logiczny Czy zezwolić na tworzenie kontenerów, jeśli jeszcze nie istnieją?
spring.cloud.azure.eventhubs.processor.checkpoint-store.account-name Struna Nazwa konta magazynu.
spring.cloud.azure.eventhubs.processor.checkpoint-store.account-key Struna Klucz dostępu do konta magazynu.
spring.cloud.azure.eventhubs.processor.checkpoint-store.container-name Struna Nazwa kontenera magazynu.

Wskazówka

Typowe opcje konfiguracji zestawu SDK usług Azure mają zastosowanie również do magazynu punktów kontrolnych Storage Blob. Obsługiwane opcje konfiguracji są wprowadzane w konfiguracji platformy Azure Spring Cloudi można je skonfigurować za pomocą ujednoliconego prefiksu spring.cloud.azure. lub prefiksu spring.cloud.azure.eventhubs.processor.checkpoint-store.

Właściwości konfiguracji powiązania usługi Azure Event Hubs

Następujące opcje są podzielone na cztery sekcje: Właściwości konsumenta, Zaawansowane konfiguracje konsumentów, Właściwości producenta i Zaawansowane konfiguracje producenta.

Właściwości konsumenta

Te właściwości są widoczne za pośrednictwem EventHubsConsumerProperties.

Uwaga

Aby uniknąć powtórzeń, od wersji 4.17.0 i 5.11.0 Spring Cloud Azure Stream Binder Event Hubs obsługuje ustawianie wartości dla wszystkich kanałów, w formacie spring.cloud.stream.eventhubs.default.consumer.<property>=<value>.

Właściwości konfigurowalne przez użytkownika końcowego elementu spring-cloud-azure-stream-binder-eventhubs:

Własność Typ Opis
spring.cloud.stream.eventhubs.bindings.binding-name.consumer.checkpoint.mode Tryb punktu kontrolnego Tryb punktu kontrolnego używany, gdy użytkownik decyduje o sposobie wyświetlania komunikatu punktu kontrolnego
spring.cloud.stream.eventhubs.bindings.binding-name.consumer.checkpoint.count Liczba całkowita Określa liczbę komunikatów w każdej partycji potrzebną do wykonania jednego punktu kontrolnego. Zacznie obowiązywać tylko wtedy, gdy jest używany PARTITION_COUNT tryb punktu kontrolnego.
spring.cloud.stream.eventhubs.bindings.binding-name.consumer.checkpoint.interval Czas trwania Decyduje o interwale czasu do wykonania jednego punktu kontrolnego. Zacznie obowiązywać tylko wtedy, gdy jest używany TIME tryb punktu kontrolnego.
spring.cloud.stream.eventhubs.bindings.<binding-name>.consumer.batch.max-size Liczba całkowita Maksymalna liczba zdarzeń w partii. Wymagany dla trybu batch-consumer.
spring.cloud.stream.eventhubs.bindings.binding-name.consumer.batch.max-wait-time Czas trwania Maksymalny czas trwania przetwarzania wsadowego. Zacznie obowiązywać tylko wtedy, gdy tryb batch-consumer jest włączony i jest opcjonalny.
spring.cloud.stream.eventhubs.bindings.binding-name.consumer.load-balancing.update-interval Czas trwania Czas trwania interwału aktualizacji.
spring.cloud.stream.eventhubs.bindings.binding-name.consumer.load-balancing.strategy Strategia równoważenia obciążenia Strategia równoważenia obciążenia.
spring.cloud.stream.eventhubs.bindings.binding-name.consumer.load-balancing.partition-ownership-expiration-interval Czas trwania Okres, po upływie którego wygasa prawo własności partycji.
spring.cloud.stream.eventhubs.bindings.binding-name.consumer.track-last-enqueued-event-properties logiczny Czy procesor zdarzeń powinien żądać informacji o ostatnio umieszczonym w kolejce zdarzeniu na przypisanej partycji i śledzić te informacje w miarę odbierania zdarzeń.
spring.cloud.stream.eventhubs.bindings.binding-name.consumer.prefetch-count Liczba całkowita Liczba używana przez konsumenta do sterowania liczbą zdarzeń, które konsument usługi Event Hub będzie aktywnie odbierał i umieszczał lokalnie w kolejce.
spring.cloud.stream.eventhubs.bindings.binding-name.consumer.initial-partition-event-position Mapowanie przy użyciu klucza jako identyfikatora partycji i wartości StartPositionProperties Mapa zawierająca pozycję zdarzenia, która ma być używana dla każdej partycji, jeśli punkt kontrolny dla danej partycji nie istnieje w magazynie punktów kontrolnych. Ta mapa jest indeksowana według identyfikatora partycji.

Uwaga

Konfiguracja initial-partition-event-position akceptuje map w celu określenia początkowej pozycji dla każdego centrum zdarzeń. Zatem jego kluczem jest identyfikator partycji, a wartością jest StartPositionProperties, który obejmuje właściwości takie jak: przesunięcie, numer sekwencji, data i godzina umieszczenia w kolejce oraz informacja, czy wartość jest uwzględniana. Można na przykład ustawić ją jako

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
Zaawansowana konfiguracja użytkownika

Powyższa konfiguracja połączenia, punktu kontrolnego i wspólna konfiguracja klienta zestawu Azure SDK obsługują dostosowywanie dla każdego konsumenta powiązania, którego można skonfigurować przy użyciu prefiksu spring.cloud.stream.eventhubs.bindings.<binding-name>.consumer..

Właściwości producenta

Te właściwości są widoczne za pośrednictwem EventHubsProducerProperties.

Uwaga

Aby uniknąć powtórzeń, od wersji 4.17.0 i 5.11.0 Spring Cloud Azure Stream Binder Event Hubs obsługuje ustawianie wartości dla wszystkich kanałów, w formacie spring.cloud.stream.eventhubs.default.producer.<property>=<value>.

Właściwości spring-cloud-azure-stream-binder-eventhubs, które może skonfigurować producent:

Własność Typ Opis
spring.cloud.stream.eventhubs.bindings.binding-name.producer.sync logiczny Flaga przełączająca synchronizację producenta. Jeśli ustawiono wartość true, producent będzie czekał na odpowiedź po operacji wysłania.
spring.cloud.stream.eventhubs.bindings.binding-name.producer.send-timeout długi Czas oczekiwania na odpowiedź po operacji wysyłania. Zacznie obowiązywać tylko wtedy, gdy producent synchronizacji jest włączony.
Zaawansowana konfiguracja producenta

Powyższe połączenia i typowe konfiguracje klienta zestawu Azure SDK obsługują dostosowywanie dla każdego producenta binder, który można skonfigurować przy użyciu prefiksu spring.cloud.stream.eventhubs.bindings.<binding-name>.producer..

Podstawowe użycie

Wysyłanie i odbieranie komunikatów z/do usługi Event Hubs

  1. Wypełnij opcje konfiguracji informacjami o poświadczeniach.

    • Jeśli poświadczenia są podawane jako parametry połączenia, skonfiguruj następujące właściwości w pliku 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
      

      Uwaga

      Firma Microsoft zaleca korzystanie z najbezpieczniejszego dostępnego przepływu uwierzytelniania. Przepływ uwierzytelniania opisany w tej procedurze, taki jak bazy danych, pamięci podręczne, komunikaty lub usługi sztucznej inteligencji, wymaga bardzo wysokiego stopnia zaufania w aplikacji i niesie ze sobą ryzyko, które nie występują w innych przepływach. Użyj tego przepływu tylko wtedy, gdy bardziej bezpieczne opcje, takie jak tożsamości zarządzane dla połączeń bez hasła lub bez kluczy, nie są opłacalne. W przypadku operacji maszyny lokalnej preferuj tożsamości użytkowników dla połączeń bez hasła lub bez klucza.

    • W przypadku poświadczeń jednostki usługi skonfiguruj następujące właściwości w pliku 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
      

Uwaga

Dozwolone wartości dla tenant-id to: common, organizations, consumerslub identyfikator dzierżawy. Aby uzyskać więcej informacji na temat tych wartości, zobacz sekcję Użyto niewłaściwego punktu końcowego (konta osobiste i organizacyjne) w artykule Błąd AADSTS50020 — konto użytkownika od dostawcy tożsamości nie istnieje w dzierżawie. Aby uzyskać informacje na temat konwersji aplikacji z jedną dzierżawą na aplikację wielodzierżawną, zobacz Konwertowanie aplikacji z jedną dzierżawą na aplikację wielodzierżawną w usłudze Microsoft Entra ID.

  • W przypadku poświadczeń jako tożsamości zarządzanych skonfiguruj następujące właściwości w pliku 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. Definiowanie dostawcy i konsumenta.

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

Obsługa partycjonowania

Tworzona jest PartitionSupplier z informacjami o partycji podanymi przez użytkownika w celu skonfigurowania informacji o partycji dla wiadomości, która ma zostać wysłana. Poniższy schemat blokowy przedstawia proces uzyskiwania różnych priorytetów dla identyfikatora partycji i klucza:

Diagram przedstawiający schemat blokowy procesu obsługi partycjonowania.

Obsługa konsumentów usługi Batch

  1. Podaj opcje konfiguracji wsadowej, jak pokazano na poniższym przykładzie:

    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. Definiowanie dostawcy i konsumenta.

    Aby w trybie checkpointingu ustawionym na BATCH można było wysyłać komunikaty i odbierać je partiami, użyj następującego kodu.

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

    W przypadku trybu checkpointingu ustawionego na MANUAL można użyć poniższego kodu do wysyłania komunikatów oraz ich odbierania i checkpointowania wsadowo.

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

Uwaga

W trybie wsadowego przetwarzania komunikatów domyślny typ treści bindera Spring Cloud Stream to application/json, dlatego upewnij się, że ładunek wiadomości jest zgodny z tym typem treści. Na przykład w przypadku używania domyślnego typu zawartości application/json do odbierania komunikatów z ładunkiem String ładunek powinien być JSON String, otoczony podwójnymi cudzysłowymi dla oryginalnego tekstu String. Natomiast w przypadku typu zawartości text/plain może to być bezpośrednio obiekt String. Aby uzyskać więcej informacji, zobacz Negocjowanie typu zawartości strumienia Spring Cloud.

Obsługa komunikatów o błędach

  • Obsługa komunikatów o błędach powiązania wychodzącego

    Domyślnie platforma Spring Integration tworzy globalny kanał błędów o nazwie errorChannel. Skonfiguruj następujący punkt końcowy komunikatu, aby obsługiwać komunikaty o błędach powiązania wychodzącego.

    @ServiceActivator(inputChannel = IntegrationContextUtils.ERROR_CHANNEL_BEAN_NAME)
    public void handleError(ErrorMessage message) {
        LOGGER.error("Handling outbound binding error: " + message);
    }
    
  • Obsłuż komunikaty o błędach powiązań przychodzących

    Binder Spring Cloud Stream Event Hubs udostępnia jeden sposób obsługi błędów dla powiązań komunikatów przychodzących: moduły obsługi błędów.

    program obsługi błędów :

    Spring Cloud Stream udostępnia mechanizm umożliwiający zdefiniowanie niestandardowego programu obsługi błędów przez dodanie elementu Consumer, przyjmującego instancje ErrorMessage. Aby uzyskać więcej informacji, zobacz Obsługa komunikatów o błędach w dokumentacji usługi Spring Cloud Stream.

    • Domyślny moduł obsługi błędów wiązania

      Skonfiguruj pojedynczy komponent Consumer bean do obsługi wszystkich komunikatów o błędach wejściowego powiązania. Poniższa funkcja domyślna subskrybuje kanał błędów dla każdego powiązania wejściowego:

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

      Należy również ustawić we właściwości spring.cloud.stream.default.error-handler-definition nazwę funkcji.

    • Procedura obsługi błędów specyficzna dla powiązania

      Skonfiguruj bean Consumer, aby obsługiwał określone komunikaty o błędach dla powiązania wejściowego. Następująca funkcja subskrybuje określony kanał błędów powiązania przychodzącego i ma wyższy priorytet niż domyślna procedura obsługi błędów powiązania:

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

      Należy również ustawić we właściwości spring.cloud.stream.bindings.<input-binding-name>.error-handler-definition nazwę funkcji.

Nagłówki komunikatów usługi Event Hubs

Aby zapoznać się z obsługiwanymi podstawowymi nagłówkami komunikatów, zobacz sekcję Nagłówki komunikatów usługi Event Hubs w Obsługa Spring Cloud Azure dla Spring Integration.

Obsługa wielu segregatorów

Obsługiwane jest również połączenie z wieloma przestrzeniami nazw usługi Event Hubs przy użyciu wielu binderów. W tym przykładzie użyto parametrów połączenia jako przykładu. Obsługiwane są również poświadczenia jednostek usługi i tożsamości zarządzanych. Możesz ustawić powiązane właściwości w ustawieniach środowiska każdego bindera.

  1. Aby użyć wielu binderów dla usługi Event Hubs, skonfiguruj następujące właściwości w pliku 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
    

    Uwaga

    W poprzednim pliku aplikacji pokazano, jak skonfigurować jeden domyślny moduł odpytywania dla aplikacji, mający zastosowanie do wszystkich powiązań. Jeśli chcesz skonfigurować narzędzie poller dla określonego powiązania, możesz użyć konfiguracji, takiej jak spring.cloud.stream.bindings.<binding-name>.producer.poller.fixed-delay=3000.

    Uwaga

    Firma Microsoft zaleca korzystanie z najbezpieczniejszego dostępnego przepływu uwierzytelniania. Przepływ uwierzytelniania opisany w tej procedurze, taki jak bazy danych, pamięci podręczne, komunikaty lub usługi sztucznej inteligencji, wymaga bardzo wysokiego stopnia zaufania w aplikacji i niesie ze sobą ryzyko, które nie występują w innych przepływach. Użyj tego przepływu tylko wtedy, gdy bardziej bezpieczne opcje, takie jak tożsamości zarządzane dla połączeń bez hasła lub bez kluczy, nie są opłacalne. W przypadku operacji maszyny lokalnej preferuj tożsamości użytkowników dla połączeń bez hasła lub bez klucza.

  2. Musimy zdefiniować dwóch dostawców i dwóch konsumentów:

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

Aprowizowanie zasobów

Binder usługi Event Hubs obsługuje aprowizację centrum zdarzeń i grupy odbiorców. Użytkownicy mogą użyć następujących właściwości w celu włączenia aprowizacji.

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

Uwaga

Dozwolone wartości dla tenant-id to: common, organizations, consumerslub identyfikator dzierżawy. Aby uzyskać więcej informacji na temat tych wartości, zobacz sekcję Użyto niewłaściwego punktu końcowego (konta osobiste i organizacyjne) w artykule Błąd AADSTS50020 — konto użytkownika od dostawcy tożsamości nie istnieje w dzierżawie. Aby uzyskać informacje na temat konwersji aplikacji z jedną dzierżawą na aplikację wielodzierżawną, zobacz Konwertowanie aplikacji z jedną dzierżawą na aplikację wielodzierżawną w usłudze Microsoft Entra ID.

Próbki

Więcej informacji można znaleźć w azure-spring-boot-samples repozytorium w serwisie GitHub.

Binder Spring Cloud Stream dla usługi Azure Service Bus

Kluczowe pojęcia

Binder Spring Cloud Stream dla Azure Service Bus zapewnia implementację mechanizmu wiązania dla platformy Spring Cloud Stream. Ta implementacja opiera się na adapterach kanałów usługi Spring Integration Service Bus.

Zaplanowana wiadomość

Ten binder obsługuje przesyłanie komunikatów do tematu w celu późniejszego przetworzenia. Użytkownicy mogą wysyłać zaplanowane komunikaty z nagłówkiem x-delay wyrażając w milisekundach czas opóźnienia komunikatu. Wiadomość zostanie dostarczona do odpowiednich tematów po x-delay milisekundach.

Grupa odbiorców

Temat usługi Service Bus zapewnia obsługę grup konsumentów podobną do tej w Apache Kafka, ale działa w nieco inny sposób. Ten binder opiera się na Subscription tematu, który działa jako grupa odbiorców.

Konfiguracja zależności

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

Alternatywnie, możesz również użyć startera Spring Cloud Azure Stream Service Bus, jak pokazano w poniższym przykładzie dla programu Maven:

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

Konfiguracja

Binder udostępnia następujące dwie części opcji konfiguracji:

Właściwości konfiguracji połączenia

Ta sekcja zawiera opcje konfiguracji używane do nawiązywania połączenia z usługą Azure Service Bus.

Uwaga

Jeśli zdecydujesz się użyć podmiotu zabezpieczeń do uwierzytelniania i autoryzacji za pomocą identyfikatora Entra firmy Microsoft na potrzeby uzyskiwania dostępu do zasobu platformy Azure, zobacz Autoryzuj dostęp za pomocą identyfikatora Entra firmy Microsoft, aby upewnić się, że podmiot zabezpieczeń otrzymał wystarczające uprawnienia dostępu do zasobu platformy Azure.

Konfigurowalne właściwości połączenia spring-cloud-azure-stream-binder-servicebus:

Własność Typ Opis
spring.cloud.azure.servicebus.enabled logiczny Określa, czy usługa Azure Service Bus jest włączona.
spring.cloud.azure.servicebus.connection-string Struna Wartość parametrów połączenia przestrzeni nazw usługi Service Bus.
spring.cloud.azure.servicebus.custom-endpoint-address Struna Niestandardowy adres punktu końcowego do użycia podczas nawiązywania połączenia z usługą Service Bus.
spring.cloud.azure.servicebus.namespace Struna Wartość przestrzeni nazw usługi Service Bus, która jest prefiksem nazwy FQDN. FQDN powinien mieć postać NamespaceName.DomainName
spring.cloud.azure.servicebus.domain-name Struna Nazwa domeny wartości przestrzeni nazw usługi Azure Service Bus.

Uwaga

Typowe opcje konfiguracji zestawu SDK usług Azure można również skonfigurować dla bindera Spring Cloud Azure Stream Service Bus. Obsługiwane opcje konfiguracji są wprowadzane w konfiguracji platformy Azure Spring Cloudi można je skonfigurować za pomocą ujednoliconego prefiksu spring.cloud.azure. lub prefiksu spring.cloud.azure.servicebus..

Binder obsługuje również Spring Could Azure Resource Manager domyślnie. Aby dowiedzieć się, jak pobrać ciąg połączenia przy użyciu jednostek zabezpieczeń, którym nie przypisano ról powiązanych z Data, zobacz sekcję Basic usage w Spring Could Azure Resource Manager.

Właściwości konfiguracji powiązania usługi Azure Service Bus

Następujące opcje są podzielone na cztery sekcje: Właściwości konsumenta, Zaawansowane konfiguracje konsumentów, Właściwości producenta i Zaawansowane konfiguracje producenta.

Właściwości konsumenta

Te właściwości są widoczne za pośrednictwem ServiceBusConsumerProperties.

Uwaga

Aby uniknąć powtórzeń, od wersji 4.17.0 i 5.11.0 usługa Spring Cloud Azure Stream Binder Service Bus obsługuje ustawianie wartości dla wszystkich kanałów w postaci spring.cloud.stream.servicebus.default.consumer.<property>=<value>.

Właściwości konfigurowalne przez użytkownika końcowego elementu spring-cloud-azure-stream-binder-servicebus:

Własność Typ Domyślny Opis
spring.cloud.stream.servicebus.bindings.binding-name.consumer.requeue-rejected logiczny fałszywy Jeśli komunikaty, które zakończyły się niepowodzeniem, są kierowane do biblioteki DLQ.
spring.cloud.stream.servicebus.bindings.binding-name.consumer.max-concurrent-calls Liczba całkowita 1 Maksymalna liczba współbieżnych komunikatów, które powinien przetworzyć klient procesora usługi Service Bus. Po włączeniu sesji ma zastosowanie do każdej sesji.
spring.cloud.stream.servicebus.bindings.binding-name.consumer.max-concurrent-sessions Liczba całkowita null Maksymalna liczba współbieżnych sesji do przetworzenia w danym momencie.
spring.cloud.stream.servicebus.bindings.binding-name.consumer.session-enabled logiczny null Czy sesja jest włączona.
spring.cloud.stream.servicebus.bindings.binding-name.consumer.session-idle-timeout Czas trwania null Ustawia maksymalny czas (czas trwania) oczekiwania na odebranie komunikatu dla aktualnie aktywnej sesji.
spring.cloud.stream.servicebus.bindings.binding-name.consumer.prefetch-count Liczba całkowita 0 Liczba komunikatów pobieranych z wyprzedzeniem przez klienta procesora usługi Service Bus.
spring.cloud.stream.servicebus.bindings.binding-name.consumer.sub-queue Podkolejka Brak Typ podkolejki, z którą należy się połączyć.
spring.cloud.stream.servicebus.bindings.binding-name.consumer.max-auto-lock-renew-duration Czas trwania 5 m Czas na kontynuowanie automatycznego odnawiania blokady.
spring.cloud.stream.servicebus.bindings.binding-name.consumer.receive-mode ServiceBusReceiveMode peek_lock Tryb odbioru klienta procesora Service Bus.
spring.cloud.stream.servicebus.bindings.binding-name.consumer.auto-complete logiczny prawdziwy Czy automatycznie rozliczać wiadomości. W przypadku ustawienia wartości false nagłówek komunikatu Checkpointer zostanie dodany w celu umożliwienia deweloperom ręcznego rozliczania komunikatów.
spring.cloud.stream.servicebus.bindings.binding-name.consumer.max-size-in-megabytes Długi 1024 Maksymalny rozmiar kolejki/tematu w megabajtach, czyli rozmiar pamięci przydzielonej do kolejki/tematu.
spring.cloud.stream.servicebus.bindings.binding-name.consumer.default-message-time-to-live Czas trwania P10675199DT2H48M5.4775807S. (10675199 dni, 2 godziny, 48 minut, 5 sekund i 477 milisekund) Czas trwania, po upływie którego komunikat wygaśnie, począwszy od momentu wysłania komunikatu do usługi Service Bus.

Ważny

W przypadku korzystania z usługi Azure Resource Manager (ARM) należy skonfigurować właściwość spring.cloud.stream.servicebus.bindings.<binding-name>.consume.entity-type. Aby uzyskać więcej informacji, zobacz przykład w witrynie servicebus-queue-binder-arm GitHub.

Zaawansowana konfiguracja użytkownika

Powyższe połączenia i typowe konfiguracje klienta zestawu Azure SDK obsługują dostosowywanie dla każdego konsumenta binder, który można skonfigurować przy użyciu prefiksu spring.cloud.stream.servicebus.bindings.<binding-name>.consumer..

Właściwości producenta

Te właściwości są widoczne za pośrednictwem ServiceBusProducerProperties.

Uwaga

Aby uniknąć powtórzeń, od wersji 4.17.0 i 5.11.0 usługa Spring Cloud Azure Stream Binder Service Bus obsługuje ustawianie wartości dla wszystkich kanałów w postaci spring.cloud.stream.servicebus.default.producer.<property>=<value>.

Właściwości spring-cloud-azure-stream-binder-servicebus, które może skonfigurować producent:

Własność Typ Domyślny Opis
spring.cloud.stream.servicebus.bindings.binding-name.producer.sync logiczny fałszywy Przełącz flagę synchronizacji producenta.
spring.cloud.stream.servicebus.bindings.binding-name.producer.send-timeout długi 10 000 Wartość limitu czasu wysyłania dla producenta.
spring.cloud.stream.servicebus.bindings.binding-name.producer.entity-type ServiceBusEntityType null Typ jednostki usługi Service Bus producenta, wymagany dla producenta powiązania.
spring.cloud.stream.servicebus.bindings.binding-name.producer.max-size-in-megabytes Długi 1024 Maksymalny rozmiar kolejki/tematu w megabajtach, czyli rozmiar pamięci przydzielonej do kolejki/tematu.
spring.cloud.stream.servicebus.bindings.binding-name.producer.default-message-time-to-live Czas trwania P10675199DT2H48M5.4775807S. (10675199 dni, 2 godziny, 48 minut, 5 sekund i 477 milisekund) Czas trwania, po upływie którego komunikat wygaśnie, począwszy od momentu wysłania komunikatu do usługi Service Bus.

Ważny

Podczas korzystania z producenta powiązania skonfigurowanie właściwości spring.cloud.stream.servicebus.bindings.<binding-name>.producer.entity-type jest wymagane.

Zaawansowana konfiguracja producenta

Powyższe połączenia i typowe konfiguracje klienta zestawu Azure SDK obsługują dostosowywanie dla każdego producenta binder, który można skonfigurować przy użyciu prefiksu spring.cloud.stream.servicebus.bindings.<binding-name>.producer..

Podstawowe użycie

Wysyłanie i odbieranie komunikatów z/do usługi Service Bus

  1. Wypełnij opcje konfiguracji informacjami o poświadczeniach.

    • Jeśli poświadczenia są podawane jako parametry połączenia, skonfiguruj następujące właściwości w pliku 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
      

      Uwaga

      Firma Microsoft zaleca korzystanie z najbezpieczniejszego dostępnego przepływu uwierzytelniania. Przepływ uwierzytelniania opisany w tej procedurze, taki jak bazy danych, pamięci podręczne, komunikaty lub usługi sztucznej inteligencji, wymaga bardzo wysokiego stopnia zaufania w aplikacji i niesie ze sobą ryzyko, które nie występują w innych przepływach. Użyj tego przepływu tylko wtedy, gdy bardziej bezpieczne opcje, takie jak tożsamości zarządzane dla połączeń bez hasła lub bez kluczy, nie są opłacalne. W przypadku operacji maszyny lokalnej preferuj tożsamości użytkowników dla połączeń bez hasła lub bez klucza.

    • W przypadku poświadczeń jednostki usługi skonfiguruj następujące właściwości w pliku 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
      

Uwaga

Dozwolone wartości dla tenant-id to: common, organizations, consumerslub identyfikator dzierżawy. Aby uzyskać więcej informacji na temat tych wartości, zobacz sekcję Użyto niewłaściwego punktu końcowego (konta osobiste i organizacyjne) w artykule Błąd AADSTS50020 — konto użytkownika od dostawcy tożsamości nie istnieje w dzierżawie. Aby uzyskać informacje na temat konwersji aplikacji z jedną dzierżawą na aplikację wielodzierżawną, zobacz Konwertowanie aplikacji z jedną dzierżawą na aplikację wielodzierżawną w usłudze Microsoft Entra ID.

  • W przypadku poświadczeń jako tożsamości zarządzanych skonfiguruj następujące właściwości w pliku 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. Definiowanie dostawcy i konsumenta.

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

Obsługa klucza partycji

Binder obsługuje partycjonowanie usługi Service Bus, umożliwiając ustawianie klucza partycji i identyfikatora sesji w nagłówku komunikatu. W tej sekcji przedstawiono sposób ustawiania klucza partycji dla komunikatów.

Usługa Spring Cloud Stream udostępnia właściwość wyrażenia SpEL klucza partycji spring.cloud.stream.bindings.<binding-name>.producer.partition-key-expression. Na przykład ustawienie tej właściwości jako "'partitionKey-' + headers[<message-header-key>]" i dodanie nagłówka o nazwie message-header-key. Usługa Spring Cloud Stream używa wartości dla tego nagłówka podczas oceniania wyrażenia w celu przypisania klucza partycji. Poniższy kod zawiera przykładowego producenta:

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

Obsługa sesji

Łącznik obsługuje sesje wiadomości usługi Service Bus. Identyfikator sesji wiadomości można ustawić za pośrednictwem nagłówka komunikatu.

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

Uwaga

Zgodnie z informacjami o partycjonowaniu usługi Service Bus identyfikator sesji ma wyższy priorytet niż klucz partycji. Dlatego po ustawieniu zarówno nagłówków ServiceBusMessageHeaders#SESSION_ID, jak i ServiceBusMessageHeaders#PARTITION_KEY wartość identyfikatora sesji zostanie ostatecznie użyta do zastąpienia wartości klucza partycji.

Obsługa komunikatów o błędach

  • Obsługa komunikatów o błędach powiązania wychodzącego

    Domyślnie platforma Spring Integration tworzy globalny kanał błędów o nazwie errorChannel. Skonfiguruj następujący punkt końcowy komunikatu, aby obsłużyć komunikat o błędzie powiązania wychodzącego.

    @ServiceActivator(inputChannel = IntegrationContextUtils.ERROR_CHANNEL_BEAN_NAME)
    public void handleError(ErrorMessage message) {
        LOGGER.error("Handling outbound binding error: " + message);
    }
    
  • Obsłuż komunikaty o błędach powiązań przychodzących

    Binder usługi Service Bus dla Spring Cloud Stream obsługuje dwa rozwiązania do obsługi błędów dla powiązań przychodzących komunikatów: program obsługi błędów bindera oraz programy obsługi.

    Obsługa błędów Binder:

    Domyślna procedura obsługi błędów bindera obsługuje powiązanie przychodzące. Ten moduł obsługi służy do wysyłania komunikatów, których przetworzenie zakończyło się niepowodzeniem, do kolejki komunikatów niedostarczonych, gdy włączono spring.cloud.stream.servicebus.bindings.<binding-name>.consumer.requeue-rejected. W przeciwnym razie komunikaty, które zakończyły się niepowodzeniem, zostaną porzucone. Procedura obsługi błędów bindera wzajemnie się wyklucza z innymi dostarczonymi procedurami obsługi błędów.

    program obsługi błędów :

    Spring Cloud Stream udostępnia mechanizm umożliwiający zdefiniowanie niestandardowego programu obsługi błędów przez dodanie elementu Consumer, przyjmującego instancje ErrorMessage. Aby uzyskać więcej informacji, zobacz Obsługa komunikatów o błędach w dokumentacji usługi Spring Cloud Stream.

    • Domyślny moduł obsługi błędów wiązania

      Skonfiguruj pojedynczy komponent Consumer bean do obsługi wszystkich komunikatów o błędach wejściowego powiązania. Poniższa funkcja domyślna subskrybuje kanał błędów dla każdego powiązania wejściowego:

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

      Należy również ustawić we właściwości spring.cloud.stream.default.error-handler-definition nazwę funkcji.

    • Procedura obsługi błędów specyficzna dla powiązania

      Skonfiguruj bean Consumer, aby obsługiwał określone komunikaty o błędach dla powiązania wejściowego. Poniższa funkcja subskrybuje określony kanał błędów powiązania przychodzącego o wyższym priorytcie niż procedura obsługi błędów powiązania domyślnego.

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

      Należy również ustawić we właściwości spring.cloud.stream.bindings.<input-binding-name>.error-handler-definition nazwę funkcji.

Nagłówki komunikatów usługi Service Bus

Aby zapoznać się z obsługiwanymi podstawowymi nagłówkami komunikatów, zobacz sekcję nagłówki komunikatów usługi Service Bus w dokumencie obsługi Spring Cloud Azure dla Spring Integration.

Uwaga

Podczas ustawiania klucza partycji priorytet nagłówka komunikatu jest wyższy niż właściwość Spring Cloud Stream. Dlatego spring.cloud.stream.bindings.<binding-name>.producer.partition-key-expression zaczynają obowiązywać tylko wtedy, gdy żaden z nagłówków ServiceBusMessageHeaders#SESSION_ID i ServiceBusMessageHeaders#PARTITION_KEY nie jest skonfigurowany.

Obsługa wielu segregatorów

Obsługiwane jest również połączenie z wieloma przestrzeniami nazw usługi Service Bus przy użyciu wielu binderów. W tym przykładzie użyto parametrów połączenia jako przykładu. Obsługiwane są również poświadczenia jednostek usługi i tożsamości zarządzanych, a użytkownicy mogą skonfigurować powiązane właściwości w ustawieniach środowiska dla każdego powiązania.

  1. Aby używać wielu binderów usługi ServiceBus, skonfiguruj następujące właściwości w swoim pliku 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
    

    Uwaga

    W poprzednim pliku aplikacji pokazano, jak skonfigurować jeden domyślny moduł odpytywania dla aplikacji, mający zastosowanie do wszystkich powiązań. Jeśli chcesz skonfigurować narzędzie poller dla określonego powiązania, możesz użyć konfiguracji, takiej jak spring.cloud.stream.bindings.<binding-name>.producer.poller.fixed-delay=3000.

    Uwaga

    Firma Microsoft zaleca korzystanie z najbezpieczniejszego dostępnego przepływu uwierzytelniania. Przepływ uwierzytelniania opisany w tej procedurze, taki jak bazy danych, pamięci podręczne, komunikaty lub usługi sztucznej inteligencji, wymaga bardzo wysokiego stopnia zaufania w aplikacji i niesie ze sobą ryzyko, które nie występują w innych przepływach. Użyj tego przepływu tylko wtedy, gdy bardziej bezpieczne opcje, takie jak tożsamości zarządzane dla połączeń bez hasła lub bez kluczy, nie są opłacalne. W przypadku operacji maszyny lokalnej preferuj tożsamości użytkowników dla połączeń bez hasła lub bez klucza.

  2. potrzebujemy zdefiniowania dwóch dostawców i dwóch konsumentów

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

Aprowizowanie zasobów

Łącznik usługi Service Bus obsługuje tworzenie kolejki, tematu i subskrypcji. Użytkownicy mogą używać następujących właściwości, aby włączyć tworzenie tych elementów.

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}

Uwaga

Dozwolone wartości dla tenant-id to: common, organizations, consumerslub identyfikator dzierżawy. Aby uzyskać więcej informacji na temat tych wartości, zobacz sekcję Użyto niewłaściwego punktu końcowego (konta osobiste i organizacyjne) w artykule Błąd AADSTS50020 — konto użytkownika od dostawcy tożsamości nie istnieje w dzierżawie. Aby uzyskać informacje na temat konwersji aplikacji z jedną dzierżawą na aplikację wielodzierżawną, zobacz Konwertowanie aplikacji z jedną dzierżawą na aplikację wielodzierżawną w usłudze Microsoft Entra ID.

Dostosowywanie właściwości klienta usługi Service Bus

Deweloperzy mogą używać AzureServiceClientBuilderCustomizer do dostosowywania właściwości klienta usługi Service Bus. Poniższy przykład dostosowuje właściwość sessionIdleTimeout w ServiceBusClientBuilder:

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

Próbki

Więcej informacji można znaleźć w azure-spring-boot-samples repozytorium w serwisie GitHub.