Spring Cloud Azure-Unterstützung für Spring Cloud Stream

Spring Cloud Stream ist ein Framework für die Erstellung hoch skalierbarer ereignisgesteuerter Mikroservices, die mit freigegebenen Messagingsystemen verbunden sind.

Das Framework bietet ein flexibles Programmiermodell, das auf bereits etablierten und vertrauten Spring-Idioms und bewährten Methoden basiert. Zu diesen bewährten Methoden zählen die Unterstützung von persistenter Pub/Sub-Semantik, Consumer-Gruppen und zustandsbehafteten Partitionen.

Zu den aktuellen Binderimplementierungen gehören:

Spring Cloud Stream Binder für Azure Event Hubs

Schlüsselkonzepte

Der Spring Cloud Stream Binder für Azure Event Hubs stellt die Bindungsimplementierung für das Spring Cloud Stream-Framework bereit. Diese Implementierung basiert auf Spring Integration Event Hubs-Kanaladaptern. Aus Der Perspektive des Designs ist Event Hubs ähnlich wie Kafka. Außerdem kann über die Kafka-API auf Event Hubs zugegriffen werden. Wenn Ihr Projekt stark von der Kafka-API abhängt, können Sie Events Hub mit Kafka-API-Beispiel ausprobieren.

Verbrauchergruppe

Event Hubs bietet ähnlich wie Apache Kafka Unterstützung für Consumergruppen, jedoch mit einer leicht abweichenden Logik. Während Kafka alle zugesicherten Offsets im Broker speichert, müssen Sie Offsets von Event Hubs-Nachrichten speichern, die manuell verarbeitet werden. Event Hubs SDK stellt die Funktion zum Speichern solcher Offsets in Azure Storage bereit.

Partitionierungsunterstützung

Event Hubs bietet ein ähnliches Konzept der physischen Partition wie Kafka. Im Gegensatz zu Kafkas automatischem Rebalancing zwischen Verbrauchern und Partitionen bietet Event Hubs jedoch eine Art präemptiven Modus. Das Speicherkonto fungiert als Lease, um zu bestimmen, welcher Verbraucher welche Partition besitzt. Wenn ein neuer Consumer startet, versucht er, den am stärksten ausgelasteten Consumern einige Partitionen zu entziehen, um eine ausgeglichene Lastverteilung zu erreichen.

Um die Lastenausgleichsstrategie anzugeben, werden Eigenschaften von spring.cloud.stream.eventhubs.bindings.<binding-name>.consumer.load-balancing.* bereitgestellt. Weitere Informationen finden Sie im Abschnitt Consumer-Eigenschaften.

Batch-Consumerunterstützung

Spring Cloud Azure Stream Event Hubs Binder unterstützt die Spring Cloud Stream Batch Consumer-Funktion.

Um mit dem Batch-Consumer-Modus zu arbeiten, legen Sie die eigenschaft spring.cloud.stream.bindings.<binding-name>.consumer.batch-mode auf truefest. Wenn dies aktiviert ist, wird eine Nachricht mit der Nutzlast einer Liste gebündelter Ereignisse empfangen und an die Funktion Consumer übergeben. Jeder Nachrichtenheader wird ebenfalls in eine Liste umgewandelt, deren Inhalt der jeweils aus jedem Ereignis geparste zugehörige Headerwert ist. Die gemeinsamen Header von Partitions-ID, checkpointer und den Eigenschaften der letzten Einreihung werden als ein einzelner Wert dargestellt, da die gesamte Ereignisgruppe denselben Wert hat. Weitere Informationen finden Sie in den Event Hubs-Nachrichtenkopfzeilen Abschnitt Spring Cloud Azure-Unterstützung für Spring Integration.

Anmerkung

Der Prüfpunktheader ist nur vorhanden, wenn der MANUAL Prüfpunktmodus verwendet wird.

Checkpointing für Batch-Consumer unterstützt zwei Modi: BATCH und MANUAL. BATCH-Modus ist ein automatischer Modus zum Setzen von Prüfpunkten, bei dem für den gesamten Batch von Ereignissen gemeinsam ein Prüfpunkt gesetzt wird, sobald der Binder sie empfängt. MANUAL-Modus dient dazu, Prüfpunkte für Benutzerereignisse zu erstellen. Wenn die Checkpointer verwendet wird, wird sie in den Nachrichtenheader eingefügt, und Benutzer können sie zum Setzen von Prüfpunkten verwenden.

Sie können die Batchgröße angeben, indem Sie die eigenschaften max-size und max-wait-time festlegen, die über ein Präfix von spring.cloud.stream.eventhubs.bindings.<binding-name>.consumer.batch.verfügen. Die max-size-Eigenschaft ist erforderlich, und die max-wait-time Eigenschaft ist optional. Weitere Informationen finden Sie im Abschnitt Consumer-Eigenschaften.

Einrichtung von Abhängigkeiten

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

Alternativ können Sie auch den Spring Cloud Azure Stream Event Hubs Starter verwenden, wie im folgenden Beispiel für Maven gezeigt:

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

Konfiguration

Der Binder bietet die folgenden drei Bereiche der Konfigurationsoptionen:

Verbindungskonfigurationseigenschaften

Dieser Abschnitt enthält die Konfigurationsoptionen, die zum Herstellen einer Verbindung mit Azure Event Hubs verwendet werden.

Anmerkung

Wenn Sie einen Sicherheitsprinzipal zum Authentifizieren und Autorisieren mit Microsoft Entra-ID für den Zugriff auf eine Azure-Ressource verwenden, lesen Sie Autorisieren des Zugriffs mit Microsoft Entra ID, um sicherzustellen, dass dem Sicherheitsprinzipal die ausreichende Berechtigung für den Zugriff auf die Azure-Ressource gewährt wurde.

