Spring Cloud Azure 對 Spring Integration 的支援

適用於 Azure 的 Spring Integration Extension 為 Azure SDK for Java 提供的各種服務提供 Spring Integration 配接器。 我們提供這些 Azure 服務的 Spring Integration 支援:事件中樞、服務總線、記憶體佇列。 以下是支援的配接器清單:

Spring 與 Azure 事件中樞整合

重要概念

Azure 事件中樞是巨量數據串流平臺和事件擷取服務。 它可以每秒接收和處理數百萬個事件。 傳送至事件中樞的數據可以使用任何即時分析提供者或批處理/記憶體配接器來轉換和儲存。

Spring Integration 可在 Spring 型應用程式中啟用輕量型傳訊,並支援透過宣告式配接器與外部系統整合。 這些配接器在 Spring 對遠端、訊息傳遞和排程的支援之上,提供了更高層級的抽象。 Spring Integration for Event Hubs 擴充專案為 Azure 事件中樞提供輸入與輸出通道配接器及閘道。

注意

RxJava 支援 API 會從 4.0.0 版卸除。 如需詳細資訊,請參閱 Javadoc。

取用者群組

事件中樞提供與 Apache Kafka 類似的取用者群組支援,但邏輯稍有不同。 雖然 Kafka 會將所有已提交的位移儲存在訊息代理中,但您必須手動儲存正在處理的 Event Hubs 訊息位移。 事件中樞 SDK 提供函式,以將這類位移儲存在 Azure 記憶體內。

資料分割支援

事件中樞提供與 Kafka 類似的實體分割區概念。 但不同於 Kafka 在取用者與分割區之間的自動重新平衡,Event Hubs 採用一種預先搶占模式。 儲存體帳戶可作為租約,用來判斷哪個分割區由哪個取用者擁有。 當新的消費者啟動時,它會嘗試從負載最重的消費者中竊取一些分割區,以達成負載平衡。

若要指定負載平衡策略,開發人員可以使用 EventHubsContainerProperties 來進行設定。 如需設定 的範例,請參閱以下章節。

批次取用端支援

EventHubsInboundChannelAdapter 支援批次消費模式。 若要啟用,用戶可以在建構 ListenerMode.BATCH 實例時,將接聽程式模式指定為 EventHubsInboundChannelAdapter。 啟用時,將會接收一則 訊息,其承載內容為批次事件清單,並將其傳遞至下游通道。 每個訊息標頭也會轉換成清單,其中內容是從每個事件剖析的相關聯標頭值。 對於分割區識別碼、checkpointer 和最後排入佇列屬性的共用標頭,若整個事件批次共用相同的值,則會顯示為單一值。 如需詳細資訊,請參閱 事件中樞訊息標頭 一節。

注意

檢查點標頭僅在使用 MANUAL 檢查點模式時才會存在。

批次取用程式的檢查點功能支援兩種模式:BATCH 和 MANUAL。 BATCH 模式是一種自動建立檢查點的模式,會在收到事件後,為整個事件批次一併建立檢查點。 MANUAL 模式用於讓使用者為事件建立檢查點。 使用時,Checkpointer 會傳遞至訊息標頭,而且使用者可以使用它來執行檢查點。

批次取用原則可透過 max-size 和 max-wait-time 的屬性指定,其中 max-size 為必要屬性,而 max-wait-time 為選用屬性。 若要指定批次取用策略,開發人員可以使用 EventHubsContainerProperties 來進行設定。 如需設定 的範例,請參閱 以下章節。

相依性設定

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

配置

此入門套件提供以下 3 部分的組態選項:

聯機組態屬性

本節包含用來連線到 Azure 事件中樞的組態選項。

注意

如果您選擇使用安全性主體向 Microsoft Entra ID 進行驗證和授權,以存取 Azure 資源,請參閱 使用 Microsoft Entra ID 授權存取,以確保安全性主體已獲得存取 Azure 資源的足夠許可權。

spring-cloud-azure-starter-integration-eventhubs 的連線可設定屬性:

財產 類型 描述
spring.cloud.azure.eventhubs.enabled 布爾 是否啟用 Azure 事件中樞。
spring.cloud.azure.eventhubs.connection-string 字串 事件中樞命名空間連接字串值。
spring.cloud.azure.eventhubs.namespace 字串 Event Hubs 命名空間值,也就是 FQDN 的前綴。 FQDN 應該由 NamespaceName.DomainName 組成
spring.cloud.azure.eventhubs.domain-name 字串 Azure 事件中樞 命名空間值所對應的網域名稱。
spring.cloud.azure.eventhubs.custom-endpoint-address 字串 自訂端點位址。
spring.cloud.azure.eventhubs.shared-connection 布林值 基礎 EventProcessorClient 和 EventHubProducerAsyncClient 是否使用相同的連線。 根據預設,系統會針對每個建立的事件中樞用戶端建構及使用新的連線。

檢查點組態屬性

本節包含記憶體 Blob 服務的組態選項,用於保存分割區擁有權和檢查點資訊。

注意

