azure-servicebus-py
microsoft/skills
使用 Python SDK 通过 Azure Service Bus 队列、主题和订阅发送和接收消息。
...展开全部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 # 订阅操作必需
AZURE_TOKEN_CREDENTIALS=prod # 仅在生产环境中使用 DefaultAzureCredential 时必需
</namespace>认证与生命周期
🔑 以下每个代码示例均适用两条规则:
- 优先使用
DefaultAzureCredential。 它在本地(Azure CLI / VS Code / Developer CLI)和 Azure(托管标识、工作负载标识)中均能正常工作,且无需更改代码。避免使用连接字符串、账户/API 密钥——它们会绕过 Entra 审计和轮换机制。
- 本地开发:
DefaultAzureCredential可直接使用。- 生产环境:设置
AZURE_TOKEN_CREDENTIALS=prod(或AZURE_TOKEN_CREDENTIALS=<specific_credential></specific_credential>)以将凭据链限制为生产安全的凭据。- 将每个客户端包装在上下文管理器中,以便确定性释放 HTTP 传输层、套接字和令牌缓存:
- 同步:
with <client>(...) as client:</client>- 异步:
async with <client>(...) as client:</client>以及async with DefaultAzureCredential() as credential:(来自azure.identity.aio)代码片段可能简化了此设置,但生产代码应始终遵循这两条规则。
from azure.identity import DefaultAzureCredential, ManagedIdentityCredential
from azure.servicebus import ServiceBusClient
# 本地开发:DefaultAzureCredential。生产环境:设置 AZURE_TOKEN_CREDENTIALS=prod 或 AZURE_TOKEN_CREDENTIALS=<specific_credential>
credential = DefaultAzureCredential(require_envvar=True)
# 或者在生产环境中直接使用特定凭据:
# 请参阅 https://learn.microsoft.com/python/api/overview/azure/identity-readme?view=azure-python#credential-classes
# credential = ManagedIdentityCredential()
namespace = "<namespace>.servicebus.windows.net"
with ServiceBusClient(
fully_qualified_namespace=namespace,
credential=credential
) as client:
# 在此处使用客户端(请参阅以下各节了解操作)
...
</namespace></specific_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())
</namespace>接收消息(异步)
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())
</namespace>接收模式
| 模式 | 行为 | 用例 |
|---|---|---|
| `PEEK_LOCK`(默认) | 消息锁定,必须完成或放弃 | 可靠处理 |
| `RECEIVE_AND_DELETE` | 接收时立即移除 | 最多一次交付 |
from azure.servicebus import ServiceBusReceiveMode
receiver = client.get_queue_receiver(
queue_name="myqueue",
receive_mode=ServiceBusReceiveMode.RECEIVE_AND_DELETE
)
消息结算
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()` | 暂存,通过序列号接收 |
主题和订阅
# 发送到主题
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)
会话(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
)
计划消息
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)
死信队列
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)
</namespace>最佳实践
- 选择同步或异步,并保持一致。 不要在同一个调用路径中混合使用
azure.xxx同步客户端和azure.xxx.aio异步客户端。每个模块选择一种模式。 - 始终对客户端和异步凭据使用上下文管理器。 将每个客户端包装在
with Client(...) as client:(同步)或async with Client(...) as client:(异步)中以实现正确清理。对于来自azure.identity.aio的异步DefaultAzureCredential,也请使用async with credential:以便清理令牌和传输层。 - 使用
DefaultAzureCredential以实现本地开发和 Azure 之间的可移植认证(尽可能避免使用连接字符串/API 密钥)。 - 生产工作负载请使用异步客户端
- 成功处理后完成消息
- 对毒消息使用死信队列
- 对有序、FIFO 处理使用会话
- 对高吞吐量场景使用消息批处理
- 设置
max_wait_time以避免无限阻塞
参考文件
| 文件 | 内容 |
|---|---|
| references/patterns.md | 竞争消费者、会话、重试模式、请求-响应、事务 |
| references/dead-letter.md | DLQ 处理、毒消息、重处理策略 |
| scripts/setup\_servicebus.py | 用于队列/主题/订阅管理和 DLQ 监控的 CLI |
---
name: azure-servicebus-py
description: Send and receive messages using Azure Service Bus queues, topics, and subscriptions with Python SDK.
license: MIT
---
# Azure Service Bus SDK for Python
Enterprise messaging for reliable cloud communication with queues and pub/sub topics.
## Installation
```bash
pip install azure-servicebus azure-identity
```
## Environment Variables
```bash
SERVICEBUS_FULLY_QUALIFIED_NAMESPACE=<namespace>.servicebus.windows.net # Required for all auth methods
SERVICEBUS_QUEUE_NAME=myqueue # Required for queue operations
SERVICEBUS_TOPIC_NAME=mytopic # Required for topic operations
SERVICEBUS_SUBSCRIPTION_NAME=mysubscription # Required for subscription operations
AZURE_TOKEN_CREDENTIALS=prod # Required only if DefaultAzureCredential is used in production
```
## Authentication & Lifecycle
> **🔑 Two rules apply to every code sample below:**
>
> 1. **Prefer `DefaultAzureCredential`.** It works locally (Azure CLI / VS Code / Developer CLI) and in Azure (managed identity, workload identity) with no code change. Avoid connection strings, account/API keys — they bypass Entra audit and rotation.
> - Local dev: `DefaultAzureCredential` works as-is.
> - Production: set `AZURE_TOKEN_CREDENTIALS=prod` (or `AZURE_TOKEN_CREDENTIALS=<specific_credential>`) to constrain the credential chain to production-safe credentials.
> 2. **Wrap every client in a context manager** so HTTP transports, sockets, and token caches are released deterministically:
> - Sync: `with <Client>(...) as client:`
> - Async: `async with <Client>(...) as client:` **and** `async with DefaultAzureCredential() as credential:` (from `azure.identity.aio`)
>
> Snippets may abbreviate this setup, but production code should always follow both rules.
```python
from azure.identity import DefaultAzureCredential, ManagedIdentityCredential
from azure.servicebus import ServiceBusClient
# Local dev: DefaultAzureCredential. Production: set AZURE_TOKEN_CREDENTIALS=prod or AZURE_TOKEN_CREDENTIALS=<specific_credential>
credential = DefaultAzureCredential(require_envvar=True)
# Or use a specific credential directly in production:
# See https://learn.microsoft.com/python/api/overview/azure/identity-readme?view=azure-python#credential-classes
# credential = ManagedIdentityCredential()
namespace = "<namespace>.servicebus.windows.net"
with ServiceBusClient(
fully_qualified_namespace=namespace,
credential=credential
) as client:
# Use client here (see following sections for operations)
...
```
## Client Types
| Client | Purpose | Get From |
|--------|---------|----------|
| `ServiceBusClient` | Connection management | Direct instantiation |
| `ServiceBusSender` | Send messages | `client.get_queue_sender()` / `get_topic_sender()` |
| `ServiceBusReceiver` | Receive messages | `client.get_queue_receiver()` / `get_subscription_receiver()` |
## Send Messages (Async)
```python
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:
# Single message
message = ServiceBusMessage("Hello, Service Bus!")
await sender.send_messages(message)
# Batch of messages
messages = [ServiceBusMessage(f"Message {i}") for i in range(10)]
await sender.send_messages(messages)
# Message batch (for size control)
batch = await sender.create_message_batch()
for i in range(100):
try:
batch.add_message(ServiceBusMessage(f"Batch message {i}"))
except ValueError: # Batch full
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())
```
## Receive Messages (Async)
```python
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:
# Receive batch
messages = await receiver.receive_messages(
max_message_count=10,
max_wait_time=5 # seconds
)
for msg in messages:
print(f"Received: {str(msg)}")
await receiver.complete_message(msg) # Remove from queue
asyncio.run(receive_messages())
```
## Receive Modes
| Mode | Behavior | Use Case |
|------|----------|----------|
| `PEEK_LOCK` (default) | Message locked, must complete/abandon | Reliable processing |
| `RECEIVE_AND_DELETE` | Removed immediately on receive | At-most-once delivery |
```python
from azure.servicebus import ServiceBusReceiveMode
receiver = client.get_queue_receiver(
queue_name="myqueue",
receive_mode=ServiceBusReceiveMode.RECEIVE_AND_DELETE
)
```
## Message Settlement
```python
async with receiver:
messages = await receiver.receive_messages(max_message_count=1)
for msg in messages:
try:
# Process message...
await receiver.complete_message(msg) # Success - remove from queue
except ProcessingError:
await receiver.abandon_message(msg) # Retry later
except PermanentError:
await receiver.dead_letter_message(
msg,
reason="ProcessingFailed",
error_description="Could not process"
)
```
| Action | Effect |
|--------|--------|
| `complete_message()` | Remove from queue (success) |
| `abandon_message()` | Release lock, retry immediately |
| `dead_letter_message()` | Move to dead-letter queue |
| `defer_message()` | Set aside, receive by sequence number |
## Topics and Subscriptions
```python
# Send to topic
sender = client.get_topic_sender(topic_name="mytopic")
async with sender:
await sender.send_messages(ServiceBusMessage("Topic message"))
# Receive from subscription
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)
```python
# Send with session
message = ServiceBusMessage("Session message")
message.session_id = "order-123"
await sender.send_messages(message)
# Receive from specific session
receiver = client.get_queue_receiver(
queue_name="session-queue",
session_id="order-123"
)
# Receive from next available session
from azure.servicebus import NEXT_AVAILABLE_SESSION
receiver = client.get_queue_receiver(
queue_name="session-queue",
session_id=NEXT_AVAILABLE_SESSION
)
```
## Scheduled Messages
```python
from datetime import datetime, timedelta, timezone
message = ServiceBusMessage("Scheduled message")
scheduled_time = datetime.now(timezone.utc) + timedelta(minutes=10)
# Schedule message
sequence_number = await sender.schedule_messages(message, scheduled_time)
# Cancel scheduled message
await sender.cancel_scheduled_messages(sequence_number)
```
## Dead-Letter Queue
```python
from azure.servicebus import ServiceBusSubQueue
# Receive from dead-letter queue
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)
```
## Sync Client (for simple scripts)
```python
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)
```
## Best Practices
1. **Pick sync OR async and stay consistent.** Do not mix `azure.xxx` sync clients with `azure.xxx.aio` async clients in the same call path. Choose one mode per module.
2. **Always use context managers for clients and async credentials.** Wrap every client in `with Client(...) as client:` (sync) or `async with Client(...) as client:` (async) for proper cleanup. For async `DefaultAzureCredential` from `azure.identity.aio`, also use `async with credential:` so tokens and transports are cleaned up.
3. **Use `DefaultAzureCredential`** for portable auth across local dev and Azure (avoid connection strings / API keys when possible).
4. **Use async client** for production workloads
5. **Complete messages** after successful processing
6. **Use dead-letter queue** for poison messages
7. **Use sessions** for ordered, FIFO processing
8. **Use message batches** for high-throughput scenarios
9. **Set `max_wait_time`** to avoid infinite blocking
## Reference Files
| File | Contents |
|------|----------|
| [references/patterns.md](references/patterns.md) | Competing consumers, sessions, retry patterns, request-response, transactions |
| [references/dead-letter.md](references/dead-letter.md) | DLQ handling, poison messages, reprocessing strategies |
| [scripts/setup_servicebus.py](scripts/setup_servicebus.py) | CLI for queue/topic/subscription management and DLQ monitoring |
所有文件
0 个文件安装 azure-servicebus-py
下载并将技能文件解压至你的 .claude/skills/ 目录。
下载ZIP克隆仓库并复制技能文件到您的项目中。
git clone https://github.com/microsoft/skills/tree/main/.github/plugins/azure-sdk-python/skills/azure-servicebus-py # Copy SKILL.md to your .claude/skills/ directory
复制





首页