Konfigurierbare Verbindungseigenschaften von spring-cloud-azure-stream-binder-eventhubs:

Eigentum Typ Beschreibung
spring.cloud.azure.eventhubs.enabled boolesch Gibt an, ob azure Event Hubs aktiviert ist.
spring.cloud.azure.eventhubs.connection-string Schnur Event Hubs Namespace-Verbindungszeichenfolgenwert.
spring.cloud.azure.eventhubs.namespace Schnur Wert des Event Hubs-Namespace, der das Präfix des FQDN darstellt. Ein FQDN sollte aus NamespaceName.DomainName bestehen.
spring.cloud.azure.eventhubs.domain-name Schnur Domänenname eines Azure Event Hubs-Namespacewerts.
spring.cloud.azure.eventhubs.custom-endpoint-address Schnur Benutzerdefinierte Endpunktadresse.

Tipp

Allgemeine Konfigurationsoptionen des Azure Service SDK sind auch für den Spring Cloud Azure Stream Event Hubs Binder konfigurierbar. Die unterstützten Konfigurationsoptionen werden in Spring Cloud Azure-Konfigurationeingeführt und können entweder mit dem einheitlichen Präfix spring.cloud.azure. oder dem Präfix von spring.cloud.azure.eventhubs.konfiguriert werden.

Der Ordner unterstützt auch Spring Could Azure Resource Manager standardmäßig. Informationen dazu, wie Sie die Verbindungszeichenfolge mit Sicherheitsprinzipalen abrufen, denen keine Data-bezogenen Rollen zugewiesen wurden, finden Sie im Abschnitt Grundlegende Verwendung von Spring Cloud Azure Resource Manager.

Eigenschaften der Prüfpunktkonfiguration

Dieser Abschnitt enthält die Konfigurationsoptionen für den Speicher-Blobs-Dienst, der zum Beibehalten des Partitionsbesitzes und der Prüfpunktinformationen verwendet wird.

Anmerkung

Ab Version 4.0.0 wird, wenn die Eigenschaft von spring.cloud.azure.eventhubs.processor.checkpoint-store.create-container-if-not-exists nicht manuell aktiviert ist, kein Speichercontainer automatisch mit dem Namen erstellt.spring.cloud.stream.bindings.binding-name.destination

Prüfpunkterstellung für konfigurierbare Eigenschaften von spring-cloud-azure-stream-binder-eventhubs:

Eigentum Typ Beschreibung
spring.cloud.azure.eventhubs.processor.checkpoint-store.create-container-if-not-exists Boolesch Gibt an, ob Container erstellt werden sollen, wenn sie noch nicht vorhanden sind.
spring.cloud.azure.eventhubs.processor.checkpoint-store.account-name Schnur Name für das Speicherkonto.
spring.cloud.azure.eventhubs.processor.checkpoint-store.account-key Schnur Zugriffsschlüssel für Speicherkonto.
spring.cloud.azure.eventhubs.processor.checkpoint-store.container-name Schnur Name des Speichercontainers.

Tipp

Allgemeine Konfigurationsoptionen des Azure Service SDK können ebenfalls für den Prüfpunktspeicher für Storage Blobs konfiguriert werden. Die unterstützten Konfigurationsoptionen werden in Spring Cloud Azure-Konfigurationeingeführt und können entweder mit dem einheitlichen Präfix spring.cloud.azure. oder dem Präfix von spring.cloud.azure.eventhubs.processor.checkpoint-storekonfiguriert werden.

Azure Event Hubs Binding-Konfigurationseigenschaften

Die folgenden Optionen sind in vier Abschnitte unterteilt: Consumer-Eigenschaften, erweiterte Consumer-Konfigurationen, Producer-Eigenschaften und erweiterte Producer-Konfigurationen.

Verbrauchereigenschaften

Diese Eigenschaften werden über EventHubsConsumerPropertiesverfügbar gemacht.

Anmerkung

Um Wiederholungen zu vermeiden, unterstützt Spring Cloud Azure Stream Binder Event Hubs seit Version 4.17.0 und 5.11.0 das Festlegen von Werten für alle Kanäle im Format von spring.cloud.stream.eventhubs.default.consumer.<property>=<value>.

Konfigurierbare Verbrauchereigenschaften von spring-cloud-azure-stream-binder-eventhubs:

Eigentum Typ Beschreibung
spring.cloud.stream.eventhubs.bindings.binding-name.consumer.checkpoint.mode CheckpointMode Prüfpunktmodus, der verwendet wird, wenn Verbraucher entscheiden, wie die Prüfpunktmeldung ausgeführt werden soll
spring.cloud.stream.eventhubs.bindings.binding-name.consumer.checkpoint.count Ganze Zahl Bestimmt die Nachrichtenmenge für jede Partition, um einen Prüfpunkt zu erledigen. Wird nur wirksam, wenn PARTITION_COUNT Prüfpunktmodus verwendet wird.
spring.cloud.stream.eventhubs.bindings.binding-name.consumer.checkpoint.interval Dauer Legt das Zeitintervall fest, um einen Prüfpunkt zu erledigen. Wird nur wirksam, wenn TIME Prüfpunktmodus verwendet wird.
spring.cloud.stream.eventhubs.bindings.<binding-name>.consumer.batch.max-size Ganze Zahl Die maximale Anzahl von Ereignissen in einem Batch. Erforderlich für den Batch-Consumer-Modus.
spring.cloud.stream.eventhubs.bindings.binding-name.consumer.batch.max-wait-time Dauer Die maximale Dauer für die Batch-Verarbeitung. Wird nur wirksam, wenn der Batch-Consumer-Modus aktiviert ist und optional ist.
spring.cloud.stream.eventhubs.bindings.binding-name.consumer.load-balancing.update-interval Dauer Die Intervallzeitdauer für die Aktualisierung.
spring.cloud.stream.eventhubs.bindings.binding-name.consumer.load-balancing.strategy Lastverteilungsstrategie Die Lastenausgleichsstrategie.
spring.cloud.stream.eventhubs.bindings.binding-name.consumer.load-balancing.partition-ownership-expiration-interval Dauer Der Zeitraum, nach dessen Ablauf die Inhaberschaft der Partition erlischt.
spring.cloud.stream.eventhubs.bindings.binding-name.consumer.track-last-enqueued-event-properties Boolesch Gibt an, ob der Ereignisprozessor Informationen zum letzten enqueued-Ereignis auf der zugehörigen Partition anfordern soll und diese Informationen nachverfolgen soll, wenn Ereignisse empfangen werden.
spring.cloud.stream.eventhubs.bindings.binding-name.consumer.prefetch-count Ganze Zahl Die Anzahl, die vom Consumer verwendet wird, um die Anzahl der Ereignisse zu steuern, die der Event Hub-Consumer aktiv empfängt und lokal in die Warteschlange stellt.
spring.cloud.stream.eventhubs.bindings.binding-name.consumer.initial-partition-event-position Zuordnung mit dem Schlüssel als Partitions-ID und Werten von StartPositionProperties Die Zuordnung, die die Ereignisposition enthält, die für jede Partition verwendet werden soll, wenn kein Prüfpunkt für die Partition im Prüfpunktspeicher vorhanden ist. Diese Zuordnung basiert auf der Partitions-ID.

Anmerkung

Die initial-partition-event-position Konfiguration akzeptiert eine map, um die Anfangsposition für jeden Event Hub anzugeben. Daher ist sein Schlüssel die Partitions-ID, und der Wert ist vom Typ StartPositionProperties, das die Eigenschaften Offset, Sequenznummer, Warteschlangenzeitpunkt und die Angabe, ob er inklusiv ist, enthält. Sie können sie z. B. als

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
Erweiterte Benutzerkonfiguration

Die oben genannten Konfigurationen für Verbindung, Prüfpunkt und gemeinsamen Azure SDK-Client unterstützen die Anpassung für jeden Binder-Consumer, die Sie mit dem Präfix spring.cloud.stream.eventhubs.bindings.<binding-name>.consumer. konfigurieren können.

Produzenteneigenschaften

Diese Eigenschaften werden über EventHubsProducerPropertiesverfügbar gemacht.

Anmerkung

Um Wiederholungen zu vermeiden, unterstützt Spring Cloud Azure Stream Binder Event Hubs seit Version 4.17.0 und 5.11.0 das Festlegen von Werten für alle Kanäle im Format von spring.cloud.stream.eventhubs.default.producer.<property>=<value>.

Vom Hersteller konfigurierbare Eigenschaften von spring-cloud-azure-stream-binder-eventhubs:

Eigentum Typ Beschreibung
spring.cloud.stream.eventhubs.bindings.binding-name.producer.sync boolesch Das Schalter-Flag für die Synchronisierung des Produzenten. Wenn dies auf „true“ gesetzt ist, wartet der Produzent nach einem Sendevorgang auf eine Antwort.
spring.cloud.stream.eventhubs.bindings.binding-name.producer.send-timeout lang Die Wartezeit auf eine Antwort nach einem Sendevorgang. Wird nur wirksam, wenn ein Synchronisierungsproduzent aktiviert ist.
Erweiterte Producer-Konfiguration

Die oben gezeigte Konfiguration für die Verbindung und den gemeinsamen Azure SDK-Client unterstützt Anpassungen für jeden Binder-Produzenten, die Sie mit dem Präfix spring.cloud.stream.eventhubs.bindings.<binding-name>.producer. konfigurieren können.

Grundlegende Nutzung

Senden und Empfangen von Nachrichten von/an Event Hubs

  1. Füllen Sie die Konfigurationsoptionen mit Anmeldeinformationen aus.

    • Konfigurieren Sie für Anmeldeinformationen als Verbindungszeichenfolge die folgenden Eigenschaften in Ihrer application.yml Datei:

      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
      

      Anmerkung

      Microsoft empfiehlt, immer den sichersten Authentifizierungsflow zu verwenden, der verfügbar ist. Der in diesem Verfahren beschriebene Authentifizierungsflow, beispielsweise für Datenbanken, Zwischenspeicher, Nachrichten oder KI-Dienste, erfordert ein sehr hohes Maß an Vertrauen in die Anwendung und birgt Risiken, die bei anderen Flows nicht vorhanden sind. Verwenden Sie diesen Flow nur, wenn sicherere Optionen wie verwaltete Identitäten für kennwortlose oder schlüssellose Verbindungen nicht geeignet sind. Bei Vorgängen des lokalen Computers bevorzugen Sie Benutzeridentitäten für kennwortlose oder schlüssellose Verbindungen.

    • Konfigurieren Sie für Anmeldeinformationen als Dienstprinzipal die folgenden Eigenschaften in Ihrer application.yml-Datei:

      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
      

Anmerkung

Die für tenant-id zulässigen Werte sind: common, organizations, consumersoder die Mandanten-ID. Weitere Informationen zu diesen Werten finden Sie im Abschnitt Falschen Endpunkt verwendet (persönliche Konten und Organisationskonten) in Fehler AADSTS50020 – Benutzerkonto des Identitätsanbieters ist im Mandanten nicht vorhanden. Informationen zum Umwandeln Ihrer Single-Tenant-App finden Sie unter Single-Tenant-App in eine Multitenant-App auf Microsoft Entra ID umwandeln.

  • Wenn Sie verwaltete Identitäten als Anmeldeinformationen verwenden, konfigurieren Sie die folgenden Eigenschaften in Ihrer application.yml-Datei:

    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. Definieren Sie Lieferanten und Verbraucher.

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