從版本 4.0.0 開始,當 的 spring.cloud.azure.eventhubs.processor.checkpoint-store.create-container-if-not-exists 屬性未手動啟用時,不會自動建立儲存容器。

spring-cloud-azure-starter-integration-eventhubs 的檢查點可設定屬性:

財產 類型 描述
spring.cloud.azure.eventhubs.processor.checkpoint-store.create-container-if-not-exists 布林值 如果不存在,是否允許建立容器。
spring.cloud.azure.eventhubs.processor.checkpoint-store.account-name 字串 記憶體帳戶的名稱。
spring.cloud.azure.eventhubs.processor.checkpoint-store.account-key 字串 儲存體帳戶存取金鑰。
spring.cloud.azure.eventhubs.processor.checkpoint-store.container-name 字串 記憶體容器名稱。

常見的 Azure 服務 SDK 組態選項也可以針對記憶體 Blob 檢查點存放區進行設定。 支援的組態選項已於 Spring Cloud Azure 組態中介紹,且可使用統一前綴 spring.cloud.azure. 或 spring.cloud.azure.eventhubs.processor.checkpoint-store. 前綴進行設定。

事件中樞處理器組態屬性

EventHubsInboundChannelAdapter 會使用 EventProcessorClient 從事件中樞取用訊息,來設定 EventProcessorClient的整體屬性,開發人員可以使用 EventHubsContainerProperties 進行設定。 請參閱下一節 如何使用 EventHubsInboundChannelAdapter。

基本用法

將訊息傳送至 Azure 事件中樞

  1. 填入認證組態選項。

    • 針對認證作為連接字串,請在 application.yml 檔案中設定下列屬性:

      spring:
        cloud:
          azure:
            eventhubs:
              connection-string: ${AZURE_EVENT_HUBS_CONNECTION_STRING}
              processor:
                checkpoint-store:
                  container-name: ${CHECKPOINT-CONTAINER}
                  account-name: ${CHECKPOINT-STORAGE-ACCOUNT}
                  account-key: ${CHECKPOINT-ACCESS-KEY}
      

      注意

      Microsoft 建議您使用最安全的可用驗證流程。 此程式中所述的驗證流程,例如資料庫、快取、傳訊或 AI 服務,在應用程式中需要高度的信任,而且不會在其他流程中帶來風險。 只有在更安全的選項(例如使用受控識別進行無密碼或無金鑰連線)不可行時,才使用此流程。 針對本機計算機作業,偏好使用無密碼或無密鑰連線的使用者身分識別。

    • 針對認證作為受控識別,請在 application.yml 檔案中設定下列屬性:

      spring:
        cloud:
          azure:
            credential:
              managed-identity-enabled: true
              client-id: ${AZURE_CLIENT_ID}
            eventhubs:
              namespace: ${AZURE_EVENT_HUBS_NAMESPACE}
              processor:
                checkpoint-store:
                  container-name: ${CONTAINER_NAME}
                  account-name: ${ACCOUNT_NAME}
      
    • 若使用服務主體認證,請在 application.yml 檔案中設定下列屬性:

      spring:
        cloud:
          azure:
            credential:
              client-id: ${AZURE_CLIENT_ID}
              client-secret: ${AZURE_CLIENT_SECRET}
            profile:
              tenant-id: <tenant>
            eventhubs:
              namespace: ${AZURE_EVENT_HUBS_NAMESPACE}
              processor:
                checkpoint-store:
                  container-name: ${CONTAINER_NAME}
                  account-name: ${ACCOUNT_NAME}
      

注意

tenant-id 允許的值包括:common、organizations、consumers或租用戶標識碼。 如需這些值的詳細資訊,請參閱 錯誤 AADSTS50020 - 來自身分識別提供者的使用者帳戶不存在於租用戶中中的 使用了錯誤的端點(個人和組織帳戶)一節。 如需有關將單一租用戶應用程式轉換為多租用戶的資訊,請參閱 在 Microsoft Entra ID 上將單一租用戶應用程式轉換為多租用戶。

  1. 使用 EventHubsTemplate Bean 建立 DefaultMessageHandler,以將訊息傳送至 Event Hubs。

    class Demo {
        private static final String OUTPUT_CHANNEL = "output";
        private static final String EVENTHUB_NAME = "eh1";
    
        @Bean
        @ServiceActivator(inputChannel = OUTPUT_CHANNEL)
        public MessageHandler messageSender(EventHubsTemplate eventHubsTemplate) {
            DefaultMessageHandler handler = new DefaultMessageHandler(EVENTHUB_NAME, eventHubsTemplate);
            handler.setSendCallback(new ListenableFutureCallback<Void>() {
                @Override
                public void onSuccess(Void result) {
                    LOGGER.info("Message was sent successfully.");
                }
                @Override
                public void onFailure(Throwable ex) {
                    LOGGER.error("There was an error sending the message.", ex);
                }
            });
            return handler;
        }
    }
    
  2. 透過訊息通道建立具有上述訊息處理程式的訊息網關係結。

    class Demo {
        @Autowired
        EventHubOutboundGateway messagingGateway;
    
        @MessagingGateway(defaultRequestChannel = OUTPUT_CHANNEL)
        public interface EventHubOutboundGateway {
            void send(String text);
        }
    }
    
  3. 使用閘道傳送訊息。

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

