Azure 存储队列 TypeScript (TS)
@azure/storage-queue (TypeScript/JavaScript)
用于 Azure Queue Storage 操作的 SDK —— 支持发送、接收、查看和管理队列中的消息。
安装
npm install @azure/storage-queue @azure/identity当前版本: 12.x
Node.js: >= 18.0.0
环境变量
AZURE_STORAGE_ACCOUNT_NAME=<account-name>
AZURE_STORAGE_ACCOUNT_KEY=<account-key>
或连接字符串
AZURE_STORAGE_CONNECTION_STRING=DefaultEndpointsProtocol=https;AccountName=...身份验证
DefaultAzureCredential (推荐)
import { QueueServiceClient } from "@azure/storage-queue";
import { DefaultAzureCredential } from "@azure/identity";
const accountName = process.env.AZURE_STORAGE_ACCOUNT_NAME!;
const client = new QueueServiceClient(
https://${accountName}.queue.core.windows.net,
new DefaultAzureCredential()
);
连接字符串
import { QueueServiceClient } from "@azure/storage-queue";
const client = QueueServiceClient.fromConnectionString(
process.env.AZURE_STORAGE_CONNECTION_STRING!
);
StorageSharedKeyCredential (仅限 Node.js)
import { QueueServiceClient, StorageSharedKeyCredential } from "@azure/storage-queue";
const accountName = process.env.AZURE_STORAGE_ACCOUNT_NAME!;
const accountKey = process.env.AZURE_STORAGE_ACCOUNT_KEY!;
const sharedKeyCredential = new StorageSharedKeyCredential(accountName, accountKey);
const client = new QueueServiceClient(
https://${accountName}.queue.core.windows.net,
sharedKeyCredential
);
SAS 令牌
import { QueueServiceClient } from "@azure/storage-queue";
const accountName = process.env.AZURE_STORAGE_ACCOUNT_NAME!;
const sasToken = process.env.AZURE_STORAGE_SAS_TOKEN!;
const client = new QueueServiceClient(
https://${accountName}.queue.core.windows.net${sasToken}
);
客户端层级
QueueServiceClient (账户级)
└── QueueClient (队列级)
└── Messages (发送, 接收, 查看, 删除)队列操作
创建队列
const queueClient = client.getQueueClient("my-queue");
await queueClient.create();
// 或如果不存在则创建
await queueClient.createIfNotExists();
列出队列
for await (const queue of client.listQueues()) {
console.log(queue.name);
}
// 使用前缀过滤
for await (const queue of client.listQueues({ prefix: "task-" })) {
console.log(queue.name);
}
删除队列
await queueClient.delete();
// 或如果存在则删除
await queueClient.deleteIfExists();
获取队列属性
const properties = await queueClient.getProperties();
console.log("Approximate message count:", properties.approximateMessagesCount);
console.log("Metadata:", properties.metadata);设置队列元数据
await queueClient.setMetadata({
department: "engineering",
priority: "high",
});消息操作
发送消息
const queueClient = client.getQueueClient("my-queue");
// 普通消息
await queueClient.sendMessage("Hello, World!");
// 带选项的消息
await queueClient.sendMessage("Delayed message", {
visibilityTimeout: 60, // 隐藏时长(秒)
60 秒
messageTimeToLive: 3600, // 1 小时后过期
});
// JSON 消息(必须为字符串)
const task = { type: "process", data: { id: 123 } };
await queueClient.sendMessage(JSON.stringify(task));
### 接收消息// 最多接收 32 条消息(默认:1 条)
const response = await queueClient.receiveMessages({
numberOfMessages: 10,
visibilityTimeout: 30, // 处理时间为 30 秒
});
for (const message of response.receivedMessageItems) {
console.log("Message ID:", message.messageId);
console.log("Content:", message.messageText);
console.log("Dequeue Count:", message.dequeueCount);
console.log("Pop Receipt:", message.popReceipt);
// 处理消息...
// 处理后删除
await queueClient.deleteMessage(message.messageId, message.popReceipt);
}
### 窥视消息 (Peek)
在不将其从队列中移除的情况下查看消息(无可见性超时)。
const response = await queueClient.peekMessages({
numberOfMessages: 5,
});
for (const message of response.peekedMessageItems) {
console.log("Message ID:", message.messageId);
console.log("Content:", message.messageText);
// 注意:没有 popReceipt - 无法删除窥视的消息
}
### 更新消息
延长可见性超时或更新内容。
// 接收一条消息
const response = await queueClient.receiveMessages();
const message = response.receivedMessageItems[0];
if (message) {
// 更新内容并延长可见性
const updateResponse = await queueClient.updateMessage(
message.messageId,
message.popReceipt,
"Updated content",
60 // 新的可见性超时时间(秒)
);
// 后续操作请使用新的 popReceipt
console.log("New pop receipt:", updateResponse.popReceipt);
}
### 删除消息// 接收后删除
const response = await queueClient.receiveMessages();
const message = response.receivedMessageItems[0];
if (message) {
await queueClient.deleteMessage(message.messageId, message.popReceipt);
}
### 清空所有消息await queueClient.clearMessages();
## 消息处理模式
基础 Worker 模式
if (response.receivedMessageItems.length === 0) {
// 没有消息,等待一段时间后再次轮询
await sleep(5000);
continue;
}
for (const message of response.receivedMessageItems) {
try {
await processMessage(message.messageText);
await queueClient.deleteMessage(message.messageId, message.popReceipt);
} catch (error) {
console.error(Failed to process message ${message.messageId}:, error);
// 消息将在超时后重新变为可见
}
}
}
}
async function processMessage(content: string): Promise<void> {
const task = JSON.parse(content);
// 处理任务...
}
function sleep(ms: number): Promise<void> {
return new Promise((resolve) => setTimeout(resolve, ms));
}
### 毒药消息 (Poison Message) 处理const MAX_DEQUEUE_COUNT = 5;
async function processWithPoisonHandling(
queueClient: QueueClient,
poisonQueueClient: QueueClient
): Promise<void> {
const response = await queueClient.receiveMessages({
numberOfMessages: 10,
visibilityTimeout: 30,
});
for (const message of response.rec
eivedMessageItems) {
if (message.dequeueCount > MAX_DEQUEUE_COUNT) {
// 移至死信队列 (poison queue)
await poisonQueueClient.sendMessage(message.messageText);
await queueClient.deleteMessage(message.messageId, message.popReceipt);
console.log(Moved message ${message.messageId} to poison queue);
continue;
}
try {
await processMessage(message.messageText);
await queueClient.deleteMessage(message.messageId, message.popReceipt);
} catch (error) {
console.error(Processing failed (attempt ${message.dequeueCount}):, error);
}
}
}
带有可见性延长的批量处理
async function processBatchWithExtension(queueClient: QueueClient): Promise<void> {
const response = await queueClient.receiveMessages({
numberOfMessages: 1,
visibilityTimeout: 60,
});
const message = response.receivedMessageItems[0];
if (!message) return;
let popReceipt = message.popReceipt;
// 启动可见性延长定时器
const extensionInterval = setInterval(async () => {
try {
const updateResponse = await queueClient.updateMessage(
message.messageId,
popReceipt,
message.messageText,
60 // 再延长 60 秒
);
popReceipt = updateResponse.popReceipt;
} catch (error) {
console.error("Failed to extend visibility:", error);
}
}, 45000); // 每 45 秒延长一次
try {
await longRunningProcess(message.messageText);
await queueClient.deleteMessage(message.messageId, popReceipt);
} finally {
clearInterval(extensionInterval);
}
}
消息编码
默认情况下,消息采用 Base64 编码。你可以对其进行自定义:
import { QueueClient } from "@azure/storage-queue";
// 用于纯文本的自定义编码器/解码器
const queueClient = new QueueClient(
https://${accountName}.queue.core.windows.net/my-queue,
credential,
{
messageEncoding: "text", // "base64" (默认) 或 "text"
}
);
// 或使用自定义编码函数
const customQueueClient = new QueueClient(
https://${accountName}.queue.core.windows.net/my-queue,
credential,
{
messageEncoding: {
encode: (message: string) => Buffer.from(message).toString("base64"),
decode: (message: string) => Buffer.from(message, "base64").toString(),
},
}
);
SAS 令牌生成 (仅限 Node.js)
生成队列 SAS
import {
QueueSASPermissions,
generateQueueSASQueryParameters,
StorageSharedKeyCredential,
} from "@azure/storage-queue";
const sharedKeyCredential = new StorageSharedKeyCredential(accountName, accountKey);
const sasToken = generateQueueSASQueryParameters(
{
queueName: "my-queue",
permissions: QueueSASPermissions.parse("raup"), // read, add, update, process
startsOn: new Date(),
expiresOn: new Date(Date.now() + 3600 * 1000), // 1 小时
},
sharedKeyCredential
).toString();
const sasUrl = https://${accountName}.queue.core.windows.net/my-queue?${sasToken};
生成账户 SAS
import {
AccountSASPermissions,
AccountSASResourceTypes,
AccountSASServices,
generateAccountSASQueryParameters,
} from "@azure/storage-queue";
const sasToken = generateAccountSASQueryParameters(
{
services: AccountSASServices.parse("q").toString(), // queue
resourceTypes: AccountSASResourceTypes.parse("sco").toString(),
permissions: AccountSASPermissions.parse("rwdlacupi"),
expiresOn: new Date(Date
.now() + 24 * 3600 * 1000),
},
sharedKeyCredential
).toString();
## 错误处理import { RestError } from "@azure/storage-queue";
try {
await queueClient.sendMessage("test");
} catch (error) {
if (error instanceof RestError) {
switch (error.statusCode) {
case 404:
console.log("队列未找到");
break;
case 400:
console.log("请求错误 - 消息过大或无效");
break;
case 403:
console.log("访问被拒绝");
break;
case 409:
console.log("队列已存在或正在删除");
break;
default:
console.error(存储错误 ${error.statusCode}: ${error.message});
}
}
throw error;
}
## TypeScript 类型参考import {
// 客户端
QueueServiceClient,
QueueClient,
// 身份验证
StorageSharedKeyCredential,
AnonymousCredential,
// SAS
QueueSASPermissions,
AccountSASPermissions,
AccountSASServices,
AccountSASResourceTypes,
generateQueueSASQueryParameters,
generateAccountSASQueryParameters,
// 消息
DequeuedMessageItem,
PeekedMessageItem,
QueueSendMessageResponse,
QueueReceiveMessageResponse,
QueueUpdateMessageResponse,
// 队列
QueueItem,
QueueGetPropertiesResponse,
// 错误
RestError,
} from "@azure/storage-queue";
``
消息限制
| 限制项 | 数值 |
|-------|-------|
| 最大消息大小 | 64 KB |
| 最大可见性超时 | 7 天 |
| 最大生存时间 (TTL) | 7 天 (或 -1 表示永久) |
| 每次接收最大消息数 | 32 |
| 默认可见性超时 | 30 秒 |
最佳实践
1. 使用 DefaultAzureCredential — 优先使用 AAD 而非连接字符串/密钥
2. 处理后立即删除 — 防止重复处理
3. 处理毒药消息 (Poison Messages) — 将失败的消息移至死信队列
4. 设置合理的可见性超时 — 根据预期处理时间进行设置
5. 为长任务延长可见性 — 更新消息以防止超时
6. 使用 JSON 传输结构化数据 — 将对象序列化为 JSON 字符串
7. 检查 dequeueCount — 检测重复失败的消息
8. 使用批量接收 — 提高接收效率
平台差异
| 功能 | Node.js | 浏览器 |
|---------|---------|---------|
| StorageSharedKeyCredential` | ✅ | ❌ |
| SAS 生成 | ✅ | ❌ |
| DefaultAzureCredential | ✅ | ❌ |
| 匿名/SAS 访问 | ✅ | ✅ |
| 所有消息操作 | ✅ | ✅ |
适用场景
本技能适用于执行概览中所描述的工作流或操作。局限性
- 仅在任务明确符合上述范围时使用此技能。
- 不要将输出结果视为环境特定验证、测试或专家评审的替代方案。
- 如果缺少必要的输入、权限、安全边界或成功标准,请停止并请求澄清。