Partitionierungsunterstützung

Es wird eine PartitionSupplier mit vom Benutzer bereitgestellten Partitionsinformationen erstellt, um die Partitionsinformationen über die zu sendende Nachricht zu konfigurieren. Das folgende Flussdiagramm zeigt den Prozess zur Ermittlung verschiedener Prioritäten für die Partitions-ID und den Schlüssel:

Diagramm, das ein Flussdiagramm des Prozesses zur Unterstützung der Partitionierung zeigt.

Batch-Consumerunterstützung

  1. Stellen Sie die Batchkonfigurationsoptionen bereit, wie im folgenden Beispiel gezeigt:

    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. Definieren Sie Lieferanten und Verbraucher.

    Wenn der Checkpointing-Modus auf BATCH gesetzt ist, können Sie den folgenden Code verwenden, um Nachrichten zu senden und stapelweise zu verarbeiten.

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

    Im Checkpointing-Modus MANUAL können Sie den folgenden Code verwenden, um Nachrichten zu senden und sie in Batches zu empfangen bzw. Checkpoints zu setzen.

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

Anmerkung

Im Batch-Nutzungsmodus ist der Standardinhaltstyp des Spring Cloud Stream-Ordners application/json. Stellen Sie daher sicher, dass die Nachrichtennutzlast am Inhaltstyp ausgerichtet ist. Wenn Sie z. B. den Standardinhaltstyp application/json verwenden, um Nachrichten mit String Nutzlast zu empfangen, sollte die Nutzlast JSON Stringsein, umgeben von doppelten Anführungszeichen für den ursprünglichen String Text. Während es sich beim text/plain-Inhaltstyp direkt um ein String-Objekt handeln kann. Weitere Informationen finden Sie unter Spring Cloud Stream Content Type Negotiation.

Behandeln von Fehlermeldungen

  • Behandeln von Ausgehenden Bindungsfehlermeldungen

    Standardmäßig erstellt Spring Integration einen globalen Fehlerkanal namens errorChannel. Konfigurieren Sie den folgenden Nachrichtenendpunkt, um ausgehende Bindungsfehlermeldungen zu behandeln.

    @ServiceActivator(inputChannel = IntegrationContextUtils.ERROR_CHANNEL_BEAN_NAME)
    public void handleError(ErrorMessage message) {
        LOGGER.error("Handling outbound binding error: " + message);
    }
    
  • Behandeln eingehender Bindungsfehlermeldungen

    Spring Cloud Stream Event Hubs Binder unterstützt eine Lösung zur Behandlung von Fehlern für eingehende Nachrichtenbindungen: Fehlerhandler.

    Fehlerhandler:

    Spring Cloud Stream macht einen Mechanismus verfügbar, mit dem Sie einen benutzerdefinierten Fehlerhandler bereitstellen können, indem sie eine Consumer hinzufügen, die ErrorMessage Instanzen akzeptiert. Weitere Informationen finden Sie unter Behandeln von Fehlermeldungen in der Spring Cloud Stream-Dokumentation.

    • Bindungsstandardfehlerhandler

      Konfigurieren Sie ein einzelnes Consumer Bean, um alle eingehenden Bindungsfehlermeldungen zu verarbeiten. Die folgende Standardfunktion abonniert jeden eingehenden Bindungsfehlerkanal:

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

      Außerdem müssen Sie die eigenschaft spring.cloud.stream.default.error-handler-definition auf den Funktionsnamen festlegen.

    • Bindungsspezifischer Fehlerhandler

      Konfigurieren Sie eine Consumer Bean, um die spezifischen eingehenden Bindungsfehlermeldungen zu nutzen. Die folgende Funktion abonniert den spezifischen Fehlerkanal der eingehenden Bindung und hat eine höhere Priorität als der Standard-Fehlerhandler der Bindung:

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

      Außerdem müssen Sie die eigenschaft spring.cloud.stream.bindings.<input-binding-name>.error-handler-definition auf den Funktionsnamen festlegen.

Event Hubs-Nachrichtenkopfzeilen

Die unterstützten grundlegenden Nachrichtenkopfzeilen finden Sie im Abschnitt Event Hubs-Nachrichtenkopfzeilen Abschnitt Spring Cloud Azure-Unterstützung für Spring Integration.

Unterstützung für mehrere Binder

Die Verbindung mit mehreren Event Hubs-Namespaces wird auch mithilfe mehrerer Ordner unterstützt. In diesem Beispiel wird eine Verbindungszeichenfolge als Beispiel verwendet. Anmeldeinformationen von Dienstprinzipalen und verwalteten Identitäten werden ebenfalls unterstützt. Sie können verwandte Eigenschaften in den Umgebungseinstellungen der einzelnen Ordner festlegen.

  1. Um mehrere Ordner mit Event Hubs zu verwenden, konfigurieren Sie die folgenden Eigenschaften in Ihrer application.yml Datei:

    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
    

    Anmerkung

    In der vorherigen Anwendungsdatei wird gezeigt, wie Sie einen einzelnen Standard-Poller zur Anwendung auf alle Bindungen konfigurieren. Wenn Sie den Poller für eine bestimmte Bindung konfigurieren möchten, können Sie eine Konfiguration wie spring.cloud.stream.bindings.<binding-name>.producer.poller.fixed-delay=3000verwenden.

    Anmerkung

    Microsoft empfiehlt, immer den sichersten Authentifizierungsflow zu verwenden, der verfügbar ist. Der in diesem Verfahren beschriebene Authentifizierungsflow, beispielsweise für Datenbanken, Zwischenspeicher, Nachrichten oder KI-Dienste, erfordert ein sehr hohes Maß an Vertrauen in die Anwendung und birgt Risiken, die bei anderen Flows nicht vorhanden sind. Verwenden Sie diesen Flow nur, wenn sicherere Optionen wie verwaltete Identitäten für kennwortlose oder schlüssellose Verbindungen nicht geeignet sind. Bei Vorgängen des lokalen Computers bevorzugen Sie Benutzeridentitäten für kennwortlose oder schlüssellose Verbindungen.

  2. Wir müssen zwei Lieferanten und zwei Verbraucher definieren:

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