從 Azure 事件中樞接收訊息

  1. 填入認證組態選項。

  2. 建立訊息通道的 Bean 做為輸入通道。

    @Configuration
    class Demo {
        @Bean
        public MessageChannel input() {
            return new DirectChannel();
        }
    }
    
  3. 使用 EventHubsMessageListenerContainer bean 建立 EventHubsInboundChannelAdapter,以接收來自事件中樞的訊息。

    @Configuration
    class Demo {
        private static final String INPUT_CHANNEL = "input";
        private static final String EVENTHUB_NAME = "eh1";
        private static final String CONSUMER_GROUP = "$Default";
    
        @Bean
        public EventHubsInboundChannelAdapter messageChannelAdapter(
                @Qualifier(INPUT_CHANNEL) MessageChannel inputChannel,
                EventHubsMessageListenerContainer listenerContainer) {
            EventHubsInboundChannelAdapter adapter = new EventHubsInboundChannelAdapter(listenerContainer);
            adapter.setOutputChannel(inputChannel);
            return adapter;
        }
    
        @Bean
        public EventHubsMessageListenerContainer messageListenerContainer(EventHubsProcessorFactory processorFactory) {
            EventHubsContainerProperties containerProperties = new EventHubsContainerProperties();
            containerProperties.setEventHubName(EVENTHUB_NAME);
            containerProperties.setConsumerGroup(CONSUMER_GROUP);
            containerProperties.setCheckpointConfig(new CheckpointConfig(CheckpointMode.MANUAL));
            return new EventHubsMessageListenerContainer(processorFactory, containerProperties);
        }
    }
    
  4. 透過之前建立的訊息通道,使用 EventHubsInboundChannelAdapter 建立訊息接收者系結。

    class Demo {
        @ServiceActivator(inputChannel = INPUT_CHANNEL)
        public void messageReceiver(byte[] payload, @Header(AzureHeaders.CHECKPOINTER) Checkpointer checkpointer) {
            String message = new String(payload);
            LOGGER.info("New message received: '{}'", message);
            checkpointer.success()
                    .doOnSuccess(s -> LOGGER.info("Message '{}' successfully checkpointed", message))
                    .doOnError(e -> LOGGER.error("Error found", e))
                    .block();
        }
    }
    

設定 EventHubsMessageConverter 以自定義 objectMapper

EventHubsMessageConverter 被設計為可配置的 Bean,讓使用者能夠自訂 ObjectMapper。

批次取用端支援

若要從 Event Hubs 以批次方式取用訊息,作法與上述範例類似;此外,使用者應在 EventHubsInboundChannelAdapter 中設定與批次取用相關的組態選項。

建立 EventHubsInboundChannelAdapter 時,接聽模式應設為 BATCH。 建立 EventHubsMessageListenerContainer的 bean 時,請將檢查點模式設定為 MANUAL 或 BATCH,並視需要設定批次選項。

@Configuration
class Demo {
    private static final String INPUT_CHANNEL = "input";
    private static final String EVENTHUB_NAME = "eh1";
    private static final String CONSUMER_GROUP = "$Default";

    @Bean
    public EventHubsInboundChannelAdapter messageChannelAdapter(
            @Qualifier(INPUT_CHANNEL) MessageChannel inputChannel,
            EventHubsMessageListenerContainer listenerContainer) {
        EventHubsInboundChannelAdapter adapter = new EventHubsInboundChannelAdapter(listenerContainer, ListenerMode.BATCH);
        adapter.setOutputChannel(inputChannel);
        return adapter;
    }

    @Bean
    public EventHubsMessageListenerContainer messageListenerContainer(EventHubsProcessorFactory processorFactory) {
        EventHubsContainerProperties containerProperties = new EventHubsContainerProperties();
        containerProperties.setEventHubName(EVENTHUB_NAME);
        containerProperties.setConsumerGroup(CONSUMER_GROUP);
        containerProperties.getBatch().setMaxSize(100);
        containerProperties.setCheckpointConfig(new CheckpointConfig(CheckpointMode.MANUAL));
        return new EventHubsMessageListenerContainer(processorFactory, containerProperties);
    }
}

Event Hubs 訊息標頭

下表說明事件中樞訊息屬性如何對應至 Spring 訊息標頭。 針對 Azure 事件中樞,訊息稱為 event。

在記錄接聽器模式中,Event Hubs 訊息 / 事件屬性與 Spring 訊息標頭之間的對應:

