Azure 服務匯流排 client library for JavaScript - version 7.10.0

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,才能在瀏覽器中正常運作:

  • buffer
  • os
  • path
  • process

結合 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 概念:

請注意,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。

有兩種方法可以選擇要打開的工作階段:

  1. sessionId指定 ,用於鎖定命名會話。

    const receiver = await serviceBusClient.acceptSession("my-session-queue", "my-session");
    
  2. 請勿指定會話 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);

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

記錄至檔案

  1. 如上所示設置 DEBUG 環境變數
  2. 按如下方式執行測試文稿:
  • 測試文稿中的記錄語句會移至 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,了解更多如何建置與測試程式碼。