Azure Service Bus TypeScript
Azure Service Bus TypeScript SDK
支持队列、主题和订阅的企业级消息传递。
安装
npm install @azure/service-bus @azure/identity环境变量
SERVICEBUS_NAMESPACE=<namespace>.servicebus.windows.net
SERVICEBUS_QUEUE_NAME=my-queue
SERVICEBUS_TOPIC_NAME=my-topic
SERVICEBUS_SUBSCRIPTION_NAME=my-subscription身份验证
import { ServiceBusClient } from "@azure/service-bus";
import { DefaultAzureCredential } from "@azure/identity";
const fullyQualifiedNamespace = process.env.SERVICEBUS_NAMESPACE!;
const client = new ServiceBusClient(fullyQualifiedNamespace, new DefaultAzureCredential());
核心工作流
向队列发送消息
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();
从队列接收消息
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();
订阅消息(事件驱动)
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);
主题与订阅
// 发送到主题
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)
// 发送会话消息
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)
// 移至死信队列
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);
}
## 定时消息 (Scheduled Messages)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);
## 消息延迟 (Message Deferral)// 将消息延迟处理
await receiver.deferMessage(message);
// 通过序列号接收延迟消息
const deferredMessage = await receiver.receiveDeferredMessages(message.sequenceNumber!);
await receiver.completeMessage(deferredMessage[0]);
## 窥视消息 (Peek Messages - 非破坏性)const receiver = client.createReceiver("my-queue");
// 窥视而不移除
const peekedMessages = await receiver.peekMessages(10);
for (const msg of peekedMessages) {
console.log(Peeked: ${msg.body});
}
## 关键类型 (Key Types)import {
ServiceBusClient,
ServiceBusSender,
ServiceBusReceiver,
ServiceBusSessionReceiver,
ServiceBusMessage,
ServiceBusReceivedMessage,
ProcessMessageCallback,
ProcessErrorCallback,
} from "@azure/service-bus";
## 接收模式 (Receive Modes)// 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,并在多个发送者/接收者之间共享。processError
3. 关闭资源 - 完成后务必关闭发送者/接收者。
4. 处理错误 - 为订阅接收者实现 回调。createMessageBatch()`。
5. 使用会话保证顺序 - 当组内消息顺序至关重要时使用 Session。
6. 配置死信队列 - 务必处理 DLQ 消息。
7. 批量发送 - 发送多条消息时使用
参考文档
详细模式请参阅:
- Queues vs Topics Patterns - 队列/主题模式、会话、接收模式、消息结算
- Error Handling and Reliability - ServiceBusError 错误代码、DLQ 处理、锁续期、优雅停机
适用场景
此技能适用于执行概览中所描述的工作流或操作。局限性
- 仅在任务明确符合上述范围时使用此技能。
- 不要将输出视为环境特定验证、测试或专家评审的替代方案。
- 如果缺少必要的输入、权限、安全边界或成功标准,请停止并请求澄清。