事件中樞的事件屬性 Spring 訊息標頭常數 類型 描述
佇列加入時間 EventHubsHeaders#ENQUEUED_TIME 瞬間 事件在 Event Hub 分割區中排入佇列的時間點,以 UTC 表示。
抵消 EventHubsHeaders#OFFSET 長 從關聯的 Event Hub 分割區接收事件時的位移。
分割區索引鍵 AzureHeaders#PARTITION_KEY 字串 如果在最初發佈事件時設定分割區哈希索引鍵,
分割區識別碼 AzureHeaders#RAW_PARTITION_ID 字串 事件中樞的分割區識別碼。
序號 EventHubsHeaders#SEQUENCE_NUMBER 長 事件在相關事件中樞分割區中排入佇列時所指派的序號。
最後排入佇列的事件屬性 EventHubsHeaders#LAST_ENQUEUED_EVENT_PROPERTIES 最後加入佇列的事件屬性 此分割區中最後一個加入佇列事件的屬性。
NA AzureHeaders#CHECKPOINTER 檢查點管理器 檢查點特定訊息的標頭。

使用者可以剖析訊息標頭,以取得每個事件的相關信息。 若要設定事件的訊息標頭,所有自定義標頭都會放置為事件的應用程式屬性,其中標頭會設定為屬性索引鍵。 從事件中樞接收事件時,所有應用程式屬性都會轉換成訊息標頭。

注意

不支援手動設定分割區索引鍵、排入佇列時間、位移與序號等訊息標頭。

啟用批次取用者模式時,批次訊息的特定標頭如下所列,其中包含每個個別 Event Hubs 事件的值清單。

在批次接聽器模式下,Event Hubs 訊息/事件屬性與 Spring 訊息標頭之間的對應關係:

事件中樞的事件屬性 Spring Batch 訊息標頭常數 類型 描述
加入佇列的時間 EventHubsHeaders#ENQUEUED_TIME 即時清單 每個事件在 Event Hub 分割區中進入佇列時的 UTC 時間點清單。
抵消 EventHubsHeaders#OFFSET Long 的清單 每個事件從對應的事件中樞分割區接收時的位移清單。
分割區索引鍵 AzureHeaders#PARTITION_KEY 字串清單 若在最初發佈各事件時已設定,則列出其分割區雜湊鍵清單。
序號 EventHubsHeaders#SEQUENCE_NUMBER Long 的清單 每個事件在排入對應的 Event Hub 分割區時所指派的各序號清單。
系統屬性 EventHubsHeaders#BATCH_CONVERTED_SYSTEM_PROPERTIES 地圖清單 每個事件的系統屬性清單。
應用程式屬性 EventHubsHeaders#BATCH_CONVERTED_APPLICATION_PROPERTIES 地圖清單 每個事件的應用程式屬性清單,其中會放置所有自定義訊息標頭或事件屬性。

注意

發佈訊息時,如果存在,則會從訊息中移除上述所有批次標頭。

樣品

欲了解更多資訊,請參閱 azure-spring-boot-samples GitHub 上的倉庫。

Spring 與 Azure 服務總線整合

重要概念

Spring Integration 可在 Spring 型應用程式中啟用輕量型傳訊,並支援透過宣告式配接器與外部系統整合。

Azure 服務總線延伸模組專案的 Spring Integration 提供 Azure 服務總線的輸入和輸出通道配接器。

注意

CompletableFuture 支援 API 已從 2.10.0 版淘汰,並由 4.0.0 版的 Reactor Core 取代。 如需詳細資訊,請參閱 Javadoc。

相依性設定

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

配置

此入門範本提供下列兩部分的設定選項:

聯機組態屬性

本節包含用來連線到 Azure 服務總線的組態選項。

注意

如果您選擇使用安全性主體向 Microsoft Entra ID 進行驗證和授權,以存取 Azure 資源,請參閱 使用 Microsoft Entra ID 授權存取,以確保安全性主體已獲得存取 Azure 資源的足夠許可權。

spring-cloud-azure-starter-integration-servicebus 的連線可設定屬性:

財產 類型 描述
spring.cloud.azure.servicebus.enabled 布爾 是否啟用 Azure 服務總線。
spring.cloud.azure.servicebus.connection-string 字串 服務總線命名空間連接字串值。
spring.cloud.azure.servicebus.custom-endpoint-address 字串 連接到服務總線時要使用的自訂端點位址。
spring.cloud.azure.servicebus.namespace 字串 服務總線命名空間值,這是 FQDN 的前置詞。 FQDN 應該由 NamespaceName.DomainName 組成
spring.cloud.azure.servicebus.domain-name 字串 Azure 服務總線命名空間值的功能變數名稱。

服務總線處理器組態屬性

ServiceBusInboundChannelAdapter 會使用 ServiceBusProcessorClient 來取用訊息,來設定 ServiceBusProcessorClient的整體屬性,開發人員可以使用 ServiceBusContainerProperties 來進行設定。 請參閱下一節 如何使用 ServiceBusInboundChannelAdapter。

基本用法

