Spring Cloud Stream 是建置與共用傳訊系統連線的高度可調整事件驅動微服務的架構。
此架構提供一個彈性的程序設計模型,建置在已建立且熟悉的 Spring 慣用語和最佳做法上。 這些最佳做法包括支援持久性發布/訂閱語意、消費者群組,以及具狀態的分割區。
目前的 binder 實作包括:
-
spring-cloud-azure-stream-binder-eventhubs- 如需詳細資訊,請參閱 適用於 Azure 事件中樞的 Spring Cloud Stream Binder -
spring-cloud-azure-stream-binder-servicebus- 如需詳細資訊,請參閱 適用於 Azure 服務總線的 Spring Cloud Stream Binder
適用於 Azure 事件中樞的 Spring Cloud Stream Binder
重要概念
適用於 Azure 事件中樞的 Spring Cloud Stream Binder 提供 Spring Cloud Stream 架構的系結實作。 此實作會在其基礎上使用 Spring Integration 事件中樞通道配接器。 從設計的觀點來看,事件中樞與 Kafka 類似。 此外,事件中樞也可以透過 Kafka API 存取。 如果您的專案與 Kafka API 有緊密的相依性,您可以嘗試使用 Kafka API 範例 事件中樞
取用者群組
事件中樞提供與 Apache Kafka 類似的取用者群組支援,但邏輯稍有不同。 雖然 Kafka 會將所有已提交的位移儲存在訊息代理中,但您必須手動儲存正在處理的 Event Hubs 訊息位移。 事件中樞 SDK 提供函式,以將這類位移儲存在 Azure 記憶體內。
資料分割支援
事件中樞提供與 Kafka 類似的實體分割區概念。 但不同於 Kafka 在取用者與分割區之間的自動重新平衡機制,Event Hubs 提供某種先占式模式。 儲存體帳戶可作為租約,用來判斷哪個取用者擁有哪個分割區。 當新的取用者啟動時,它會嘗試從負載最重的取用者竊取一些分割區,以達到工作負載平衡。
若要指定負載平衡策略,則會提供 spring.cloud.stream.eventhubs.bindings.<binding-name>.consumer.load-balancing.* 的屬性。 如需詳細資訊,請參閱 取用者屬性 一節。
批次取用端支援
Spring Cloud Azure Stream 事件中樞系結器支援 Spring Cloud Stream Batch 取用者功能。
若要使用批次取用者模式,請將 spring.cloud.stream.bindings.<binding-name>.consumer.batch-mode 屬性設定為 true。 啟用時,會收到具有批次事件清單承載的訊息,並傳遞至 Consumer 函式。 每個訊息標頭也會轉換成清單,其中內容是從每個事件剖析的相關聯標頭值。 分割區識別碼、檢查點器和最後排入佇列屬性的共同標頭會以單一值顯示,因為整個事件批次共用相同的值。 如需詳細資訊,請參閱 Spring Cloud Azure 支援 Spring Integration 的 事件中樞訊息標頭 一節。
注意
只有在使用 MANUAL 檢查點模式時,才會有檢查點標頭。
批次取用程式的檢查點功能支援兩種模式:BATCH 和 MANUAL。
BATCH 模式是一種自動檢查點模式,一旦繫結器收到事件,就會為整個事件批次一併建立檢查點。
MANUAL 模式用於讓使用者為事件建立檢查點。 使用時,Checkpointer 會傳遞至訊息標頭,使用者可以使用它進行檢查點檢查。
您可以藉由設定前置詞為 max-size的 max-wait-time 和 spring.cloud.stream.eventhubs.bindings.<binding-name>.consumer.batch. 屬性來指定批次大小。
max-size 屬性是必要的,而且 max-wait-time 屬性是選擇性的。 如需詳細資訊,請參閱 取用者屬性 一節。
相依性設定
<dependency>
<groupId>com.azure.spring</groupId>
<artifactId>spring-cloud-azure-stream-binder-eventhubs</artifactId>
</dependency>
或者,您也可以使用 Spring Cloud Azure Stream 事件中樞入門版,如下列 Maven 範例所示:
<dependency>
<groupId>com.azure.spring</groupId>
<artifactId>spring-cloud-azure-starter-stream-eventhubs</artifactId>
</dependency>
配置
系結器提供下列三個組態選項部分:
聯機組態屬性
本節包含用來連線到 Azure 事件中樞的組態選項。
注意
如果您選擇使用安全性主體向 Microsoft Entra ID 進行驗證和授權,以存取 Azure 資源,請參閱 使用 Microsoft Entra ID 授權存取,以確保安全性主體已獲得存取 Azure 資源的足夠許可權。
spring-cloud-azure-stream-binder-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 |
字串 | 自訂端點位址。 |
提示
一般 Azure 服務 SDK 組態選項也可以針對 Spring Cloud Azure Stream 事件中樞系結器進行設定。 支援的組態選項已於 Spring Cloud Azure 組態中介紹,且可使用統一前綴 spring.cloud.azure. 或 spring.cloud.azure.eventhubs. 前綴進行設定。
系結器預設也支援 spring Could Azure Resource Manager 。 若要了解如何擷取未獲授與 相關角色的安全性主體之連接字串,請參閱 Spring Could Azure Resource Manager 的 基本用法 一節。
檢查點組態屬性
本節包含記憶體 Blob 服務的組態選項,用於保存分割區擁有權和檢查點資訊。
注意
從版本 4.0.0 開始,當 的 spring.cloud.azure.eventhubs.processor.checkpoint-store.create-container-if-not-exists 屬性未手動啟用時,不會自動建立帶有 的 spring.cloud.stream.bindings.binding-name.destination儲存容器名稱。
spring-cloud-azure-stream-binder-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 前綴進行設定。
Azure 事件中樞系結組態屬性
下列選項分為四個區段:取用者屬性、進階取用者組態、產生者屬性和進階生產者設定。
消費者屬性
這些屬性會透過 EventHubsConsumerProperties公開。
注意
為避免重複,自 4.17.0 與 5.11.0 版起,Spring Cloud Azure Stream Binder Event Hubs 支援為所有通道設定值,格式為 spring.cloud.stream.eventhubs.default.consumer.<property>=<value>。
spring-cloud-azure-stream-binder-eventhubs 的使用者可設定屬性:
| 財產 | 類型 | 描述 |
|---|---|---|
spring.cloud.stream.eventhubs.bindings.binding-name.consumer.checkpoint.mode |
檢查點模式 | 當取用者決定如何為訊息建立檢查點時所使用的檢查點模式 |
spring.cloud.stream.eventhubs.bindings.binding-name.consumer.checkpoint.count |
整數 | 決定每個分割區進行一次檢查點所需的訊息數量。 只有在使用 PARTITION_COUNT 檢查點模式時才會生效。 |
spring.cloud.stream.eventhubs.bindings.binding-name.consumer.checkpoint.interval |
期間 | 決定執行一次檢查點的時間間隔。 只有在使用 TIME 檢查點模式時才會生效。 |
spring.cloud.stream.eventhubs.bindings.<binding-name>.consumer.batch.max-size |
整數 | 批次中的事件數目上限。 為批次取用模式所必需。 |
spring.cloud.stream.eventhubs.bindings.binding-name.consumer.batch.max-wait-time |
期間 | 批次耗用的最大持續時間。 只有在啟用批次取用者模式且為選擇性時,才會生效。 |
spring.cloud.stream.eventhubs.bindings.binding-name.consumer.load-balancing.update-interval |
期間 | 更新的間隔時間。 |
spring.cloud.stream.eventhubs.bindings.binding-name.consumer.load-balancing.strategy |
LoadBalancingStrategy (負載均衡策略) | 負載平衡策略。 |
spring.cloud.stream.eventhubs.bindings.binding-name.consumer.load-balancing.partition-ownership-expiration-interval |
期間 | 分割區擁有權在多久後到期。 |
spring.cloud.stream.eventhubs.bindings.binding-name.consumer.track-last-enqueued-event-properties |
布林值 | 事件處理器是否應在其所屬分割區上要求取得最後一個已排入佇列事件的資訊,並在接收事件時追蹤該資訊。 |
spring.cloud.stream.eventhubs.bindings.binding-name.consumer.prefetch-count |
整數 | 供消費者用來控制事件中樞消費者會主動接收並在本機排入佇列的事件數量。 |
spring.cloud.stream.eventhubs.bindings.binding-name.consumer.initial-partition-event-position |
將索引鍵對應為分割區標識碼,以及 StartPositionProperties 的值 |
如果分割區在檢查點存放區中不存在檢查點,則使用包含各分割區事件位置的對應表。 此對應表是以分割區 ID 作為鍵。 |
注意
initial-partition-event-position 組態接受 map 來指定每個事件中樞的初始位置。 因此,其索引鍵是分割區標識符,而值是 StartPositionProperties,其中包含位移、序號、加入佇列日期時間和是否包含的屬性。 例如,您可以將它設定為
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
進階使用者設定
上述 連線、檢查點和 通用 Azure SDK 用戶端組態支援針對每個繫結器取用者進行自訂,您可以使用前置詞 spring.cloud.stream.eventhubs.bindings.<binding-name>.consumer. 來進行設定。
生產者屬性
這些屬性會透過 EventHubsProducerProperties公開。
注意
為避免重複,自 4.17.0 與 5.11.0 版起,Spring Cloud Azure Stream Binder Event Hubs 支援為所有通道設定值,格式為 spring.cloud.stream.eventhubs.default.producer.<property>=<value>。
spring-cloud-azure-stream-binder-eventhubs 的生產者可設定屬性:
| 財產 | 類型 | 描述 |
|---|---|---|
spring.cloud.stream.eventhubs.bindings.binding-name.producer.sync |
布爾 | 生產者同步的開關旗標。 如果為 true,產生者會在傳送作業之後等候回應。 |
spring.cloud.stream.eventhubs.bindings.binding-name.producer.send-timeout |
長 | 在傳送作業之後等候回應的時間量。 只有在啟用同步產生者時才會生效。 |
進階生產者設定
上述 連線和 通用 Azure SDK 用戶端組態支援針對每個繫結器生產者進行自訂,您可以使用前綴 spring.cloud.stream.eventhubs.bindings.<binding-name>.producer. 進行設定。
基本用法
向/從 Event Hubs 傳送及接收訊息
使用認證資訊填入組態選項。
針對認證作為連接字串,請在 application.yml 檔案中設定下列屬性:
spring: cloud: azure: eventhubs: connection-string: ${EVENTHUB_NAMESPACE_CONNECTION_STRING} processor: checkpoint-store: container-name: ${CHECKPOINT_CONTAINER} account-name: ${CHECKPOINT_STORAGE_ACCOUNT} account-key: ${CHECKPOINT_ACCESS_KEY} function: definition: consume;supply stream: bindings: consume-in-0: destination: ${EVENTHUB_NAME} group: ${CONSUMER_GROUP} supply-out-0: destination: ${THE_SAME_EVENTHUB_NAME_AS_ABOVE} eventhubs: bindings: consume-in-0: consumer: checkpoint: mode: MANUAL注意
Microsoft 建議您使用最安全的可用驗證流程。 此程式中所述的驗證流程,例如資料庫、快取、傳訊或 AI 服務,在應用程式中需要高度的信任,而且不會在其他流程中帶來風險。 只有在更安全的選項(例如使用受控識別進行無密碼或無金鑰連線)不可行時,才使用此流程。 針對本機計算機作業,偏好使用無密碼或無密鑰連線的使用者身分識別。
若使用服務主體認證,請在 application.yml 檔案中設定下列屬性:
spring: cloud: azure: credential: client-id: ${AZURE_CLIENT_ID} client-secret: ${AZURE_CLIENT_SECRET} profile: tenant-id: <tenant> eventhubs: namespace: ${EVENTHUB_NAMESPACE} processor: checkpoint-store: container-name: ${CONTAINER_NAME} account-name: ${ACCOUNT_NAME} function: definition: consume;supply stream: bindings: consume-in-0: destination: ${EVENTHUB_NAME} group: ${CONSUMER_GROUP} supply-out-0: destination: ${THE_SAME_EVENTHUB_NAME_AS_ABOVE} eventhubs: bindings: consume-in-0: consumer: checkpoint: mode: MANUAL
注意
tenant-id 允許的值包括:common、organizations、consumers或租用戶標識碼。 如需這些值的詳細資訊,請參閱 錯誤 AADSTS50020 - 來自身分識別提供者的使用者帳戶不存在於租用戶中中的 使用了錯誤的端點(個人和組織帳戶)一節。 如需有關將單一租用戶應用程式轉換為多租用戶的資訊,請參閱 在 Microsoft Entra ID 上將單一租用戶應用程式轉換為多租用戶。
針對認證作為受控識別,請在 application.yml 檔案中設定下列屬性:
spring: cloud: azure: credential: managed-identity-enabled: true client-id: ${AZURE_MANAGED_IDENTITY_CLIENT_ID} # Only needed when using a user-assigned managed identity eventhubs: namespace: ${EVENTHUB_NAMESPACE} processor: checkpoint-store: container-name: ${CONTAINER_NAME} account-name: ${ACCOUNT_NAME} function: definition: consume;supply stream: bindings: consume-in-0: destination: ${EVENTHUB_NAME} group: ${CONSUMER_GROUP} supply-out-0: destination: ${THE_SAME_EVENTHUB_NAME_AS_ABOVE} eventhubs: bindings: consume-in-0: consumer: checkpoint: mode: MANUAL
定義供應商和取用者。
@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(); }; }
資料分割支援
系統會建立一個包含使用者提供之分割區資訊的 PartitionSupplier,用來設定即將傳送之訊息的分割區資訊。 下列流程圖顯示取得分割區識別碼和索引鍵之不同優先順序的程式:
批次取用端支援
提供批次組態選項,如下列範例所示:
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定義供應商和取用者。
若檢查點模式為
BATCH,您可以使用下列程式碼來傳送訊息,並批次取用訊息。@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(); }; }若檢查點模式為
MANUAL,您可以使用下列程式碼來傳送訊息,並以批次方式取用及建立檢查點。@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(); }; }
注意
在批次取用模式中,Spring Cloud Stream 系結器的預設內容類型是 application/json,因此請確定訊息承載與內容類型一致。 例如,當使用預設內容類型 application/json 來接收含有 String 承載的訊息時,承載應為 JSON String,而原始的 String 文字需以雙引號括住。 至於 text/plain 內容類型,它可直接為 String 物件。 如需詳細資訊,請參閱 Spring Cloud Stream 內容類型交涉。
處理錯誤訊息
處理輸出系結錯誤訊息
根據預設,Spring Integration 會建立名為
errorChannel的全域錯誤通道。 設定下列訊息端點來處理輸出系結錯誤訊息。@ServiceActivator(inputChannel = IntegrationContextUtils.ERROR_CHANNEL_BEAN_NAME) public void handleError(ErrorMessage message) { LOGGER.error("Handling outbound binding error: " + message); }處理輸入系結錯誤訊息
Spring Cloud Stream 事件中樞系結器支援一個解決方案來處理輸入訊息系結的錯誤:錯誤處理程式。
錯誤處理程式:
Spring Cloud Stream 會公開一種機制,讓您藉由新增可接受
Consumer實例的ErrorMessage來提供自定義錯誤處理程式。 如需詳細資訊,請參閱 Spring Cloud Stream 檔中 處理錯誤訊息。系結預設錯誤處理程式
設定單一
ConsumerBean 以處理所有傳入繫結錯誤訊息。 下列預設函式會訂閱每個輸入系結錯誤通道:@Bean public Consumer<ErrorMessage> myDefaultHandler() { return message -> { // consume the error message }; }您也需要將
spring.cloud.stream.default.error-handler-definition屬性設定為函式名稱。系結特定錯誤處理程式
設定
ConsumerBean,以處理特定傳入繫結的錯誤訊息。 下列函式會訂閱特定的輸入系結錯誤通道,且優先順序高於系結預設錯誤處理程式:@Bean public Consumer<ErrorMessage> myErrorHandler() { return message -> { // consume the error message }; }您也需要將
spring.cloud.stream.bindings.<input-binding-name>.error-handler-definition屬性設定為函式名稱。
Event Hubs 訊息標頭
如需了解支援的基本訊息標頭,請參閱 Spring Cloud Azure 對 Spring Integration 的支援中的 事件中樞訊息標頭一節。
多個系結器支援
透過使用多個系結器,也支援連線到多個 Event Hubs 命名空間。 此範例會採用連接字串作為範例。 也支援服務主體和受控識別的認證。 您可以在每個系結器的環境設定中設定相關的屬性。
若要搭配事件中樞使用多個系結器,請在 application.yml 檔案中設定下列屬性:
spring: cloud: function: definition: consume1;supply1;consume2;supply2 stream: bindings: consume1-in-0: destination: ${EVENTHUB_NAME_01} group: ${CONSUMER_GROUP_01} supply1-out-0: destination: ${THE_SAME_EVENTHUB_NAME_01_AS_ABOVE} consume2-in-0: binder: eventhub-2 destination: ${EVENTHUB_NAME_02} group: ${CONSUMER_GROUP_02} supply2-out-0: binder: eventhub-2 destination: ${THE_SAME_EVENTHUB_NAME_02_AS_ABOVE} binders: eventhub-1: type: eventhubs default-candidate: true environment: spring: cloud: azure: eventhubs: connection-string: ${EVENTHUB_NAMESPACE_01_CONNECTION_STRING} processor: checkpoint-store: container-name: ${CHECKPOINT_CONTAINER_01} account-name: ${CHECKPOINT_STORAGE_ACCOUNT} account-key: ${CHECKPOINT_ACCESS_KEY} eventhub-2: type: eventhubs default-candidate: false environment: spring: cloud: azure: eventhubs: connection-string: ${EVENTHUB_NAMESPACE_02_CONNECTION_STRING} processor: checkpoint-store: container-name: ${CHECKPOINT_CONTAINER_02} account-name: ${CHECKPOINT_STORAGE_ACCOUNT} account-key: ${CHECKPOINT_ACCESS_KEY} eventhubs: bindings: consume1-in-0: consumer: checkpoint: mode: MANUAL consume2-in-0: consumer: checkpoint: mode: MANUAL poller: initial-delay: 0 fixed-delay: 1000注意
前述的應用程式檔案顯示如何設定單一預設輪詢器,並將其套用至所有繫結。 如果您想為特定繫結設定輪詢器,可以使用如下設定,例如
spring.cloud.stream.bindings.<binding-name>.producer.poller.fixed-delay=3000。注意
Microsoft 建議您使用最安全的可用驗證流程。 此程式中所述的驗證流程,例如資料庫、快取、傳訊或 AI 服務,在應用程式中需要高度的信任,而且不會在其他流程中帶來風險。 只有在更安全的選項(例如使用受控識別進行無密碼或無金鑰連線)不可行時,才使用此流程。 針對本機計算機作業,偏好使用無密碼或無密鑰連線的使用者身分識別。
我們需要定義兩個供應商和兩個取用者:
@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(); }; }
資源布建
Event Hubs 繫結器支援佈建事件中樞和取用者群組,使用者可使用下列屬性來啟用此功能。
spring:
cloud:
azure:
credential:
tenant-id: <tenant>
profile:
subscription-id: ${AZURE_SUBSCRIPTION_ID}
eventhubs:
resource:
resource-group: ${AZURE_EVENTHUBS_RESOURCE_GROUP}
注意
tenant-id 允許的值包括:common、organizations、consumers或租用戶標識碼。 如需這些值的詳細資訊,請參閱 錯誤 AADSTS50020 - 來自身分識別提供者的使用者帳戶不存在於租用戶中中的 使用了錯誤的端點(個人和組織帳戶)一節。 如需有關將單一租用戶應用程式轉換為多租用戶的資訊,請參閱 在 Microsoft Entra ID 上將單一租用戶應用程式轉換為多租用戶。
樣品
欲了解更多資訊,請參閱 azure-spring-boot-samples GitHub 上的倉庫。
適用於 Azure 服務總線的 Spring Cloud Stream Binder
重要概念
適用於 Azure 服務總線的 Spring Cloud Stream Binder 提供 Spring Cloud Stream 架構的系結實作。 此實作會在其基礎上使用 Spring Integration Service 總線通道配接器。
排定訊息
此繫結器支援將訊息傳送到主題,以供延後處理。 用戶可以傳送具有標頭 x-delay 以毫秒表示訊息延遲時間的排程訊息。 訊息會在 x-delay 毫秒之後傳遞至個別主題。
取用者群組
服務總線主題提供與 Apache Kafka 類似的取用者群組支援,但邏輯稍有不同。
此系結器依賴主題 Subscription 作為取用者群組。
相依性設定
<dependency>
<groupId>com.azure.spring</groupId>
<artifactId>spring-cloud-azure-stream-binder-servicebus</artifactId>
</dependency>
或者,您也可以使用 Spring Cloud Azure Stream 服務總線入門版,如下列 Maven 範例所示:
<dependency>
<groupId>com.azure.spring</groupId>
<artifactId>spring-cloud-azure-starter-stream-servicebus</artifactId>
</dependency>
配置
系結器提供下列兩個組態選項部分:
聯機組態屬性
本節包含用來連線到 Azure 服務總線的組態選項。
注意
如果您選擇使用安全性主體向 Microsoft Entra ID 進行驗證和授權,以存取 Azure 資源,請參閱 使用 Microsoft Entra ID 授權存取,以確保安全性主體已獲得存取 Azure 資源的足夠許可權。
spring-cloud-azure-stream-binder-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 服務總線命名空間值的功能變數名稱。 |
注意
一般 Azure 服務 SDK 組態選項也可以針對 Spring Cloud Azure Stream 服務總線系結器進行設定。 支援的組態選項已於 Spring Cloud Azure 組態中介紹,且可使用統一前綴 spring.cloud.azure. 或 spring.cloud.azure.servicebus. 前綴進行設定。
系結器預設也支援 spring Could Azure Resource Manager 。 若要了解如何擷取未獲授與 相關角色的安全性主體之連接字串,請參閱 Spring Could Azure Resource Manager 的 基本用法 一節。
Azure 服務總線系結組態屬性
下列選項分為四個區段:取用者屬性、進階取用者組態、產生者屬性和進階生產者設定。
消費者屬性
這些屬性會透過 ServiceBusConsumerProperties公開。
注意
為了避免重複,自 4.17.0 和 5.11.0 版起,Spring Cloud Azure Stream Binder 服務匯流排 支援為所有通道設定值,格式為 spring.cloud.stream.servicebus.default.consumer.<property>=<value>。
spring-cloud-azure-stream-binder-servicebus 的使用者可設定屬性:
| 財產 | 類型 | 預設 | 描述 |
|---|---|---|---|
spring.cloud.stream.servicebus.bindings.binding-name.consumer.requeue-rejected |
布爾 | 假 | 如果失敗訊息會被路由至 DLQ。 |
spring.cloud.stream.servicebus.bindings.binding-name.consumer.max-concurrent-calls |
整數 | 1 | 服務總線處理器客戶端應該處理的最大並行訊息。 啟用工作階段時,它會套用至每個工作階段。 |
spring.cloud.stream.servicebus.bindings.binding-name.consumer.max-concurrent-sessions |
整數 | null | 在任何指定時間要處理的並行會話數目上限。 |
spring.cloud.stream.servicebus.bindings.binding-name.consumer.session-enabled |
布林值 | null | 工作階段是否已啟用。 |
spring.cloud.stream.servicebus.bindings.binding-name.consumer.session-idle-timeout |
期間 | null | 設定等待目前作用中的工作階段接收到訊息的最長時間(持續時間)。 |
spring.cloud.stream.servicebus.bindings.binding-name.consumer.prefetch-count |
整數 | 0 | 服務匯流排 處理器用戶端的預先擷取計數。 |
spring.cloud.stream.servicebus.bindings.binding-name.consumer.sub-queue |
子佇列 | 沒有 | 要連接的子佇列類型。 |
spring.cloud.stream.servicebus.bindings.binding-name.consumer.max-auto-lock-renew-duration |
期間 | 5分鐘 | 持續自動延長鎖定的時間長度。 |
spring.cloud.stream.servicebus.bindings.binding-name.consumer.receive-mode |
ServiceBusReceiveMode | peek_lock | 服務總線處理器用戶端的接收模式。 |
spring.cloud.stream.servicebus.bindings.binding-name.consumer.auto-complete |
布林值 | true | 是否要自動解決訊息。 如果設定為 false,將會新增 Checkpointer 的訊息標頭,讓開發人員手動解決訊息。 |
spring.cloud.stream.servicebus.bindings.binding-name.consumer.max-size-in-megabytes |
長 | 1024 | 佇列/主題的大小上限,以 MB 為單位,這是為佇列/主題配置的記憶體大小。 |
spring.cloud.stream.servicebus.bindings.binding-name.consumer.default-message-time-to-live |
期間 | P10675199DT2H48M5.4775807S。 (10675199天、2小時、48分、5秒和477毫秒) | 訊息到期的持續時間,從訊息傳送至服務總線時開始。 |
重要
當您使用 Azure Resource Manager (ARM) 時,您必須設定 屬性。 欲了解更多資訊,請參閱 servicebus-queue-binder-arm GitHub 上的範例。
進階使用者設定
上述 連線組態和 通用 Azure SDK 用戶端組態支援為各個繫結器取用者進行自訂,您可以使用前置詞 spring.cloud.stream.servicebus.bindings.<binding-name>.consumer. 來設定。
生產者屬性
這些屬性會透過 ServiceBusProducerProperties公開。
注意
為了避免重複,自 4.17.0 和 5.11.0 版起,Spring Cloud Azure Stream Binder 服務匯流排 支援為所有通道設定值,格式為 spring.cloud.stream.servicebus.default.producer.<property>=<value>。
spring-cloud-azure-stream-binder-servicebus 的生產者可設定屬性:
| 財產 | 類型 | 預設 | 描述 |
|---|---|---|---|
spring.cloud.stream.servicebus.bindings.binding-name.producer.sync |
布爾 | 假 | 用於切換生產者同步的旗標。 |
spring.cloud.stream.servicebus.bindings.binding-name.producer.send-timeout |
長 | 一萬 | 生產者傳送的逾時值。 |
spring.cloud.stream.servicebus.bindings.binding-name.producer.entity-type |
ServiceBusEntityType | null | 產生者的 服務匯流排 實體類型,為繫結產生者所必需。 |
spring.cloud.stream.servicebus.bindings.binding-name.producer.max-size-in-megabytes |
長 | 1024 | 佇列/主題的大小上限,以 MB 為單位,這是為佇列/主題配置的記憶體大小。 |
spring.cloud.stream.servicebus.bindings.binding-name.producer.default-message-time-to-live |
期間 | P10675199DT2H48M5.4775807S。 (10675199天、2小時、48分、5秒和477毫秒) | 訊息到期的持續時間,從訊息傳送至服務總線時開始。 |
重要
使用系結產生者時,必須設定 spring.cloud.stream.servicebus.bindings.<binding-name>.producer.entity-type 的屬性。
進階生產者設定
上述 連線和 通用 Azure SDK 用戶端組態支援針對每個繫結器生產者進行自訂,您可以使用前綴 spring.cloud.stream.servicebus.bindings.<binding-name>.producer. 進行設定。
基本用法
從/到服務總線傳送和接收訊息
使用認證資訊填入組態選項。
針對認證作為連接字串,請在 application.yml 檔案中設定下列屬性:
spring: cloud: azure: servicebus: connection-string: ${SERVICEBUS_NAMESPACE_CONNECTION_STRING} function: definition: consume;supply stream: bindings: consume-in-0: destination: ${SERVICEBUS_ENTITY_NAME} # If you use Service Bus Topic, add the following configuration # group: ${SUBSCRIPTION_NAME} supply-out-0: destination: ${SERVICEBUS_ENTITY_NAME_SAME_AS_ABOVE} servicebus: bindings: consume-in-0: consumer: auto-complete: false supply-out-0: producer: entity-type: queue # set as "topic" if you use Service Bus Topic注意
Microsoft 建議您使用最安全的可用驗證流程。 此程式中所述的驗證流程,例如資料庫、快取、傳訊或 AI 服務,在應用程式中需要高度的信任,而且不會在其他流程中帶來風險。 只有在更安全的選項(例如使用受控識別進行無密碼或無金鑰連線)不可行時,才使用此流程。 針對本機計算機作業,偏好使用無密碼或無密鑰連線的使用者身分識別。
若使用服務主體認證,請在 application.yml 檔案中設定下列屬性:
spring: cloud: azure: credential: client-id: ${AZURE_CLIENT_ID} client-secret: ${AZURE_CLIENT_SECRET} profile: tenant-id: <tenant> servicebus: namespace: ${SERVICEBUS_NAMESPACE} function: definition: consume;supply stream: bindings: consume-in-0: destination: ${SERVICEBUS_ENTITY_NAME} # If you use Service Bus Topic, add the following configuration # group: ${SUBSCRIPTION_NAME} supply-out-0: destination: ${SERVICEBUS_ENTITY_NAME_SAME_AS_ABOVE} servicebus: bindings: consume-in-0: consumer: auto-complete: false supply-out-0: producer: entity-type: queue # set as "topic" if you use Service Bus Topic
注意
tenant-id 允許的值包括:common、organizations、consumers或租用戶標識碼。 如需這些值的詳細資訊,請參閱 錯誤 AADSTS50020 - 來自身分識別提供者的使用者帳戶不存在於租用戶中中的 使用了錯誤的端點(個人和組織帳戶)一節。 如需有關將單一租用戶應用程式轉換為多租用戶的資訊,請參閱 在 Microsoft Entra ID 上將單一租用戶應用程式轉換為多租用戶。
針對認證作為受控識別,請在 application.yml 檔案中設定下列屬性:
spring: cloud: azure: credential: managed-identity-enabled: true client-id: ${MANAGED_IDENTITY_CLIENT_ID} # Only needed when using a user-assigned managed identity servicebus: namespace: ${SERVICEBUS_NAMESPACE} function: definition: consume;supply stream: bindings: consume-in-0: destination: ${SERVICEBUS_ENTITY_NAME} # If you use Service Bus Topic, add the following configuration # group: ${SUBSCRIPTION_NAME} supply-out-0: destination: ${SERVICEBUS_ENTITY_NAME_SAME_AS_ABOVE} servicebus: bindings: consume-in-0: consumer: auto-complete: false supply-out-0: producer: entity-type: queue # set as "topic" if you use Service Bus Topic
定義供應商和取用者。
@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(); }; }
分割區索引鍵支援
系結器支援 服務總線分割,方法是在訊息標頭中設定分割區索引鍵和會話標識符。 本節將介紹如何為訊息設定分割區索引鍵。
Spring Cloud Stream 提供分區鍵 SpEL 運算式屬性 spring.cloud.stream.bindings.<binding-name>.producer.partition-key-expression。 例如,將此屬性設定為 "'partitionKey-' + headers[<message-header-key>]",並新增名為 message-header-key 的標頭。 在評估表達式以指派分割區索引鍵時,Spring Cloud Stream 會使用此標頭的值。 下列程式碼提供了一個生產者範例:
@Bean
public Supplier<Message<String>> generate() {
return () -> {
String value = "random payload";
return MessageBuilder.withPayload(value)
.setHeader("<message-header-key>", value.length() % 4)
.build();
};
}
會話支援
系結器支援 服務匯流排 的訊息工作階段。 訊息的會話標識碼可以透過訊息標頭來設定。
@Bean
public Supplier<Message<String>> generate() {
return () -> {
String value = "random payload";
return MessageBuilder.withPayload(value)
.setHeader(ServiceBusMessageHeaders.SESSION_ID, "Customize session ID")
.build();
};
}
注意
根據 服務總線分割,會話標識符的優先順序高於分割區索引鍵。 因此,當同時設定 ServiceBusMessageHeaders#SESSION_ID 和 ServiceBusMessageHeaders#PARTITION_KEY 標頭時,會話標識碼的值最終會用來覆寫分割區索引鍵的值。
處理錯誤訊息
處理輸出系結錯誤訊息
根據預設,Spring Integration 會建立名為
errorChannel的全域錯誤通道。 設定下列訊息端點來處理輸出系結錯誤訊息。@ServiceActivator(inputChannel = IntegrationContextUtils.ERROR_CHANNEL_BEAN_NAME) public void handleError(ErrorMessage message) { LOGGER.error("Handling outbound binding error: " + message); }處理輸入系結錯誤訊息
Spring Cloud Stream 服務總線系結器支援兩個解決方案來處理輸入訊息系結的錯誤:系結器錯誤處理程式和處理程式。
Binder 錯誤處理程式:
默認系結器錯誤處理程式會處理輸入系結。 當啟用
spring.cloud.stream.servicebus.bindings.<binding-name>.consumer.requeue-rejected時,您可以使用此處理常式將失敗的訊息傳送至死信佇列。 否則,這些失敗的訊息會被捨棄。 系結器錯誤處理程式與其他提供的錯誤處理程式互斥。錯誤處理程式:
Spring Cloud Stream 會公開一種機制,讓您藉由新增可接受
Consumer實例的ErrorMessage來提供自定義錯誤處理程式。 如需詳細資訊,請參閱 Spring Cloud Stream 檔中 處理錯誤訊息。系結預設錯誤處理程式
設定單一
ConsumerBean 以處理所有傳入繫結錯誤訊息。 下列預設函式會訂閱每個輸入系結錯誤通道:@Bean public Consumer<ErrorMessage> myDefaultHandler() { return message -> { // consume the error message }; }您也需要將
spring.cloud.stream.default.error-handler-definition屬性設定為函式名稱。系結特定錯誤處理程式
設定
ConsumerBean,以處理特定傳入繫結的錯誤訊息。 下列函式會訂閱優先順序高於系結預設錯誤處理程式的特定輸入系結錯誤通道。@Bean public Consumer<ErrorMessage> myDefaultHandler() { return message -> { // consume the error message }; }您也需要將
spring.cloud.stream.bindings.<input-binding-name>.error-handler-definition屬性設定為函式名稱。
服務總線訊息標頭
如需了解所支援的基本訊息標頭,請參閱 Spring Cloud Azure 對 Spring Integration 的支援中的 服務匯流排 訊息標頭一節。
注意
設定分割區索引鍵時,訊息標頭的優先順序高於 Spring Cloud Stream 屬性。 因此,只有當未設定任何 spring.cloud.stream.bindings.<binding-name>.producer.partition-key-expression 和 ServiceBusMessageHeaders#SESSION_ID 標頭時,ServiceBusMessageHeaders#PARTITION_KEY 才會生效。
多個系結器支援
也支援透過使用多個繫結器來連線至多個 服務匯流排 命名空間。 此範例會採用連接字串作為範例。 也支援服務主體和受控識別的認證,用戶可以在每個系結器的環境設定中設定相關的屬性。
若要使用 ServiceBus 的多個系結器,請在 application.yml 檔案中設定下列屬性:
spring: cloud: function: definition: consume1;supply1;consume2;supply2 stream: bindings: consume1-in-0: destination: ${SERVICEBUS_TOPIC_NAME} group: ${SUBSCRIPTION_NAME} supply1-out-0: destination: ${SERVICEBUS_TOPIC_NAME_SAME_AS_ABOVE} consume2-in-0: binder: servicebus-2 destination: ${SERVICEBUS_QUEUE_NAME} supply2-out-0: binder: servicebus-2 destination: ${SERVICEBUS_QUEUE_NAME_SAME_AS_ABOVE} binders: servicebus-1: type: servicebus default-candidate: true environment: spring: cloud: azure: servicebus: connection-string: ${SERVICEBUS_NAMESPACE_01_CONNECTION_STRING} servicebus-2: type: servicebus default-candidate: false environment: spring: cloud: azure: servicebus: connection-string: ${SERVICEBUS_NAMESPACE_02_CONNECTION_STRING} servicebus: bindings: consume1-in-0: consumer: auto-complete: false supply1-out-0: producer: entity-type: topic consume2-in-0: consumer: auto-complete: false supply2-out-0: producer: entity-type: queue poller: initial-delay: 0 fixed-delay: 1000注意
前述的應用程式檔案顯示如何設定單一預設輪詢器,並將其套用至所有繫結。 如果您想為特定繫結設定輪詢器,可以使用如下設定,例如
spring.cloud.stream.bindings.<binding-name>.producer.poller.fixed-delay=3000。注意
Microsoft 建議您使用最安全的可用驗證流程。 此程式中所述的驗證流程,例如資料庫、快取、傳訊或 AI 服務,在應用程式中需要高度的信任,而且不會在其他流程中帶來風險。 只有在更安全的選項(例如使用受控識別進行無密碼或無金鑰連線)不可行時,才使用此流程。 針對本機計算機作業,偏好使用無密碼或無密鑰連線的使用者身分識別。
我們需要定義兩個供應商和兩個消費者
@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(); }; }
資源布建
Service Bus 繫結器支援佈建佇列、主題和訂用帳戶,使用者可使用下列屬性來啟用佈建功能。
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}
注意
tenant-id 允許的值包括:common、organizations、consumers或租用戶標識碼。 如需這些值的詳細資訊,請參閱 錯誤 AADSTS50020 - 來自身分識別提供者的使用者帳戶不存在於租用戶中中的 使用了錯誤的端點(個人和組織帳戶)一節。 如需有關將單一租用戶應用程式轉換為多租用戶的資訊,請參閱 在 Microsoft Entra ID 上將單一租用戶應用程式轉換為多租用戶。
自訂服務總線客戶端屬性
開發人員可以使用 AzureServiceClientBuilderCustomizer 來自定義服務總線用戶端屬性。 下列範例會自訂 sessionIdleTimeout中的 ServiceBusClientBuilder 屬性:
@Bean
public AzureServiceClientBuilderCustomizer<ServiceBusClientBuilder.ServiceBusSessionProcessorClientBuilder> customizeBuilder() {
return builder -> builder.sessionIdleTimeout(Duration.ofSeconds(10));
}
樣品
欲了解更多資訊,請參閱 azure-spring-boot-samples GitHub 上的倉庫。