Ressourcenbereitstellung

Der Event Hubs-Binder unterstützt die Bereitstellung eines Event Hubs und einer Consumergruppe. Benutzer können die folgenden Eigenschaften verwenden, um die Bereitstellung zu aktivieren.

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

Anmerkung

Die für tenant-id zulässigen Werte sind: common, organizations, consumersoder die Mandanten-ID. Weitere Informationen zu diesen Werten finden Sie im Abschnitt Falschen Endpunkt verwendet (persönliche Konten und Organisationskonten) in Fehler AADSTS50020 – Benutzerkonto des Identitätsanbieters ist im Mandanten nicht vorhanden. Informationen zum Umwandeln Ihrer Single-Tenant-App finden Sie unter Single-Tenant-App in eine Multitenant-App auf Microsoft Entra ID umwandeln.

Proben

Weitere Informationen finden Sie im azure-spring-boot-samples Repository auf GitHub.

Spring Cloud Stream Binder für Azure Service Bus

Schlüsselkonzepte

Der Spring Cloud Stream Binder für Azure Service Bus stellt die Bindungsimplementierung für das Spring Cloud Stream Framework bereit. Diese Implementierung basiert auf Spring Integration Service Bus Channel-Adaptern.

Geplante Nachricht

Dieser Binder unterstützt das Übermitteln von Nachrichten an ein Topic zur verzögerten Verarbeitung. Benutzer können geplante Nachrichten mit dem Header x-delay senden, der eine Verzögerungszeit für die Nachricht in Millisekunden angibt. Die Nachricht wird nach x-delay Millisekunden an die jeweiligen Themen übermittelt.

Verbrauchergruppe

Service Bus Topic bietet ähnliche Unterstützung für Consumergruppen wie Apache Kafka, allerdings mit einer etwas anderen Logik. Dieser Ordner basiert auf Subscription eines Themas, das als Verbrauchergruppe fungiert.

Einrichtung von Abhängigkeiten

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

Alternativ können Sie auch den Spring Cloud Azure Stream Service Bus Starter verwenden, wie im folgenden Beispiel für Maven gezeigt:

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

Konfiguration

Der Binder stellt die folgenden zwei Bereiche von Konfigurationsoptionen bereit:

Verbindungskonfigurationseigenschaften

Dieser Abschnitt enthält die Konfigurationsoptionen, die zum Herstellen einer Verbindung mit Azure Service Bus verwendet werden.

Anmerkung

Wenn Sie einen Sicherheitsprinzipal zum Authentifizieren und Autorisieren mit Microsoft Entra-ID für den Zugriff auf eine Azure-Ressource verwenden, lesen Sie Autorisieren des Zugriffs mit Microsoft Entra ID, um sicherzustellen, dass dem Sicherheitsprinzipal die ausreichende Berechtigung für den Zugriff auf die Azure-Ressource gewährt wurde.

Konfigurierbare Verbindungseigenschaften von spring-cloud-azure-stream-binder-servicebus:

Eigentum Typ Beschreibung
spring.cloud.azure.servicebus.enabled boolesch Gibt an, ob ein Azure Service Bus aktiviert ist.
spring.cloud.azure.servicebus.connection-string Schnur Service Bus-Namespace-Verbindungszeichenfolgenwert.
spring.cloud.azure.servicebus.custom-endpoint-address Schnur Die benutzerdefinierte Endpunktadresse, die beim Herstellen einer Verbindung mit Service Bus verwendet werden soll.
spring.cloud.azure.servicebus.namespace Schnur Wert des Service Bus-Namespace, der das Präfix der FQDN ist. Ein FQDN sollte aus NamespaceName.DomainName bestehen.
spring.cloud.azure.servicebus.domain-name Schnur Domänenname eines Azure Service Bus-Namespacewerts.

Anmerkung

Allgemeine Konfigurationsoptionen des Azure Service SDK können auch für den Spring Cloud Azure Stream Service Bus Binder konfiguriert werden. Die unterstützten Konfigurationsoptionen werden in Spring Cloud Azure-Konfigurationeingeführt und können entweder mit dem einheitlichen Präfix spring.cloud.azure. oder dem Präfix von spring.cloud.azure.servicebus.konfiguriert werden.

Der Ordner unterstützt auch Spring Could Azure Resource Manager standardmäßig. Informationen dazu, wie Sie die Verbindungszeichenfolge mit Sicherheitsprinzipalen abrufen, denen keine Data-bezogenen Rollen zugewiesen wurden, finden Sie im Abschnitt Grundlegende Verwendung von Spring Cloud Azure Resource Manager.

Azure Service Bus-Bindungskonfigurationseigenschaften

Die folgenden Optionen sind in vier Abschnitte unterteilt: Consumer-Eigenschaften, erweiterte Consumer-Konfigurationen, Producer-Eigenschaften und erweiterte Producer-Konfigurationen.