將訊息傳送至 Azure 服務總線

  1. 填入認證組態選項。

    • 針對認證作為連接字串,請在 application.yml 檔案中設定下列屬性:

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

      注意

      Microsoft 建議您使用最安全的可用驗證流程。 此程式中所述的驗證流程,例如資料庫、快取、傳訊或 AI 服務,在應用程式中需要高度的信任,而且不會在其他流程中帶來風險。 只有在更安全的選項(例如使用受控識別進行無密碼或無金鑰連線)不可行時,才使用此流程。 針對本機計算機作業,偏好使用無密碼或無密鑰連線的使用者身分識別。

    • 針對認證作為受控識別,請在 application.yml 檔案中設定下列屬性:

      spring:
        cloud:
          azure:
            credential:
              managed-identity-enabled: true
              client-id: ${AZURE_CLIENT_ID}
            profile:
              tenant-id: <tenant>
            servicebus:
              namespace: ${AZURE_SERVICE_BUS_NAMESPACE}
      

注意

tenant-id 允許的值包括:common、organizations、consumers或租用戶標識碼。 如需這些值的詳細資訊,請參閱 錯誤 AADSTS50020 - 來自身分識別提供者的使用者帳戶不存在於租用戶中中的 使用了錯誤的端點(個人和組織帳戶)一節。 如需有關將單一租用戶應用程式轉換為多租用戶的資訊,請參閱 在 Microsoft Entra ID 上將單一租用戶應用程式轉換為多租用戶。

  • 若使用服務主體認證,請在 application.yml 檔案中設定下列屬性:

    spring:
      cloud:
        azure:
          credential:
            client-id: ${AZURE_CLIENT_ID}
            client-secret: ${AZURE_CLIENT_SECRET}
          profile:
            tenant-id: <tenant>
          servicebus:
            namespace: ${AZURE_SERVICE_BUS_NAMESPACE}
    

注意

tenant-id 允許的值包括:common、organizations、consumers或租用戶標識碼。 如需這些值的詳細資訊,請參閱 錯誤 AADSTS50020 - 來自身分識別提供者的使用者帳戶不存在於租用戶中中的 使用了錯誤的端點(個人和組織帳戶)一節。 如需有關將單一租用戶應用程式轉換為多租用戶的資訊,請參閱 在 Microsoft Entra ID 上將單一租用戶應用程式轉換為多租用戶。

  1. 使用 ServiceBusTemplate bean 建立 DefaultMessageHandler 以將訊息傳送至 服務匯流排,並為 ServiceBusTemplate 設定實體類型。 此範例會採用服務總線佇列作為範例。

    class Demo {
        private static final String OUTPUT_CHANNEL = "queue.output";
    
        @Bean
        @ServiceActivator(inputChannel = OUTPUT_CHANNEL)
        public MessageHandler queueMessageSender(ServiceBusTemplate serviceBusTemplate) {
            serviceBusTemplate.setDefaultEntityType(ServiceBusEntityType.QUEUE);
            DefaultMessageHandler handler = new DefaultMessageHandler(QUEUE_NAME, serviceBusTemplate);
            handler.setSendCallback(new ListenableFutureCallback<Void>() {
                @Override
                public void onSuccess(Void result) {
                    LOGGER.info("Message was sent successfully.");
                }
    
                @Override
                public void onFailure(Throwable ex) {
                    LOGGER.error("There was an error sending the message.", ex);
                }
            });
    
            return handler;
        }
    }
    
  2. 透過訊息通道建立具有上述訊息處理程式的訊息網關係結。

    class Demo {
        @Autowired
        QueueOutboundGateway messagingGateway;
    
        @MessagingGateway(defaultRequestChannel = OUTPUT_CHANNEL)
        public interface QueueOutboundGateway {
            void send(String text);
        }
    }
    
  3. 使用閘道傳送訊息。

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

從 Azure 服務總線接收訊息

  1. 填入認證組態選項。

  2. 建立訊息通道的 Bean 做為輸入通道。

    @Configuration
    class Demo {
        private static final String INPUT_CHANNEL = "input";
    
        @Bean
        public MessageChannel input() {
            return new DirectChannel();
        }
    }
    
  3. 使用 ServiceBusMessageListenerContainer Bean 建立 ServiceBusInboundChannelAdapter,以接收來自 服務匯流排 的訊息。 此範例會採用服務總線佇列作為範例。

    @Configuration
    class Demo {
        private static final String QUEUE_NAME = "queue1";
    
        @Bean
        public ServiceBusMessageListenerContainer messageListenerContainer(ServiceBusProcessorFactory processorFactory) {
            ServiceBusContainerProperties containerProperties = new ServiceBusContainerProperties();
            containerProperties.setEntityName(QUEUE_NAME);
            containerProperties.setAutoComplete(false);
            return new ServiceBusMessageListenerContainer(processorFactory, containerProperties);
        }
    
        @Bean
        public ServiceBusInboundChannelAdapter queueMessageChannelAdapter(
            @Qualifier(INPUT_CHANNEL) MessageChannel inputChannel,
            ServiceBusMessageListenerContainer listenerContainer) {
            ServiceBusInboundChannelAdapter adapter = new ServiceBusInboundChannelAdapter(listenerContainer);
            adapter.setOutputChannel(inputChannel);
            return adapter;
        }
    }
    
  4. 透過我們先前建立的訊息通道,建立具有 ServiceBusInboundChannelAdapter 的訊息接收者系結。

    class Demo {
        @ServiceActivator(inputChannel = INPUT_CHANNEL)
        public void messageReceiver(byte[] payload, @Header(AzureHeaders.CHECKPOINTER) Checkpointer checkpointer) {
            String message = new String(payload);
            LOGGER.info("New message received: '{}'", message);
            checkpointer.success()
                    .doOnSuccess(s -> LOGGER.info("Message '{}' successfully checkpointed", message))
                    .doOnError(e -> LOGGER.error("Error found", e))
                    .block();
        }
    }
    

