Azure 存储队列 Python SDK

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

Azure Queue Storage Python SDK

简单且具有成本效益的消息队列,用于异步通信。

安装

bash
pip install azure-storage-queue azure-identity

环境变量

bash
AZURE_STORAGE_ACCOUNT_URL=https://<account>.queue.core.windows.net

身份验证

python
from azure.identity import DefaultAzureCredential
from azure.storage.queue import QueueServiceClient, QueueClient

credential = DefaultAzureCredential()
account_url = "https://<account>.queue.core.windows.net"

服务客户端

service_client = QueueServiceClient(account_url=account_url, credential=credential)

队列客户端

queue_client = QueueClient(account_url=account_url, queue_name="myqueue", credential=credential)

队列操作

python
# 创建队列
service_client.create_queue("myqueue")

获取队列客户端

queue_client = service_client.get_queue_client("myqueue")

删除队列

service_client.delete_queue("myqueue")

列出队列

for queue in service_client.list_queues(): print(queue.name)

发送消息

python
# 发送消息 (字符串)
queue_client.send_message("Hello, Queue!")

使用选项发送

queue_client.send_message( content="Delayed message", visibility_timeout=60, # 隐藏 60 秒 time_to_live=3600 # 1 小时后过期 )

发送 JSON

import json data = {"task": "process", "id": 123} queue_client.send_message(json.dumps(data))

接收消息

python
# 接收消息 (消息将暂时不可见)
messages = queue_client.receive_messages(
    messages_per_page=10,
    visibility_timeout=30  # 处理时间为 30 秒
)

for message in messages:
print(f"ID: {message.id}")
print(f"Content: {message.content}")
print(f"Dequeue count: {message.dequeue_count}")

# 处理消息...

# 处理后删除
queue_client.delete_message(message)

查看消息 (Peek)

python
# 查看消息但不隐藏 (不影响可见性)
messages = queue_client.peek_messages(max_messages=5)

for message in messages:
print(message.content)

更新消息

python
# 延长可见性超时或更新内容
messages = queue_client.receive_messages()
for message in messages:
    # 延长超时时间 (需要更多处理时间)
    queue_client.update_message(
        message,
        visibility_timeout=60
    )
    
    # 更新内容和超时时间
    queue_client.update_message(
        message,
        content="Updated content",
        visibility_timeout=60
    )

删除消息

python
# 处理成功后删除
messages = queue_client.receive_messages()
for message in messages:
    try:
        # 处理...
        queue_client.delete_message(message)
    except Exception:
        # 消息在超时后将重新变为可见
        pass

清空队列

python
# 删除所有消息
queue_client.clear_messages()

队列属性

python
# 获取队列属性
properties = queue_client.get_queue_properties()
print(f"Approximate message count: {properties.approximate_message_count}")

设置/获取元数据

queue_client.set_queue_metadata(metadata={"environment": "production"}) properties = queue_client.get_queue_properties() print(properties.metadata)

异步客户端 (Async Client)

python
from azure.storage.queue.aio import QueueServiceClient, QueueClient
from azure.identity.aio import DefaultAzureCredential

async def queue_operations():
credential = DefaultAzureCredential()

async with QueueClient(
account_url="https://<account>.queue.core.windows.net",
queue_name="myqueue",
credential=credential
) as client:
# 发送
await client.send_message("Async message")

# 接收
async for message in client.receive_messages():
print(message.content)
await client.delete_message(message)

import asyncio
asyncio.run(queue_operations())

Base64 编码

python
from azure.storage.queue import QueueClient, BinaryBase64EncodePolicy, BinaryBase64DecodePolicy

用于二进制数据

queue_client = QueueClient( account_url=account_url, queue_name="myqueue", credential=credential, message_encode_policy=BinaryBase64EncodePolicy(), message_decode_policy=BinaryBase64DecodePolicy() )

发送字节流

queue_client.send_message(b"Binary content")

最佳实践

1. 处理后删除消息,以防止重复处理。
2. 根据处理时间设置合适的可见性超时 (visibility timeout)
3. 处理 dequeue_count 以检测毒丸消息 (poison message)。
4. 在高吞吐量场景下使用异步客户端
5. 使用 peek_messages 进行监控,且不影响队列状态。
6. 设置 time_to_live 以防止消息过期。
7. 如需高级功能(如会话、主题),请考虑使用 Service Bus

适用场景

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

局限性

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