Verbrauchereigenschaften

Diese Eigenschaften werden über ServiceBusConsumerPropertiesverfügbar gemacht.

Anmerkung

Um Wiederholungen zu vermeiden, unterstützt Spring Cloud Azure Stream Binder Service Bus seit Version 4.17.0 und 5.11.0 das Festlegen von Werten für alle Kanäle im Format von spring.cloud.stream.servicebus.default.consumer.<property>=<value>.

Konfigurierbare Verbrauchereigenschaften von spring-cloud-azure-stream-binder-servicebus:

Eigentum Typ Vorgabe Beschreibung
spring.cloud.stream.servicebus.bindings.binding-name.consumer.requeue-rejected boolesch FALSCH Wenn die fehlgeschlagenen Nachrichten in die DLQ weitergeleitet werden.
spring.cloud.stream.servicebus.bindings.binding-name.consumer.max-concurrent-calls Ganze Zahl 1 Maximale Anzahl gleichzeitig verarbeiteter Nachrichten, die der Service Bus-Prozessorclient verarbeiten soll. Wenn die Sitzung aktiviert ist, gilt sie für jede Sitzung.
spring.cloud.stream.servicebus.bindings.binding-name.consumer.max-concurrent-sessions Ganze Zahl NULL Maximale Anzahl gleichzeitiger Sitzungen, die zu einem bestimmten Zeitpunkt verarbeitet werden sollen.
spring.cloud.stream.servicebus.bindings.binding-name.consumer.session-enabled Boolesch NULL Gibt an, ob die Sitzung aktiviert ist.
spring.cloud.stream.servicebus.bindings.binding-name.consumer.session-idle-timeout Dauer NULL Legt die maximale Zeitspanne (Dauer) fest, die auf den Empfang einer Nachricht für die aktuell aktive Sitzung wartet.
spring.cloud.stream.servicebus.bindings.binding-name.consumer.prefetch-count Ganze Zahl 0 Die Prefetch-Anzahl des Service Bus-Prozessorclients.
spring.cloud.stream.servicebus.bindings.binding-name.consumer.sub-queue SubQueue nichts Der Typ der Unterwarteschlange, zu der eine Verbindung hergestellt werden soll.
spring.cloud.stream.servicebus.bindings.binding-name.consumer.max-auto-lock-renew-duration Dauer 5m Die Zeitspanne, um die automatische Verlängerung der Sperre fortzusetzen.
spring.cloud.stream.servicebus.bindings.binding-name.consumer.receive-mode ServiceBusReceiveMode peek_lock Der Empfangsmodus des Service Bus-Prozessorclients.
spring.cloud.stream.servicebus.bindings.binding-name.consumer.auto-complete Boolesch wahr Ob Nachrichten automatisch abgewickelt werden sollen. Wenn dies auf „false“ gesetzt ist, wird ein Nachrichtenheader mit dem Wert Checkpointer hinzugefügt, damit Entwickler Nachrichten manuell bestätigen können.
spring.cloud.stream.servicebus.bindings.binding-name.consumer.max-size-in-megabytes Lang 1024 Die maximale Größe der Warteschlange/des Themas in Megabyte, bei der es sich um die Größe des arbeitsspeichers handelt, der für die Warteschlange/das Thema zugeordnet ist.
spring.cloud.stream.servicebus.bindings.binding-name.consumer.default-message-time-to-live Dauer P10675199DT2H48M5.4775807S. (10675199 Tage, 2 Stunden, 48 Minuten, 5 Sekunden und 477 Millisekunden) Die Dauer, nach der die Nachricht abläuft, beginnend ab dem Zeitpunkt, an den die Nachricht an Service Bus gesendet wird.

Wichtig

Wenn Sie den Azure Resource Manager (ARM) verwenden, müssen Sie die eigenschaft spring.cloud.stream.servicebus.bindings.<binding-name>.consume.entity-type konfigurieren. Weitere Informationen finden Sie im servicebus-queue-binder-arm Beispiel auf GitHub.

Erweiterte Benutzerkonfiguration

Die obige Verbindungs- und allgemeine Azure SDK-Client-Konfiguration unterstützt die Anpassung für jeden Binder-Consumer, die Sie mit dem Präfix spring.cloud.stream.servicebus.bindings.<binding-name>.consumer. konfigurieren können.

Produzenteneigenschaften

Diese Eigenschaften werden über ServiceBusProducerPropertiesverfügbar gemacht.

Anmerkung

Um Wiederholungen zu vermeiden, unterstützt Spring Cloud Azure Stream Binder Service Bus seit Version 4.17.0 und 5.11.0 das Festlegen von Werten für alle Kanäle im Format von spring.cloud.stream.servicebus.default.producer.<property>=<value>.

Vom Hersteller konfigurierbare Eigenschaften von spring-cloud-azure-stream-binder-servicebus:

Eigentum Typ Vorgabe Beschreibung
spring.cloud.stream.servicebus.bindings.binding-name.producer.sync boolesch FALSCH Switch-Flag für Synchronisation des Producers.
spring.cloud.stream.servicebus.bindings.binding-name.producer.send-timeout lang 10.000 Timeoutwert für das Senden des Produzenten.
spring.cloud.stream.servicebus.bindings.binding-name.producer.entity-type ServiceBusEntityType NULL Service Bus-Entitätstyp des Produzenten, erforderlich für den Bindungshersteller.
spring.cloud.stream.servicebus.bindings.binding-name.producer.max-size-in-megabytes Lang 1024 Die maximale Größe der Warteschlange/des Themas in Megabyte, bei der es sich um die Größe des arbeitsspeichers handelt, der für die Warteschlange/das Thema zugeordnet ist.
spring.cloud.stream.servicebus.bindings.binding-name.producer.default-message-time-to-live Dauer P10675199DT2H48M5.4775807S. (10675199 Tage, 2 Stunden, 48 Minuten, 5 Sekunden und 477 Millisekunden) Die Dauer, nach der die Nachricht abläuft, beginnend ab dem Zeitpunkt, an den die Nachricht an Service Bus gesendet wird.