設定 ServiceBusMessageConverter 以自定義 objectMapper

ServiceBusMessageConverter 被設計為可配置的 Bean,讓使用者能夠自訂 ObjectMapper。

服務總線訊息標頭

對於一些可以對應至多個 Spring 標頭常數的服務總線標頭,會列出不同 Spring 標頭的優先順序。

服務總線標頭與 Spring 標頭之間的對應:

服務總線訊息標頭和屬性 Spring 訊息標頭常數 類型 可設定 描述
內容類型 MessageHeaders#CONTENT_TYPE 字串 是的 訊息的RFC2045內容類型描述元。
關聯識別碼 ServiceBusMessageHeaders#CORRELATION_ID 字串 是的 訊息的相互關聯標識碼
訊息標識碼 ServiceBusMessageHeaders#MESSAGE_ID 字串 是的 訊息的訊息識別碼,此標頭的優先順序高於 MessageHeaders#ID。
訊息標識碼 MessageHeaders#ID 通用唯一識別碼 (UUID) 是的 訊息的訊息識別碼,此標頭的優先順序低於 ServiceBusMessageHeaders#MESSAGE_ID。
分割區索引鍵 ServiceBusMessageHeaders#PARTITION_KEY 字串 是的 用於將訊息傳送至分割區實體的分割區索引鍵。
回覆 MessageHeaders#REPLY_CHANNEL 字串 是的 要傳送回復之實體的位址。
回復會話標識碼 ServiceBusMessageHeaders#REPLY_TO_SESSION_ID 字串 是的 訊息的 ReplyToGroupId 屬性值。
排程排入佇列時間 UTC ServiceBusMessageHeaders#SCHEDULED_ENQUEUE_TIME 偏移日期時間 是的 訊息應在 服務匯流排 中排入佇列的日期與時間,此標頭的優先順序高於 AzureHeaders#SCHEDULED_ENQUEUE_MESSAGE。
已排程的加入佇列時間 (UTC) AzureHeaders#SCHEDULED_ENQUEUE_MESSAGE 整數 是的 訊息應排入 服務匯流排 佇列的日期與時間;此標頭的優先順序低於 ServiceBusMessageHeaders#SCHEDULED_ENQUEUE_TIME。
會話標識碼 ServiceBusMessageHeaders#SESSION_ID 字串 是的 具工作階段感知能力之實體的工作階段識別碼。
存活時間 ServiceBusMessageHeaders#TIME_TO_LIVE 期間 是的 此訊息到期之前的持續時間。
到 ServiceBusMessageHeaders#TO 字串 是的 訊息的「to」位址,保留供未來在路由情境中使用,且目前會被訊息代理程式本身忽略。
主題 ServiceBusMessageHeaders#SUBJECT 字串 是的 訊息的主旨。
死信錯誤描述 ServiceBusMessageHeaders#DEAD_LETTER_ERROR_DESCRIPTION 字串 不 已成為死信之訊息的描述。
死信原因 ServiceBusMessageHeaders#DEAD_LETTER_REASON 字串 不 訊息成為死信的原因。
死信來源 ServiceBusMessageHeaders#DEAD_LETTER_SOURCE 字串 不 訊息被移至死信佇列的實體。
傳送次數 ServiceBusMessageHeaders#DELIVERY_COUNT 長 不 此訊息傳遞至客戶端的次數。
已加入佇列的序號 ServiceBusMessageHeaders#ENQUEUED_SEQUENCE_NUMBER 長 不 由 服務匯流排 指派給訊息的已排入佇列序號。
加入佇列的時間 ServiceBusMessageHeaders#ENQUEUED_TIME 偏移日期時間 不 此訊息在 服務匯流排 中排入佇列的日期和時間。
到期時間: ServiceBusMessageHeaders#EXPIRES_AT 偏移日期時間 不 此訊息到期的日期時間。
鎖定令牌 ServiceBusMessageHeaders#LOCK_TOKEN 字串 不 目前訊息的鎖定令牌。
鎖定直到 ServiceBusMessageHeaders#LOCKED_UNTIL 偏移日期時間 不 此訊息的鎖定失效日期時間。
序號 ServiceBusMessageHeaders#SEQUENCE_NUMBER 長 不 服務總線指派給訊息的唯一號碼。
州 ServiceBusMessageHeaders#STATE ServiceBusMessageState 不 訊息的狀態,可以是 [作用中]、[延遲] 或 [已排程]。

