適用於 Azure 的 Spring Integration Extension 為 Azure SDK for Java 提供的各種服務提供 Spring Integration 配接器。 我們提供這些 Azure 服務的 Spring Integration 支援:事件中樞、服務總線、記憶體佇列。 以下是支援的配接器清單:
-
spring-cloud-azure-starter-integration-eventhubs- 如需詳細資訊,請參閱 Spring Integration with Azure 事件中樞 -
spring-cloud-azure-starter-integration-servicebus- 如需詳細資訊,請參閱 Spring Integration with Azure 服務匯流排 -
spring-cloud-azure-starter-integration-storage-queue- 如需詳細資訊,請參閱 Spring Integration with Azure 儲存體 Queue
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 事件中樞
填入認證組態選項。
針對認證作為連接字串,請在 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 上將單一租用戶應用程式轉換為多租用戶。
使用
EventHubsTemplateBean 建立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; } }透過訊息通道建立具有上述訊息處理程式的訊息網關係結。
class Demo { @Autowired EventHubOutboundGateway messagingGateway; @MessagingGateway(defaultRequestChannel = OUTPUT_CHANNEL) public interface EventHubOutboundGateway { void send(String text); } }使用閘道傳送訊息。
class Demo { public void demo() { this.messagingGateway.send(message); } }
從 Azure 事件中樞接收訊息
填入認證組態選項。
建立訊息通道的 Bean 做為輸入通道。
@Configuration class Demo { @Bean public MessageChannel input() { return new DirectChannel(); } }使用
EventHubsMessageListenerContainerbean 建立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); } }透過之前建立的訊息通道,使用 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 服務總線
填入認證組態選項。
針對認證作為連接字串,請在 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 上將單一租用戶應用程式轉換為多租用戶。
使用
ServiceBusTemplatebean 建立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; } }透過訊息通道建立具有上述訊息處理程式的訊息網關係結。
class Demo { @Autowired QueueOutboundGateway messagingGateway; @MessagingGateway(defaultRequestChannel = OUTPUT_CHANNEL) public interface QueueOutboundGateway { void send(String text); } }使用閘道傳送訊息。
class Demo { public void demo() { this.messagingGateway.send(message); } }
從 Azure 服務總線接收訊息
填入認證組態選項。
建立訊息通道的 Bean 做為輸入通道。
@Configuration class Demo { private static final String INPUT_CHANNEL = "input"; @Bean public MessageChannel input() { return new DirectChannel(); } }使用
ServiceBusMessageListenerContainerBean 建立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; } }透過我們先前建立的訊息通道,建立具有
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 記憶體佇列
填入認證組態選項。
針對認證作為連接字串,請在 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 上將單一租用戶應用程式轉換為多租用戶。
使用
StorageQueueTemplateBean 建立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; } }透過訊息通道,使用上述訊息處理程式建立訊息網關係結。
class Demo { @Autowired StorageQueueOutboundGateway storageQueueOutboundGateway; @MessagingGateway(defaultRequestChannel = OUTPUT_CHANNEL) public interface StorageQueueOutboundGateway { void send(String text); } }使用閘道傳送訊息。
class Demo { public void demo() { this.storageQueueOutboundGateway.send(message); } }
從 Azure 記憶體佇列接收訊息
填入認證組態選項。
建立訊息通道的 Bean 做為輸入通道。
class Demo { private static final String INPUT_CHANNEL = "input"; @Bean public MessageChannel input() { return new DirectChannel(); } }使用
StorageQueueTemplateBean 建立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); } }透過我們先前建立的訊息通道,使用上一個步驟中建立的 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 上的倉庫。