Notatka
Dostęp do tej strony wymaga autoryzacji. Może spróbować zalogować się lub zmienić katalogi.
Dostęp do tej strony wymaga autoryzacji. Możesz spróbować zmienić katalogi.
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-azure-stream-binder-eventhubs— aby uzyskać więcej informacji, zobacz Łącznik Spring Cloud Stream dla usługi Azure Event Hubs -
spring-cloud-azure-stream-binder-servicebus— aby uzyskać więcej informacji, zobacz Spring Cloud Stream Binder for Azure Service Bus
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
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: MANUALUwaga
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
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:
Obsługa konsumentów usługi Batch
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 neededDefiniowanie dostawcy i konsumenta.
Aby w trybie checkpointingu ustawionym na
BATCHmoż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
MANUALmoż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 instancjeErrorMessage. 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
Consumerbean 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-definitionnazwę 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-definitionnazwę 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.
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: 1000Uwaga
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.
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
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 TopicUwaga
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
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 instancjeErrorMessage. 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
Consumerbean 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-definitionnazwę 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-definitionnazwę 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.
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: 1000Uwaga
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.
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.