Azure Service Bus TypeScript

azure-servicebus-ts
分类通用
作者Agentic Awesome Skills 社区
许可MIT
评分4.50/5
使用10.1K

Azure Service Bus TypeScript SDK

支持队列、主题和订阅的企业级消息传递。

安装

bash
npm install @azure/service-bus @azure/identity

环境变量

bash
SERVICEBUS_NAMESPACE=<namespace>.servicebus.windows.net
SERVICEBUS_QUEUE_NAME=my-queue
SERVICEBUS_TOPIC_NAME=my-topic
SERVICEBUS_SUBSCRIPTION_NAME=my-subscription

身份验证

typescript
import { ServiceBusClient } from "@azure/service-bus";
import { DefaultAzureCredential } from "@azure/identity";

const fullyQualifiedNamespace = process.env.SERVICEBUS_NAMESPACE!;
const client = new ServiceBusClient(fullyQualifiedNamespace, new DefaultAzureCredential());

核心工作流

向队列发送消息

typescript
const sender = client.createSender("my-queue");

// 单条消息
await sender.sendMessages({
body: { orderId: "12345", amount: 99.99 },
contentType: "application/json",
});

// 批量消息
const batch = await sender.createMessageBatch();
batch.tryAddMessage({ body: "Message 1" });
batch.tryAddMessage({ body: "Message 2" });
await sender.sendMessages(batch);

await sender.close();

从队列接收消息

typescript
const receiver = client.createReceiver("my-queue");

// 批量接收
const messages = await receiver.receiveMessages(10, { maxWaitTimeInMs: 5000 });
for (const message of messages) {
console.log(Received: ${message.body});
await receiver.completeMessage(message);
}

await receiver.close();

订阅消息(事件驱动)

typescript
const receiver = client.createReceiver("my-queue");

const subscription = receiver.subscribe({
processMessage: async (message) => {
console.log(Processing: ${message.body});
// 成功后消息将自动完成
},
processError: async (args) => {
console.error(Error: ${args.error});
},
});

// 一段时间后停止
setTimeout(async () => {
await subscription.close();
await receiver.close();
}, 60000);

主题与订阅

typescript
// 发送到主题
const topicSender = client.createSender("my-topic");
await topicSender.sendMessages({
  body: { event: "order.created", data: { orderId: "123" } },
  applicationProperties: { eventType: "order.created" },
});

// 从订阅中接收
const subscriptionReceiver = client.createReceiver("my-topic", "my-subscription");
const messages = await subscriptionReceiver.receiveMessages(10);

消息会话 (Message Sessions)

typescript
// 发送会话消息
const sender = client.createSender("session-queue");
await sender.sendMessages({
  body: { step: 1, data: "First step" },
  sessionId: "workflow-123",
});

// 接收会话消息
const sessionReceiver = await client.acceptSession("session-queue", "workflow-123");
const messages = await sessionReceiver.receiveMessages(10);

// 获取/设置会话状态
const state = await sessionReceiver.getSessionState();
await sessionReceiver.setSessionState(Buffer.from(JSON.stringify({ progress: 50 })));

await sessionReceiver.close();

死信处理 (Dead-Letter Handling)

typescript
// 移至死信队列
await receiver.deadLetterMessage(message, {
  deadLetterReason: "Validation failed",
  deadLetterErrorDescription: "Missing required field: orderId",
});

// 处理死信队列
const dlqReceiver = client.createReceiver("my-que


ue", { subQueueType: "deadLetter" });
const dlqMessages = await dlqReceiver.receiveMessages(10);
for (const msg of dlqMessages) {
console.log(DLQ Reason: ${msg.deadLetterReason});
// 重新处理或记录日志
await dlqReceiver.completeMessage(msg);
}
code
## 定时消息 (Scheduled Messages)
typescript
const sender = client.createSender("my-queue");

// 计划未来发送
const scheduledTime = new Date(Date.now() + 60000); // 从现在起 1 分钟后
const sequenceNumber = await sender.scheduleMessages(
{ body: "Delayed message" },
scheduledTime
);

// 取消定时消息
await sender.cancelScheduledMessages(sequenceNumber);

code
## 消息延迟 (Message Deferral)
typescript
// 将消息延迟处理
await receiver.deferMessage(message);

// 通过序列号接收延迟消息
const deferredMessage = await receiver.receiveDeferredMessages(message.sequenceNumber!);
await receiver.completeMessage(deferredMessage[0]);

code
## 窥视消息 (Peek Messages - 非破坏性)
typescript
const receiver = client.createReceiver("my-queue");

// 窥视而不移除
const peekedMessages = await receiver.peekMessages(10);
for (const msg of peekedMessages) {
console.log(Peeked: ${msg.body});
}

code
## 关键类型 (Key Types)
typescript
import {
ServiceBusClient,
ServiceBusSender,
ServiceBusReceiver,
ServiceBusSessionReceiver,
ServiceBusMessage,
ServiceBusReceivedMessage,
ProcessMessageCallback,
ProcessErrorCallback,
} from "@azure/service-bus";
code
## 接收模式 (Receive Modes)
typescript
// Peek-Lock (默认) - 消息被锁定直到完成或放弃
const receiver = client.createReceiver("my-queue", { receiveMode: "peekLock" });
await receiver.completeMessage(message); // 从队列中移除
await receiver.abandonMessage(message); // 返回队列
await receiver.deferMessage(message); // 延迟处理
await receiver.deadLetterMessage(message); // 移至 DLQ

// Receive-and-Delete - 消息立即被移除
const receiver = client.createReceiver("my-queue", { receiveMode: "receiveAndDelete" });
``

最佳实践

1. 使用 Entra ID 认证 - 生产环境应避免使用连接字符串。
2. 复用客户端 - 仅创建一次
ServiceBusClient,并在多个发送者/接收者之间共享。
3. 关闭资源 - 完成后务必关闭发送者/接收者。
4. 处理错误 - 为订阅接收者实现
processError 回调。
5. 使用会话保证顺序 - 当组内消息顺序至关重要时使用 Session。
6. 配置死信队列 - 务必处理 DLQ 消息。
7. 批量发送 - 发送多条消息时使用
createMessageBatch()`。

参考文档

详细模式请参阅:

  • Queues vs Topics Patterns - 队列/主题模式、会话、接收模式、消息结算
  • Error Handling and Reliability - ServiceBusError 错误代码、DLQ 处理、锁续期、优雅停机

适用场景

此技能适用于执行概览中所描述的工作流或操作。

局限性

  • 仅在任务明确符合上述范围时使用此技能。
  • 不要将输出视为环境特定验证、测试或专家评审的替代方案。
  • 如果缺少必要的输入、权限、安全边界或成功标准,请停止并请求澄清。