Azure Service Bus Python SDK
Azure Service Bus Python SDK
通过队列和发布/订阅主题实现可靠云通信的企业级消息传递。
安装
pip install azure-servicebus azure-identity环境变量
SERVICEBUS_FULLY_QUALIFIED_NAMESPACE=<namespace>.servicebus.windows.net
SERVICEBUS_QUEUE_NAME=myqueue
SERVICEBUS_TOPIC_NAME=mytopic
SERVICEBUS_SUBSCRIPTION_NAME=mysubscription身份验证
from azure.identity import DefaultAzureCredential
from azure.servicebus import ServiceBusClient
credential = DefaultAzureCredential()
namespace = "<namespace>.servicebus.windows.net"
client = ServiceBusClient(
fully_qualified_namespace=namespace,
credential=credential
)
客户端类型
| 客户端 | 用途 | 获取方式 |
|--------|---------|----------|
| ServiceBusClient | 连接管理 | 直接实例化 |
| ServiceBusSender | 发送消息 | client.get_queue_sender() / get_topic_sender() |
| ServiceBusReceiver | 接收消息 | client.get_queue_receiver() / get_subscription_receiver() |
发送消息 (异步)
import asyncio
from azure.servicebus.aio import ServiceBusClient
from azure.servicebus import ServiceBusMessage
from azure.identity.aio import DefaultAzureCredential
async def send_messages():
credential = DefaultAzureCredential()
async with ServiceBusClient(
fully_qualified_namespace="<namespace>.servicebus.windows.net",
credential=credential
) as client:
sender = client.get_queue_sender(queue_name="myqueue")
async with sender:
# 单条消息
message = ServiceBusMessage("Hello, Service Bus!")
await sender.send_messages(message)
# 消息列表
messages = [ServiceBusMessage(f"Message {i}") for i in range(10)]
await sender.send_messages(messages)
# 消息批次 (用于控制大小)
batch = await sender.create_message_batch()
for i in range(100):
try:
batch.add_message(ServiceBusMessage(f"Batch message {i}"))
except ValueError: # 批次已满
await sender.send_messages(batch)
batch = await sender.create_message_batch()
batch.add_message(ServiceBusMessage(f"Batch message {i}"))
await sender.send_messages(batch)
asyncio.run(send_messages())
接收消息 (异步)
async def receive_messages():
credential = DefaultAzureCredential()
async with ServiceBusClient(
fully_qualified_namespace="<namespace>.servicebus.windows.net",
credential=credential
) as client:
receiver = client.get_queue_receiver(queue_name="myqueue")
async with receiver:
# 批量接收
messages = await receiver.receive_messages(
max_message_count=10,
max_wait_time=5 # 秒
)
for msg in messages:
print(f"Received: {str(msg)}")
await receiver.complete_message(msg) # 从队列中移除
asyncio.run(receive_messages())
接收模式
| 模式
| 行为 | 使用场景 |
|------|----------|
| PEEK_LOCK (默认) | 消息被锁定,必须完成 (complete) 或放弃 (abandon) | 可靠处理 |
| RECEIVE_AND_DELETE | 接收后立即删除 | 最多一次交付 (At-most-once) |
from azure.servicebus import ServiceBusReceiveMode
receiver = client.get_queue_receiver(
queue_name="myqueue",
receive_mode=ServiceBusReceiveMode.RECEIVE_AND_DELETE
)
消息结算 (Message Settlement)
async with receiver:
messages = await receiver.receive_messages(max_message_count=1)
for msg in messages:
try:
# 处理消息...
await receiver.complete_message(msg) # 成功 - 从队列中删除
except ProcessingError:
await receiver.abandon_message(msg) # 稍后重试
except PermanentError:
await receiver.dead_letter_message(
msg,
reason="ProcessingFailed",
error_description="Could not process"
)| 操作 | 效果 |
|--------|--------|
| complete_message() | 从队列中删除 (成功) |
| abandon_message() | 释放锁,立即重试 |
| dead_letter_message() | 移至死信队列 |
| defer_message() | 暂存,通过序列号接收 |
主题与订阅 (Topics and Subscriptions)
# 发送到主题
sender = client.get_topic_sender(topic_name="mytopic")
async with sender:
await sender.send_messages(ServiceBusMessage("Topic message"))
从订阅中接收
receiver = client.get_subscription_receiver(
topic_name="mytopic",
subscription_name="mysubscription"
)
async with receiver:
messages = await receiver.receive_messages(max_message_count=10)会话 (Sessions - FIFO)
# 发送带会话的消息
message = ServiceBusMessage("Session message")
message.session_id = "order-123"
await sender.send_messages(message)
从特定会话接收
receiver = client.get_queue_receiver(
queue_name="session-queue",
session_id="order-123"
)
从下一个可用会话接收
from azure.servicebus import NEXT_AVAILABLE_SESSION
receiver = client.get_queue_receiver(
queue_name="session-queue",
session_id=NEXT_AVAILABLE_SESSION
)计划消息 (Scheduled Messages)
from datetime import datetime, timedelta, timezone
message = ServiceBusMessage("Scheduled message")
scheduled_time = datetime.now(timezone.utc) + timedelta(minutes=10)
计划发送消息
sequence_number = await sender.schedule_messages(message, scheduled_time)
取消计划消息
await sender.cancel_scheduled_messages(sequence_number)死信队列 (Dead-Letter Queue)
from azure.servicebus import ServiceBusSubQueue
从死信队列接收
dlq_receiver = client.get_queue_receiver(
queue_name="myqueue",
sub_queue=ServiceBusSubQueue.DEAD_LETTER
)
async with dlq_receiver:
messages = await dlq_receiver.receive_messages(max_message_count=10)
for msg in messages:
print(f"Dead-lettered: {msg.dead_letter_reason}")
await dlq_receiver.complete_message(msg)
同步客户端 (适用于简单脚本)
from azure.servicebus import ServiceBusClient, ServiceBusMessage
from azure.identity import DefaultAzureCredential
with ServiceBusClient(
fully_qualified_namespace="<namespace>.servicebus.windows.net",
credential=DefaultAzureCredential()
) as client:
with client.get_queue_sender("myqueue") as sender:
sender.send_messages(ServiceBusMessage("Sync message"))
with client.get_queue_receiver("myqueue") as receiver:
for msg in receiver:
print(str(msg))
receiver.complete_message(msg)最佳实践
1. 生产环境建议使用异步客户端 (async client)
2. 使用上下文管理器 (async with) 以确保正确清理资源
3. 在成功处理后完成消息 (Complete messages)
4. 针对毒丸消息使用死信队列 (Dead-letter queue)
5. 使用会话 (Sessions) 实现有序的 FIFO 处理
6. 在高吞吐量场景下使用消息批处理
7. 设置 max_wait_time 以避免无限阻塞
参考文件
| 文件 | 内容 |
|------|----------|
| references/patterns.md | 竞争消费者、会话、重试模式、请求-响应、事务 |
| references/dead-letter.md | DLQ 处理、毒丸消息、重新处理策略 |
| scripts/setup_servicebus.py | 用于队列/主题/订阅管理及 DLQ 监控的 CLI 工具 |
适用场景
本技能适用于执行概览中所描述的工作流或操作。局限性
- 仅在任务明确符合上述范围时使用此技能。
- 不要将输出结果视为针对特定环境的验证、测试或专家评审的替代方案。
- 如果缺少必要的输入、权限、安全边界或成功标准,请停止操作并请求澄清。