Wichtig

Bei Verwendung des Bindungsherstellers muss die Eigenschaft spring.cloud.stream.servicebus.bindings.<binding-name>.producer.entity-type konfiguriert werden.

Erweiterte Producer-Konfiguration

Die oben gezeigte Konfiguration für die Verbindung und den gemeinsamen Azure SDK-Client unterstützt Anpassungen für jeden Binder-Produzenten, die Sie mit dem Präfix spring.cloud.stream.servicebus.bindings.<binding-name>.producer. konfigurieren können.

Grundlegende Nutzung

Senden und Empfangen von Nachrichten von/an Service Bus

  1. Füllen Sie die Konfigurationsoptionen mit Anmeldeinformationen aus.

    • Konfigurieren Sie für Anmeldeinformationen als Verbindungszeichenfolge die folgenden Eigenschaften in Ihrer application.yml Datei:

      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
      

      Anmerkung

      Microsoft empfiehlt, immer den sichersten Authentifizierungsflow zu verwenden, der verfügbar ist. Der in diesem Verfahren beschriebene Authentifizierungsflow, beispielsweise für Datenbanken, Zwischenspeicher, Nachrichten oder KI-Dienste, erfordert ein sehr hohes Maß an Vertrauen in die Anwendung und birgt Risiken, die bei anderen Flows nicht vorhanden sind. Verwenden Sie diesen Flow nur, wenn sicherere Optionen wie verwaltete Identitäten für kennwortlose oder schlüssellose Verbindungen nicht geeignet sind. Bei Vorgängen des lokalen Computers bevorzugen Sie Benutzeridentitäten für kennwortlose oder schlüssellose Verbindungen.

    • Konfigurieren Sie für Anmeldeinformationen als Dienstprinzipal die folgenden Eigenschaften in Ihrer application.yml-Datei:

      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
      

Anmerkung

Die für tenant-id zulässigen Werte sind: common, organizations, consumersoder die Mandanten-ID. Weitere Informationen zu diesen Werten finden Sie im Abschnitt Falschen Endpunkt verwendet (persönliche Konten und Organisationskonten) in Fehler AADSTS50020 – Benutzerkonto des Identitätsanbieters ist im Mandanten nicht vorhanden. Informationen zum Umwandeln Ihrer Single-Tenant-App finden Sie unter Single-Tenant-App in eine Multitenant-App auf Microsoft Entra ID umwandeln.

  • Wenn Sie verwaltete Identitäten als Anmeldeinformationen verwenden, konfigurieren Sie die folgenden Eigenschaften in Ihrer application.yml-Datei:

    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. Definieren Sie Lieferanten und Verbraucher.

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

Unterstützung für Partitionsschlüssel

Der Ordner unterstützt Service Bus-Partitionierung, indem die Einstellung des Partitionsschlüssels und der Sitzungs-ID im Nachrichtenkopf zugelassen wird. In diesem Abschnitt wird erläutert, wie Sie den Partitionsschlüssel für Nachrichten festlegen.

Spring Cloud Stream stellt eine SpEL-Ausdruckseigenschaft des Partitionsschlüssels spring.cloud.stream.bindings.<binding-name>.producer.partition-key-expressionbereit. Legen Sie diese Eigenschaft beispielsweise als "'partitionKey-' + headers[<message-header-key>]" fest, und fügen Sie eine Kopfzeile namens "Message-header-key" hinzu. Spring Cloud Stream verwendet den Wert für diesen Header beim Auswerten des Ausdrucks, um einen Partitionsschlüssel zuzuweisen. Der folgende Code stellt einen Beispielhersteller bereit:

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

Sitzungsunterstützung

Der Binder unterstützt Nachrichtensitzungen von Service Bus. Die Sitzungs-ID einer Nachricht kann über den Nachrichtenkopf festgelegt werden.

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

Anmerkung

Gemäß Dienstbuspartitionierunghat die Sitzungs-ID eine höhere Priorität als Partitionsschlüssel. Wenn also sowohl ServiceBusMessageHeaders#SESSION_ID- als auch ServiceBusMessageHeaders#PARTITION_KEY Header festgelegt werden, wird der Wert der Sitzungs-ID schließlich verwendet, um den Wert des Partitionsschlüssels zu überschreiben.