分割區索引鍵支援

此入門支援 服務總線分割,方法是允許在訊息標頭中設定分割區索引鍵和會話標識符。 本節將介紹如何為訊息設定分割區索引鍵。

建議:使用 ServiceBusMessageHeaders.PARTITION_KEY 作為標頭的鍵。

public class SampleController {
    @PostMapping("/messages")
    public ResponseEntity<String> sendMessage(@RequestParam String message) {
        LOGGER.info("Going to add message {} to Sinks.Many.", message);
        many.emitNext(MessageBuilder.withPayload(message)
                                    .setHeader(ServiceBusMessageHeaders.PARTITION_KEY, "Customize partition key")
                                    .build(), Sinks.EmitFailureHandler.FAIL_FAST);
        return ResponseEntity.ok("Sent!");
    }
}

不建議,但目前仍支援將 AzureHeaders.PARTITION_KEY 作為標頭的鍵。

public class SampleController {
    @PostMapping("/messages")
    public ResponseEntity<String> sendMessage(@RequestParam String message) {
        LOGGER.info("Going to add message {} to Sinks.Many.", message);
        many.emitNext(MessageBuilder.withPayload(message)
                                    .setHeader(AzureHeaders.PARTITION_KEY, "Customize partition key")
                                    .build(), Sinks.EmitFailureHandler.FAIL_FAST);
        return ResponseEntity.ok("Sent!");
    }
}

注意

當訊息標頭中同時設定 ServiceBusMessageHeaders.PARTITION_KEY 和 AzureHeaders.PARTITION_KEY 時,建議使用 ServiceBusMessageHeaders.PARTITION_KEY。

會話支援

此範例示範如何在應用程式中手動設定訊息的會話標識碼。

public class SampleController {
    @PostMapping("/messages")
    public ResponseEntity<String> sendMessage(@RequestParam String message) {
        LOGGER.info("Going to add message {} to Sinks.Many.", message);
        many.emitNext(MessageBuilder.withPayload(message)
                                    .setHeader(ServiceBusMessageHeaders.SESSION_ID, "Customize session ID")
                                    .build(), Sinks.EmitFailureHandler.FAIL_FAST);
        return ResponseEntity.ok("Sent!");
    }
}

注意

當訊息標頭中設定 ServiceBusMessageHeaders.SESSION_ID,而且也會設定不同的 ServiceBusMessageHeaders.PARTITION_KEY 標頭時,會話標識碼的值最終將用來覆寫分割區索引鍵的值。

自訂服務總線客戶端屬性

開發人員可以使用 AzureServiceClientBuilderCustomizer 來自定義服務總線用戶端屬性。 下列範例會自訂 sessionIdleTimeout中的 ServiceBusClientBuilder 屬性:

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

樣品

欲了解更多資訊,請參閱 azure-spring-boot-samples GitHub 上的倉庫。

Spring 與 Azure 記憶體佇列整合

重要概念

Azure 佇列記憶體是用來儲存大量訊息的服務。 您可以使用 HTTP 或 HTTPS 透過已驗證的呼叫,從世界各地存取訊息。 佇列訊息的大小最多可達 64 KB。 佇列可能包含數百萬則訊息,最多可達記憶體帳戶的總容量限制。 佇列通常用於建立待處理工作的積壓,供非同步處理。

相依性設定

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

配置

此入門提供下列組態選項:

聯機組態屬性

本節包含用來連線到 Azure 記憶體佇列的組態選項。

注意

如果您選擇使用安全性主體向 Microsoft Entra ID 進行驗證和授權,以存取 Azure 資源,請參閱 使用 Microsoft Entra ID 授權存取,以確保安全性主體已獲得存取 Azure 資源的足夠許可權。

spring-cloud-azure-starter-integration-storage-queue 的連線可設定屬性:

財產 類型 描述
spring.cloud.azure.storage.queue.enabled 布爾 是否啟用 Azure 記憶體佇列。
spring.cloud.azure.storage.queue.connection-string 字串 記憶體佇列命名空間連接字串值。
spring.cloud.azure.storage.queue.accountName 字串 記憶體佇列帳戶名稱。
spring.cloud.azure.storage.queue.accountKey 字串 儲存體佇列帳戶金鑰。
spring.cloud.azure.storage.queue.endpoint 字串 記憶體佇列服務端點。
spring.cloud.azure.storage.queue.sasToken 字串 Sas 令牌認證
spring.cloud.azure.storage.queue.serviceVersion QueueServiceVersion 發出 API 要求時所使用的 QueueServiceVersion。
spring.cloud.azure.storage.queue.messageEncoding 字串 佇列訊息編碼。

基本用法

