Azure 存储队列 TypeScript (TS)

azure-storage-queue-ts
分类编程
作者Agentic Awesome Skills 社区
许可MIT
评分4.50/5
使用14.0K

@azure/storage-queue (TypeScript/JavaScript)

用于 Azure Queue Storage 操作的 SDK —— 支持发送、接收、查看和管理队列中的消息。

安装

bash
npm install @azure/storage-queue @azure/identity

当前版本: 12.x
Node.js: >= 18.0.0

环境变量

bash
AZURE_STORAGE_ACCOUNT_NAME=<account-name>
AZURE_STORAGE_ACCOUNT_KEY=<account-key>

或连接字符串

AZURE_STORAGE_CONNECTION_STRING=DefaultEndpointsProtocol=https;AccountName=...

身份验证

DefaultAzureCredential (推荐)

typescript
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()
);

连接字符串

typescript
import { QueueServiceClient } from "@azure/storage-queue";

const client = QueueServiceClient.fromConnectionString(
process.env.AZURE_STORAGE_CONNECTION_STRING!
);

StorageSharedKeyCredential (仅限 Node.js)

typescript
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 令牌

typescript
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}
);

客户端层级

code
QueueServiceClient (账户级)
└── QueueClient (队列级)
    └── Messages (发送, 接收, 查看, 删除)

队列操作

创建队列

typescript
const queueClient = client.getQueueClient("my-queue");
await queueClient.create();

// 或如果不存在则创建
await queueClient.createIfNotExists();

列出队列

typescript
for await (const queue of client.listQueues()) {
  console.log(queue.name);
}

// 使用前缀过滤
for await (const queue of client.listQueues({ prefix: "task-" })) {
console.log(queue.name);
}

删除队列

typescript
await queueClient.delete();

// 或如果存在则删除
await queueClient.deleteIfExists();

获取队列属性

typescript
const properties = await queueClient.getProperties();
console.log("Approximate message count:", properties.approximateMessagesCount);
console.log("Metadata:", properties.metadata);

设置队列元数据

typescript
await queueClient.setMetadata({
  department: "engineering",
  priority: "high",
});

消息操作

发送消息

typescript
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));

code
### 接收消息
typescript
// 最多接收 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);
}

code
### 窥视消息 (Peek)

在不将其从队列中移除的情况下查看消息(无可见性超时)。

typescript
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 - 无法删除窥视的消息
}

code
### 更新消息

延长可见性超时或更新内容。

typescript
// 接收一条消息
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);
}

code
### 删除消息
typescript
// 接收后删除
const response = await queueClient.receiveMessages();
const message = response.receivedMessageItems[0];

if (message) {
await queueClient.deleteMessage(message.messageId, message.popReceipt);
}

code
### 清空所有消息
typescript
await queueClient.clearMessages();
code
## 消息处理模式

基础 Worker 模式

typescript async function processQueue(queueClient: QueueClient): Promise<void> { while (true) { const response = await queueClient.receiveMessages({ numberOfMessages: 10, visibilityTimeout: 30, });

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));
}

code
### 毒药消息 (Poison Message) 处理
typescript
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

code
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);
}
}
}

带有可见性延长的批量处理

typescript
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 编码。你可以对其进行自定义:

typescript
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

typescript
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

typescript
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();
code
## 错误处理
typescript
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;
}

code
## TypeScript 类型参考
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 访问 | ✅ | ✅ |
| 所有消息操作 | ✅ | ✅ |

适用场景

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

局限性

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