Behandeln von Fehlermeldungen

  • Behandeln von Ausgehenden Bindungsfehlermeldungen

    Standardmäßig erstellt Spring Integration einen globalen Fehlerkanal namens errorChannel. Konfigurieren Sie den folgenden Nachrichtenendpunkt so, dass er die Fehlermeldung zur ausgehenden Bindung verarbeitet.

    @ServiceActivator(inputChannel = IntegrationContextUtils.ERROR_CHANNEL_BEAN_NAME)
    public void handleError(ErrorMessage message) {
        LOGGER.error("Handling outbound binding error: " + message);
    }
    
  • Behandeln eingehender Bindungsfehlermeldungen

    Spring Cloud Stream Service Bus Binder unterstützt zwei Lösungen zur Behandlung von Fehlern für die eingehenden Nachrichtenbindungen: den Binder-Fehlerhandler und Handler.

    Binder-Fehlerhandler:

    Der Standardordnerfehlerhandler behandelt die eingehende Bindung. Sie verwenden diesen Handler, um fehlgeschlagene Nachrichten an die Dead-Letter-Warteschlange zu senden, wenn spring.cloud.stream.servicebus.bindings.<binding-name>.consumer.requeue-rejected aktiviert ist. Andernfalls werden die fehlgeschlagenen Nachrichten verworfen. Der Binder-Fehlerhandler ist mit anderen verfügbaren Fehlerhandlern gegenseitig ausschließend.

    Fehlerhandler:

    Spring Cloud Stream macht einen Mechanismus verfügbar, mit dem Sie einen benutzerdefinierten Fehlerhandler bereitstellen können, indem sie eine Consumer hinzufügen, die ErrorMessage Instanzen akzeptiert. Weitere Informationen finden Sie unter Behandeln von Fehlermeldungen in der Spring Cloud Stream-Dokumentation.

    • Bindungsstandardfehlerhandler

      Konfigurieren Sie ein einzelnes Consumer Bean, um alle eingehenden Bindungsfehlermeldungen zu verarbeiten. Die folgende Standardfunktion abonniert jeden eingehenden Bindungsfehlerkanal:

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

      Außerdem müssen Sie die eigenschaft spring.cloud.stream.default.error-handler-definition auf den Funktionsnamen festlegen.

    • Bindungsspezifischer Fehlerhandler

      Konfigurieren Sie eine Consumer Bean, um die spezifischen eingehenden Bindungsfehlermeldungen zu nutzen. Die folgende Funktion abonniert den spezifischen Eingehenden Bindungsfehlerkanal mit einer höheren Priorität als der Bindungsstandardfehlerhandler.

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

      Außerdem müssen Sie die eigenschaft spring.cloud.stream.bindings.<input-binding-name>.error-handler-definition auf den Funktionsnamen festlegen.

Dienstbus-Nachrichtenkopfzeilen

Die unterstützten grundlegenden Nachrichtenkopfzeilen finden Sie im Abschnitt Service Bus-Nachrichtenkopfzeilen Abschnitt Spring Cloud Azure-Unterstützung für Spring Integration.

Anmerkung

Beim Festlegen des Partitionsschlüssels ist die Priorität des Nachrichtenkopfs höher als die Spring Cloud Stream-Eigenschaft. Daher wird spring.cloud.stream.bindings.<binding-name>.producer.partition-key-expression nur wirksam, wenn keines der ServiceBusMessageHeaders#SESSION_ID und ServiceBusMessageHeaders#PARTITION_KEY Header konfiguriert ist.

Unterstützung für mehrere Binder

Die Verbindung mit mehreren ServiceBus-Namespaces wird auch mithilfe mehrerer Ordner unterstützt. In diesem Beispiel wird die Verbindungszeichenfolge als Beispiel verwendet. Anmeldeinformationen von Dienstprinzipalen und verwalteten Identitäten werden ebenfalls unterstützt, und Benutzer können die entsprechenden Eigenschaften in den Umgebungseinstellungen der einzelnen Binder festlegen.

  1. Um mehrere Ordner von ServiceBus zu verwenden, konfigurieren Sie die folgenden Eigenschaften in Ihrer application.yml Datei:

    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
    

    Anmerkung

    In der vorherigen Anwendungsdatei wird gezeigt, wie Sie einen einzelnen Standard-Poller zur Anwendung auf alle Bindungen konfigurieren. Wenn Sie den Poller für eine bestimmte Bindung konfigurieren möchten, können Sie eine Konfiguration wie spring.cloud.stream.bindings.<binding-name>.producer.poller.fixed-delay=3000verwenden.

    Anmerkung

    Microsoft empfiehlt, immer den sichersten Authentifizierungsflow zu verwenden, der verfügbar ist. Der in diesem Verfahren beschriebene Authentifizierungsflow, beispielsweise für Datenbanken, Zwischenspeicher, Nachrichten oder KI-Dienste, erfordert ein sehr hohes Maß an Vertrauen in die Anwendung und birgt Risiken, die bei anderen Flows nicht vorhanden sind. Verwenden Sie diesen Flow nur, wenn sicherere Optionen wie verwaltete Identitäten für kennwortlose oder schlüssellose Verbindungen nicht geeignet sind. Bei Vorgängen des lokalen Computers bevorzugen Sie Benutzeridentitäten für kennwortlose oder schlüssellose Verbindungen.

  2. wir müssen zwei Lieferanten und zwei Verbraucher definieren

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

Ressourcenbereitstellung

Der Servicebusordner unterstützt die Bereitstellung von Warteschlangen, Themen und Abonnements. Benutzer können die folgenden Eigenschaften verwenden, um die Bereitstellung zu aktivieren.

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}

Anmerkung

Die für tenant-id zulässigen Werte sind: common, organizations, consumersoder die Mandanten-ID. Weitere Informationen zu diesen Werten finden Sie im Abschnitt Falschen Endpunkt verwendet (persönliche Konten und Organisationskonten) in Fehler AADSTS50020 – Benutzerkonto des Identitätsanbieters ist im Mandanten nicht vorhanden. Informationen zum Umwandeln Ihrer Single-Tenant-App finden Sie unter Single-Tenant-App in eine Multitenant-App auf Microsoft Entra ID umwandeln.

Anpassen von ServiceBus-Clienteigenschaften

Entwickler können AzureServiceClientBuilderCustomizer zum Anpassen von Service Bus-Clienteigenschaften verwenden. Im folgenden Beispiel wird die eigenschaft sessionIdleTimeout in ServiceBusClientBuilderangepasst:

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

Proben

Weitere Informationen finden Sie im azure-spring-boot-samples Repository auf GitHub.