將訊息傳送至 Azure 記憶體佇列

  1. 填入認證組態選項。

    • 針對認證作為連接字串,請在 application.yml 檔案中設定下列屬性:

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

      注意

      Microsoft 建議您使用最安全的可用驗證流程。 此程式中所述的驗證流程,例如資料庫、快取、傳訊或 AI 服務,在應用程式中需要高度的信任,而且不會在其他流程中帶來風險。 只有在更安全的選項(例如使用受控識別進行無密碼或無金鑰連線)不可行時,才使用此流程。 針對本機計算機作業,偏好使用無密碼或無密鑰連線的使用者身分識別。

    • 針對認證作為受控識別,請在 application.yml 檔案中設定下列屬性:

      spring:
        cloud:
          azure:
            credential:
              managed-identity-enabled: true
              client-id: ${AZURE_CLIENT_ID}
            profile:
              tenant-id: <tenant>
            storage:
              queue:
                account-name: ${AZURE_STORAGE_QUEUE_ACCOUNT_NAME}
      

注意

tenant-id 允許的值包括:common、organizations、consumers或租用戶標識碼。 如需這些值的詳細資訊,請參閱 錯誤 AADSTS50020 - 來自身分識別提供者的使用者帳戶不存在於租用戶中中的 使用了錯誤的端點(個人和組織帳戶)一節。 如需有關將單一租用戶應用程式轉換為多租用戶的資訊,請參閱 在 Microsoft Entra ID 上將單一租用戶應用程式轉換為多租用戶。

  • 若使用服務主體認證,請在 application.yml 檔案中設定下列屬性:

    spring:
      cloud:
        azure:
          credential:
            client-id: ${AZURE_CLIENT_ID}
            client-secret: ${AZURE_CLIENT_SECRET}
          profile:
            tenant-id: <tenant>
          storage:
            queue:
              account-name: ${AZURE_STORAGE_QUEUE_ACCOUNT_NAME}
    

注意

tenant-id 允許的值包括:common、organizations、consumers或租用戶標識碼。 如需這些值的詳細資訊,請參閱 錯誤 AADSTS50020 - 來自身分識別提供者的使用者帳戶不存在於租用戶中中的 使用了錯誤的端點(個人和組織帳戶)一節。 如需有關將單一租用戶應用程式轉換為多租用戶的資訊,請參閱 在 Microsoft Entra ID 上將單一租用戶應用程式轉換為多租用戶。

  1. 使用 StorageQueueTemplate Bean 建立 DefaultMessageHandler,以將訊息傳送至儲存體佇列。

    class Demo {
        private static final String STORAGE_QUEUE_NAME = "example";
        private static final String OUTPUT_CHANNEL = "output";
    
        @Bean
        @ServiceActivator(inputChannel = OUTPUT_CHANNEL)
        public MessageHandler messageSender(StorageQueueTemplate storageQueueTemplate) {
            DefaultMessageHandler handler = new DefaultMessageHandler(STORAGE_QUEUE_NAME, storageQueueTemplate);
            handler.setSendCallback(new ListenableFutureCallback<Void>() {
                @Override
                public void onSuccess(Void result) {
                    LOGGER.info("Message was sent successfully.");
                }
    
                @Override
                public void onFailure(Throwable ex) {
                    LOGGER.error("There was an error sending the message.", ex);
                }
            });
            return handler;
        }
    }
    
  2. 透過訊息通道,使用上述訊息處理程式建立訊息網關係結。

    class Demo {
        @Autowired
        StorageQueueOutboundGateway storageQueueOutboundGateway;
    
        @MessagingGateway(defaultRequestChannel = OUTPUT_CHANNEL)
        public interface StorageQueueOutboundGateway {
            void send(String text);
        }
    }
    
  3. 使用閘道傳送訊息。

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

從 Azure 記憶體佇列接收訊息

  1. 填入認證組態選項。

  2. 建立訊息通道的 Bean 做為輸入通道。

    class Demo {
        private static final String INPUT_CHANNEL = "input";
    
        @Bean
        public MessageChannel input() {
            return new DirectChannel();
        }
    }
    
  3. 使用 StorageQueueTemplate Bean 建立 StorageQueueMessageSource,以接收來自儲存體佇列的訊息。

    class Demo {
        private static final String STORAGE_QUEUE_NAME = "example";
    
        @Bean
        @InboundChannelAdapter(channel = INPUT_CHANNEL, poller = @Poller(fixedDelay = "1000"))
        public StorageQueueMessageSource storageQueueMessageSource(StorageQueueTemplate storageQueueTemplate) {
            return new StorageQueueMessageSource(STORAGE_QUEUE_NAME, storageQueueTemplate);
        }
    }
    
  4. 透過我們先前建立的訊息通道,使用上一個步驟中建立的 StorageQueueMessageSource 建立訊息接收器繫結。

    class Demo {
        @ServiceActivator(inputChannel = INPUT_CHANNEL)
        public void messageReceiver(byte[] payload, @Header(AzureHeaders.CHECKPOINTER) Checkpointer checkpointer) {
            String message = new String(payload);
            LOGGER.info("New message received: '{}'", message);
            checkpointer.success()
                .doOnError(Throwable::printStackTrace)
                .doOnSuccess(t -> LOGGER.info("Message '{}' successfully checkpointed", message))
                .block();
        }
    }
    

樣品

欲了解更多資訊,請參閱 azure-spring-boot-samples GitHub 上的倉庫。