Azure 服務總線 是 Microsoft 提供的一種高度可靠的雲消息傳送服務。
使用應用程式中的用戶端庫@azure/service-bus
- 將消息發送到 Azure 服務總線佇列或主題
- 從 Azure 服務總線佇列或訂閱接收消息
- 在 Azure 服務總線命名空間中創建/獲取/刪除/更新/列出佇列/主題/訂閱/規則。
版本 7 的資源 @azure/service-bus :
主要連結:
注意:如果您使用的是 1.1.10 或更低版本,並且想要遷移到此包的最新版本,請查看我們的 遷移指南,從 服務匯流排 V1 遷移到 服務匯流排 V7
開始使用
安裝套件
使用 npm 安裝 Azure 服務總線用戶端庫的最新版本。
npm install @azure/service-bus
目前支援的環境
Prerequisites
設定 TypeScript
TypeScript 用戶必須安裝 Node 類型定義:
npm install @types/node
您也需要在 tsconfig.json中啟用 compilerOptions.allowSyntheticDefaultImports。 請注意,如果您已啟用 compilerOptions.esModuleInterop,預設會啟用 allowSyntheticDefaultImports。 如需詳細資訊,請參閱 TypeScript 的編譯程式選項手冊。
JavaScript 套件組合
若要在瀏覽器中使用此用戶端連結庫,您必須先使用配套程式。 如需如何執行這項操作的詳細資訊,請參閱我們的 組合檔。
除了該處所描述的內容之外,此連結庫還需要下列 NodeJS 核心內建模組的額外 polyfills,才能在瀏覽器中正常運作:
bufferospathprocess
結合 Webpack
如果您使用 Webpack v5,您可以安裝下列開發相依性
npm install --save-dev os-browserify path-browserify
然後將下列內容新增至您的 webpack.config.js
const path = require("path");
+const webpack = require("webpack");
module.exports = {
entry: "./src/index.ts",
@@ -12,8 +13,21 @@ module.exports = {
},
],
},
+ plugins: [
+ new webpack.ProvidePlugin({
+ process: "process/browser",
+ }),
+ new webpack.ProvidePlugin({
+ Buffer: ["buffer", "Buffer"],
+ }),
+ ],
resolve: {
extensions: [".ts", ".js"],
+ fallback: {
+ buffer: require.resolve("buffer/"),
+ os: require.resolve("os-browserify"),
+ path: require.resolve("path-browserify"),
+ },
},
組合與匯總
如果您使用匯總套件組合器,請安裝下列開發相依性
npm install --save-dev @rollup/plugin-commonjs @rollup/plugin-inject @rollup/plugin-node-resolve
然後在您的 rollup.config.js 中包含下列內容
+import nodeResolve from "@rollup/plugin-node-resolve";
+import cjs from "@rollup/plugin-commonjs";
+import shim from "rollup-plugin-shim";
+import inject from "@rollup/plugin-inject";
export default {
// other configs
plugins: [
+ shim({
+ fs: `export default {}`,
+ net: `export default {}`,
+ tls: `export default {}`,
+ path: `export default {}`,
+ dns: `export function resolve() { }`,
+ }),
+ nodeResolve({
+ mainFields: ["module", "browser"],
+ preferBuiltins: false,
+ }),
+ cjs(),
+ inject({
+ modules: {
+ Buffer: ["buffer", "Buffer"],
+ process: "process",
+ },
+ exclude: ["./**/package.json"],
+ }),
]
};
如需使用 polyfills 的詳細資訊,請參閱您慣用套件組合的檔。
React 原生支援
與瀏覽器類似,React Native 不支援此 SDK 庫使用的某些 JavaScript API,因此您需要為它們提供 polyfill。 如需詳細資訊,請參閱 Messaging React Native 範例與 Expo。
驗證客戶端
與 服務匯流排 的交互從 ServiceBusClient 類的實例開始。 可以使用連接字串或使用 Azure Active Directory 憑據向 服務匯流排 進行身份驗證。
使用連接字串
此方法將連接字串引入 服務匯流排 實例。 您可以從 Azure 入口網站取得連接字串。
import { ServiceBusClient } from "@azure/service-bus";
const serviceBusClient = new ServiceBusClient("<connectionString>");
有關此構造函數的更多資訊,請參閱 API 文件。
使用 Azure Active Directory 認證
使用 Azure Active Directory 進行身份驗證使用 Azure 標識庫。
下面的示例使用 DefaultAzureCredential,這是庫中許多可用憑據提供程式 @azure/identity 之一。
import { DefaultAzureCredential } from "@azure/identity";
import { ServiceBusClient } from "@azure/service-bus";
const fullyQualifiedNamespace = "<name-of-service-bus-namespace>.servicebus.windows.net";
const credential = new DefaultAzureCredential();
const serviceBusClient = new ServiceBusClient(fullyQualifiedNamespace, credential);
注意:如果您使用自己的介面實現
TokenCredential來對抗 AAD,請將 service-bus 的“scopes”設置為以下內容以獲取相應的令牌:
["https://servicebus.azure.net//user_impersonation"];
有關此構造函數的更多資訊,請參閱 API 文件
關鍵概念
初始化 后, ServiceBusClient您可以在 服務匯流排 命名空間中與以下資源進行交互:
- 佇列:允許發送和接收消息。 通常用於點對點通信。
- 主題:與佇列相反,主題更適合發佈/訂閱場景。 主題可以發送到主題,但需要一個訂閱,其中可以有多個並行訂閱才能使用。
- Subscriptions:從 Topic 消費的機制。 每個訂閱都是獨立的,並接收發送到主題的每條消息的副本。 Rules 和 Filters 可用於定製特定訂閱接收的消息。
有關這些資源的詳細資訊,請參閱 什麼是 Azure 服務總線?。
要與這些資源交互,應熟悉以下 SDK 概念:
- 使用
ServiceBusSender建立的ServiceBusClient.createSender()將消息發送到佇列或主題。 - 使用
ServiceBusReceiver使用 創建的 從ServiceBusClient.createReceiver()佇列或訂閱接收消息。 - 使用
ServiceBusSessionReceiver使用 或ServiceBusClient.acceptSession()創建的 從ServiceBusClient.acceptNextSession()啟用會話的佇列或訂閱接收消息。
請注意,Queues、Topics 和 Subscriptions 應在使用此庫之前創建。
Examples
以下部分提供了涵蓋使用 Azure 服務總線的一些常見任務的代碼片段
傳送訊息
創建類的ServiceBusClient實例后,可以使用 ServiceBusSender 方法獲取可用於發送消息的實例。
import { DefaultAzureCredential } from "@azure/identity";
import { ServiceBusClient } from "@azure/service-bus";
const fullyQualifiedNamespace = "<name-of-service-bus-namespace>.servicebus.windows.net";
const credential = new DefaultAzureCredential();
const serviceBusClient = new ServiceBusClient(fullyQualifiedNamespace, credential);
const sender = serviceBusClient.createSender("my-queue");
const messages = [
{ body: "Albert Einstein" },
{ body: "Werner Heisenberg" },
{ body: "Marie Curie" },
{ body: "Steven Hawking" },
{ body: "Isaac Newton" },
{ body: "Niels Bohr" },
{ body: "Michael Faraday" },
{ body: "Galileo Galilei" },
{ body: "Johannes Kepler" },
{ body: "Nikolaus Kopernikus" },
];
// sending a single message
await sender.sendMessages(messages[0]);
// sending multiple messages in a single call
// this will fail if the messages cannot fit in a batch
await sender.sendMessages(messages);
// Sends multiple messages using one or more ServiceBusMessageBatch objects as required
let batch = await sender.createMessageBatch();
for (let i = 0; i < messages.length; i++) {
const message = messages[i];
if (!batch.tryAddMessage(message)) {
// Send the current batch as it is full and create a new one
await sender.sendMessages(batch);
batch = await sender.createMessageBatch();
if (!batch.tryAddMessage(messages[i])) {
throw new Error("Message too big to fit in a batch");
}
}
}
// Send the batch
await sender.sendMessages(batch);
接收訊息
建立類的ServiceBusClient實例後,可以使用 ServiceBusReceiver 方法獲取 。
import { DefaultAzureCredential } from "@azure/identity";
import { ServiceBusClient } from "@azure/service-bus";
const fullyQualifiedNamespace = "<name-of-service-bus-namespace>.servicebus.windows.net";
const credential = new DefaultAzureCredential();
const serviceBusClient = new ServiceBusClient(fullyQualifiedNamespace, credential);
const receiver = serviceBusClient.createReceiver("my-queue");
receiveMode有兩種可用。
- “peekLock” - 在 peekLock 模式下,接收方在佇列上指定的持續時間內對消息具有鎖定。
- “receiveAndDelete” - 在 receiveAndDelete 模式下,收到消息時,將從服務總線中刪除消息。
如果選項中未提供 receiveMode,則預設為 「peekLock」 模式。 您還可以在 「peekLock」 模式下 結算收到的消息 。
您可以透過以下 3 種方式之一使用此接收器來接收訊息:
獲取消息陣列
使用 receiveMessages 函數,該函數返回一個解析為消息陣列的 Promise。
import { DefaultAzureCredential } from "@azure/identity";
import { ServiceBusClient } from "@azure/service-bus";
const fullyQualifiedNamespace = "<name-of-service-bus-namespace>.servicebus.windows.net";
const credential = new DefaultAzureCredential();
const serviceBusClient = new ServiceBusClient(fullyQualifiedNamespace, credential);
const receiver = serviceBusClient.createReceiver("my-queue");
const myMessages = await receiver.receiveMessages(10);
使用消息處理程式訂閱
使用 subscribe 方法設置消息處理程式,並使其在需要時運行。
完成後,致電 receiver.close() 停止接收更多消息。
import { DefaultAzureCredential } from "@azure/identity";
import { ServiceBusClient } from "@azure/service-bus";
const fullyQualifiedNamespace = "<name-of-service-bus-namespace>.servicebus.windows.net";
const credential = new DefaultAzureCredential();
const serviceBusClient = new ServiceBusClient(fullyQualifiedNamespace, credential);
const receiver = serviceBusClient.createReceiver("my-queue");
const myMessageHandler = async (message) => {
// your code here
console.log(`message.body: ${message.body}`);
};
const myErrorHandler = async (args) => {
console.log(
`Error occurred with ${args.entityPath} within ${args.fullyQualifiedNamespace}: `,
args.error,
);
};
receiver.subscribe({
processMessage: myMessageHandler,
processError: myErrorHandler,
});
使用 async iterator
使用 getMessageIterator 獲取消息的異步反覆運算器
import { DefaultAzureCredential } from "@azure/identity";
import { ServiceBusClient } from "@azure/service-bus";
const fullyQualifiedNamespace = "<name-of-service-bus-namespace>.servicebus.windows.net";
const credential = new DefaultAzureCredential();
const serviceBusClient = new ServiceBusClient(fullyQualifiedNamespace, credential);
const receiver = serviceBusClient.createReceiver("my-queue");
for await (const message of receiver.getMessageIterator()) {
// your code here
}
結算消息
收到消息后,您可以根據您希望如何結算消息來調用completeMessage()abandonMessage()deferMessage()、 、 或 deadLetterMessage() 接收方。
要瞭解更多資訊,請閱讀 結算收到的消息
無效信件佇列
死信佇列是一個 子佇列。 每個佇列或訂閱都有自己的死信佇列。 死信佇列存儲已明確死信 (via receiver.deadLetterMessage()) 的消息或已超過其最大送達計數的消息。
為死信子佇列建立接收器類似於為訂閱或佇列建立接收器:
import { DefaultAzureCredential } from "@azure/identity";
import { ServiceBusClient } from "@azure/service-bus";
const fullyQualifiedNamespace = "<name-of-service-bus-namespace>.servicebus.windows.net";
const credential = new DefaultAzureCredential();
const serviceBusClient = new ServiceBusClient(fullyQualifiedNamespace, credential);
// To receive from a queue's dead letter sub-queue
const deadLetterReceiverForQueue = serviceBusClient.createReceiver("queue", {
subQueueType: "deadLetter",
});
// To receive from a subscription's dead letter sub-queue
const deadLetterReceiverForSubscription = serviceBusClient.createReceiver("topic", "subscription", {
subQueueType: "deadLetter",
});
// Dead letter receivers work like any other receiver connected to a queue
// ex:
const messages = await deadLetterReceiverForQueue.receiveMessages(5);
for (const message of messages) {
console.log(`Dead lettered message: ${message.body}`);
}
更徹底地演示死信佇列的完整示例:
使用 Sessions 發送消息
使用會話需要您創建啟用會話的 Queue 或 Subscription。 您可以 在此處閱讀有關如何在門戶中配置此功能的更多資訊。
為了向會話發送消息,請使用 ServiceBusClient 創建發件人 createSender.
發送消息時,請在消息中設置 sessionId 屬性,以確保您的消息到達正確的會話。
import { DefaultAzureCredential } from "@azure/identity";
import { ServiceBusClient } from "@azure/service-bus";
const fullyQualifiedNamespace = "<name-of-service-bus-namespace>.servicebus.windows.net";
const credential = new DefaultAzureCredential();
const serviceBusClient = new ServiceBusClient(fullyQualifiedNamespace, credential);
const sender = serviceBusClient.createSender("my-session-queue");
await sender.sendMessages({
body: "my-message-body",
sessionId: "my-session",
});
您可以 在此處閱讀有關工作階段如何運作的更多資訊。
從 Session 接收消息
使用會話需要您創建啟用會話的 Queue 或 Subscription。 您可以 在此處閱讀有關如何在門戶中配置此功能的更多資訊。
與未啟用會話的 Queues 或 Subscriptions 不同,任何時候都只有一個接收器可以從會話中讀取數據。 這是通過 鎖定 會話來強制實施的,該會話由 服務匯流排 處理。 從概念上講,這類似於使用 peekLock mode時消息鎖定的工作原理 - 當消息(或會話)被鎖定時,您的接收者擁有對它的獨佔訪問許可權。
為了打開和鎖定會話,請使用的 ServiceBusClient 實例來創建 SessionReceiver。
有兩種方法可以選擇要打開的工作階段:
sessionId指定 ,用於鎖定命名會話。const receiver = await serviceBusClient.acceptSession("my-session-queue", "my-session");請勿指定會話 ID。在這種情況下,服務總線將查找下一個尚未鎖定的可用會話。
const receiver = await serviceBusClient.acceptNextSession("my-session-queue");您可以通過 上的屬性
sessionIdSessionReceiver找到工作階段的名稱。 如果選項中未提供 receiveMode,則預設為 「peekLock」 模式。 您還可以在 「peekLock」 模式下 結算收到的消息 。
建立接收者後,您可以選擇 3 種方式來接收消息:
您可以 在此處閱讀有關工作階段如何運作的更多資訊。
列表訊息會話
要發現佇列或訂閱中哪些會話有活躍訊息或會話狀態,請使用 listMessageSessions():
import { DefaultAzureCredential } from "@azure/identity";
import { ServiceBusClient } from "@azure/service-bus";
const fullyQualifiedNamespace = "<name-of-service-bus-namespace>.servicebus.windows.net";
const credential = new DefaultAzureCredential();
const serviceBusClient = new ServiceBusClient(fullyQualifiedNamespace, credential);
// List all sessions with active messages or session state in a queue
for await (const sessionId of serviceBusClient.listMessageSessions("my-session-queue")) {
console.log("Session ID:", sessionId);
}
// List sessions in a subscription
for await (const sessionId of serviceBusClient.listMessageSessions("my-topic", "my-subscription")) {
console.log("Session ID:", sessionId);
}
// List only sessions whose stored session state was set or updated in the last seven days
const sessionStateUpdatedAfter = new Date(Date.now() - 7 * 24 * 60 * 60 * 1000);
for await (const sessionId of serviceBusClient.listMessageSessions("my-session-queue", {
sessionStateUpdatedAfter,
})) {
console.log("Recently updated session ID:", sessionId);
}
管理服務總線命名空間的資源
ServiceBusAdministrationClient 允許您對實體(佇列、主題和訂閱)和訂閱的規則執行 CRUD作來管理命名空間。
- 支援使用服務總線連接字串以及類似於
@azure/identityServiceBusClient的 AAD 憑據進行身份驗證。
注意:服務總線尚不支援為命名空間設置 CORS 規則,因此 ServiceBusAdministrationClient 如果不禁用 Web 安全性,則無法在瀏覽器中工作。 有關更多資訊,請參閱 此處。
import { ServiceBusAdministrationClient } from "@azure/service-bus";
const queueName = "my-session-queue";
// Get the connection string from the portal
// OR
// use the token credential overload, provide the host name of your Service Bus instance and the AAD credentials from the @azure/identity library
const serviceBusAdministrationClient = new ServiceBusAdministrationClient("<connectionString>");
// Similarly, you can create topics and subscriptions as well.
const createQueueResponse = await serviceBusAdministrationClient.createQueue(queueName);
console.log("Created queue with name - ", createQueueResponse.name);
const queueRuntimeProperties =
await serviceBusAdministrationClient.getQueueRuntimeProperties(queueName);
console.log(`Number of messages in the queue = ${queueRuntimeProperties.totalMessageCount}`);
// Topic runtime properties additionally report the total number of SQL and correlation filters
// across all of the topic's subscriptions. These counts are served by the 2024-05 service API
// version and later; on an older version they are `undefined`.
const topicName = "my-topic";
const subscriptionName = "my-subscription";
await serviceBusAdministrationClient.createTopic(topicName);
// A new subscription carries a default rule with a SQL TrueFilter. Adding a correlation rule
// gives the topic one of each, so the counts below aggregate across the subscription's rules.
await serviceBusAdministrationClient.createSubscription(topicName, subscriptionName);
await serviceBusAdministrationClient.createRule(
topicName,
subscriptionName,
"my-correlation-rule",
{ correlationId: "my-correlation-id" },
);
const topicRuntimeProperties =
await serviceBusAdministrationClient.getTopicRuntimeProperties(topicName);
console.log(`SQL filter count = ${topicRuntimeProperties.sqlFilterCount}`);
console.log(`Correlation filter count = ${topicRuntimeProperties.correlationFilterCount}`);
await serviceBusAdministrationClient.deleteTopic(topicName);
await serviceBusAdministrationClient.deleteQueue(queueName);
- 參考示例 - administrationClient.ts
Troubleshooting
以下是開始診斷問題的一些初始步驟。 有關更多資訊,請參閱 服務匯流排 故障排除指南。
AMQP 相依性
服務總線庫依賴於 rhea-promise 庫來管理連接、通過 AMQP 協定發送和接收消息。
Logging
你可以設定以下環境變數,使用這個函式庫時取得除錯日誌。
- 從服務總線 SDK 獲取調試日誌
export DEBUG=azure*
- 從服務總線 SDK 和協定級庫獲取調試日誌。
export DEBUG=azure*,rhea*
- 如果您對 查看消息轉換( 佔用大量控制台/磁碟空間)不感興趣,則可以按如下方式設置
DEBUG環境變數:
export DEBUG=azure*,rhea*,-rhea:raw,-rhea:message,-azure:core-amqp:datatransformer
- 如果您只對 錯誤感興趣,則可以按如下方式設定
DEBUG環境變數:
export DEBUG=azure:service-bus:error,azure:core-amqp:error,rhea-promise:error,rhea:events,rhea:frames,rhea:io,rhea:flow
記錄至檔案
- 如上所示設置
DEBUG環境變數 - 按如下方式執行測試文稿:
- 測試文稿中的記錄語句會移至
out.log,而 sdk 的記錄語句會移至debug.log。node your-test-script.js > out.log 2>debug.log - 從測試腳本和 sdk 記錄語句,將 stderr 重新導向至 stdout (&1),然後將 stdout 重新導向至檔案,以移至相同的檔案
out.log:node your-test-script.js >out.log 2>&1 - 從測試文稿記錄語句,sdk 會移至相同的檔案
out.log。node your-test-script.js &> out.log
下一步
請查看 範例 目錄,瞭解有關如何使用此庫向 服務總線佇列、主題和訂閱發送消息和從服務總線佇列、主題和訂閱發送消息的詳細示例。
Contributing
如果你想為這個函式庫貢獻,請閱讀 contributing guide,了解更多如何建置與測試程式碼。