azure-servicebus-py
Azure Service Bus SDK for Python messaging. Use for queues, topics, subscriptions, and enterprise messaging patterns. Triggers: "service bus", "ServiceBusClient", "queue", "topic", "subscription", "message broker".
Install
npx skills add https://github.com/microsoft/skills/tree/main/.github/plugins/azure-sdk-python/skills/azure-servicebus-py
claude plugin marketplace add https://llmmart.ai/marketplace.json && claude plugin install microsoft-skills@llmmart
git clone https://github.com/microsoft/skills.git
The skills CLI installs just this skill, for any of its supported agents. Claude Code installs the whole microsoft/skills collection as a plugin from our marketplace. Git is the plain clone.
Skill manifest
Azure Service Bus SDK for Python
Enterprise messaging for reliable cloud communication with queues and pub/sub topics.
Installation
pip install azure-servicebus azure-identity
Environment Variables
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:
- 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:
DefaultAzureCredentialworks as-is.- Production: set
AZURE_TOKEN_CREDENTIALS=prod(orAZURE_TOKEN_CREDENTIALS=<specific_credential>) to constrain the credential chain to production-safe credentials.- 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:andasync with DefaultAzureCredential() as credential:(fromazure.identity.aio)Snippets may abbreviate this setup, but production code should always follow both rules.
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)
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)
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 |
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:
# 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
# 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)
# 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
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
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)
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
- Pick sync OR async and stay consistent. Do not mix
azure.xxxsync clients withazure.xxx.aioasync clients in the same call path. Choose one mode per module. - Always use context managers for clients and async credentials. Wrap every client in
with Client(...) as client:(sync) orasync with Client(...) as client:(async) for proper cleanup. For asyncDefaultAzureCredentialfromazure.identity.aio, also useasync with credential:so tokens and transports are cleaned up. - Use
DefaultAzureCredentialfor portable auth across local dev and Azure (avoid connection strings / API keys when possible). - Use async client for production workloads
- Complete messages after successful processing
- Use dead-letter queue for poison messages
- Use sessions for ordered, FIFO processing
- Use message batches for high-throughput scenarios
- Set
max_wait_timeto avoid infinite blocking
Reference Files
| File | Contents |
|---|---|
| references/patterns.md | Competing consumers, sessions, retry patterns, request-response, transactions |
| references/dead-letter.md | DLQ handling, poison messages, reprocessing strategies |
| scripts/setup_servicebus.py | CLI for queue/topic/subscription management and DLQ monitoring |
Files (skills)
-
references
-
dead-letter.md 13.8 KB
# Dead-Letter Queue Reference Handling poison messages and dead-letter queue processing in Azure Service Bus. ## Dead-Letter Queue Overview The dead-letter queue (DLQ) is a secondary sub-queue for messages that cannot be processed: ``` Main Queue: myqueue └── Dead-Letter Queue: myqueue/$deadletterqueue ``` ## Why Messages Get Dead-Lettered | Reason | Description | |--------|-------------| | `MaxDeliveryCountExceeded` | Message abandoned too many times | | `HeaderSizeExceeded` | Message headers too large | | `TTLExpiration` | Message expired before delivery | | `SessionIdMismatch` | Session ID doesn't match | | `MessageSizeExceeded` | Message body too large | | Custom reason | Explicitly dead-lettered by application | ## Receiving from Dead-Letter Queue ```python from azure.servicebus import ServiceBusSubQueue from azure.servicebus.aio import ServiceBusClient from azure.identity.aio import DefaultAzureCredential async def receive_dead_letters(namespace: str, queue_name: str): """Receive messages from dead-letter queue.""" credential = DefaultAzureCredential() async with ServiceBusClient( fully_qualified_namespace=namespace, credential=credential ) as client: # Get dead-letter queue receiver dlq_receiver = client.get_queue_receiver( queue_name=queue_name, sub_queue=ServiceBusSubQueue.DEAD_LETTER ) async with dlq_receiver: messages = await dlq_receiver.receive_messages( max_message_count=10, max_wait_time=5 ) for msg in messages: print(f"Dead-letter message: {str(msg)}") print(f" Reason: {msg.dead_letter_reason}") print(f" Description: {msg.dead_letter_error_description}") print(f" Enqueued: {msg.enqueued_time_utc}") print(f" Delivery count: {msg.delivery_count}") # Process or complete await dlq_receiver.complete_message(msg) ``` ## Explicit Dead-Lettering Dead-letter messages programmatically for non-retryable errors: ```python async def process_with_dead_letter(receiver, msg): """Process message and dead-letter on permanent failures.""" try: result = await process_message(msg) await receiver.complete_message(msg) return result except ValidationError as e: # Invalid message format - don't retry await receiver.dead_letter_message( msg, reason="ValidationFailed", error_description=f"Invalid format: {e}" ) except DuplicateError as e: # Already processed - dead-letter with context await receiver.dead_letter_message( msg, reason="DuplicateDetected", error_description=f"Message already processed: {e}" ) except ExternalServiceUnavailable: # Temporary - abandon for retry await receiver.abandon_message(msg) except Exception as e: # Unknown error - dead-letter with full context await receiver.dead_letter_message( msg, reason="UnhandledException", error_description=f"{type(e).__name__}: {str(e)}" ) ``` ## Dead-Letter Queue Processor Automated processing of dead-lettered messages: ```python import json from datetime import datetime, timezone class DeadLetterProcessor: """Process and analyze dead-lettered messages.""" def __init__(self, client, queue_name: str): self.client = client self.queue_name = queue_name async def process_dlq(self, handler_map: dict = None): """Process DLQ messages with reason-specific handlers.""" handler_map = handler_map or {} receiver = self.client.get_queue_receiver( queue_name=self.queue_name, sub_queue=ServiceBusSubQueue.DEAD_LETTER ) async with receiver: while True: messages = await receiver.receive_messages( max_message_count=10, max_wait_time=5 ) if not messages: break for msg in messages: reason = msg.dead_letter_reason or "Unknown" handler = handler_map.get(reason, self.default_handler) try: await handler(msg, receiver) except Exception as e: print(f"Error handling DLQ message: {e}") # Leave message in DLQ for manual review async def default_handler(self, msg, receiver): """Default: log and complete.""" print(f"DLQ Message: {msg.message_id}") print(f" Reason: {msg.dead_letter_reason}") print(f" Body: {str(msg.body)[:100]}...") await receiver.complete_message(msg) async def retry_handler(self, msg, receiver): """Retry message by sending back to main queue.""" sender = self.client.get_queue_sender(queue_name=self.queue_name) async with sender: # Create new message from DLQ message retry_msg = ServiceBusMessage( body=msg.body, application_properties={ **(msg.application_properties or {}), "dlq_retry": True, "dlq_reason": msg.dead_letter_reason, "dlq_retry_time": datetime.now(timezone.utc).isoformat() } ) await sender.send_messages(retry_msg) await receiver.complete_message(msg) print(f"Retried message: {msg.message_id}") async def archive_handler(self, msg, receiver): """Archive message to storage for analysis.""" archive_data = { "message_id": msg.message_id, "body": str(msg.body), "dead_letter_reason": msg.dead_letter_reason, "dead_letter_error_description": msg.dead_letter_error_description, "enqueued_time": str(msg.enqueued_time_utc), "delivery_count": msg.delivery_count, "application_properties": dict(msg.application_properties or {}) } # Save to blob storage, database, etc. await archive_to_storage(archive_data) await receiver.complete_message(msg) # Usage processor = DeadLetterProcessor(client, "myqueue") await processor.process_dlq({ "MaxDeliveryCountExceeded": processor.retry_handler, "ValidationFailed": processor.archive_handler, }) ``` ## Reprocessing Strategies ### Selective Retry Based on Age ```python from datetime import datetime, timedelta, timezone async def retry_recent_messages(client, queue_name: str, max_age_hours: int = 24): """Retry only recently dead-lettered messages.""" cutoff_time = datetime.now(timezone.utc) - timedelta(hours=max_age_hours) dlq_receiver = client.get_queue_receiver( queue_name=queue_name, sub_queue=ServiceBusSubQueue.DEAD_LETTER ) sender = client.get_queue_sender(queue_name=queue_name) retried = 0 expired = 0 async with dlq_receiver, sender: while True: messages = await dlq_receiver.receive_messages( max_message_count=50, max_wait_time=5 ) if not messages: break for msg in messages: if msg.enqueued_time_utc > cutoff_time: # Recent enough to retry retry_msg = ServiceBusMessage(body=msg.body) await sender.send_messages(retry_msg) retried += 1 else: # Too old - just complete (discard) expired += 1 await dlq_receiver.complete_message(msg) print(f"Retried: {retried}, Expired: {expired}") ``` ### Retry with Fix ```python async def retry_with_transform(client, queue_name: str, transformer): """Retry messages after applying a fix/transformation.""" dlq_receiver = client.get_queue_receiver( queue_name=queue_name, sub_queue=ServiceBusSubQueue.DEAD_LETTER ) sender = client.get_queue_sender(queue_name=queue_name) async with dlq_receiver, sender: messages = await dlq_receiver.receive_messages(max_message_count=100) for msg in messages: try: # Transform/fix the message fixed_body = transformer(msg.body) retry_msg = ServiceBusMessage( body=fixed_body, application_properties={ "retried_from_dlq": True, "original_message_id": msg.message_id } ) await sender.send_messages(retry_msg) await dlq_receiver.complete_message(msg) except Exception as e: print(f"Could not fix message {msg.message_id}: {e}") # Leave in DLQ # Example transformer: fix JSON encoding issue def fix_json_encoding(body: bytes) -> bytes: text = body.decode('utf-8', errors='replace') data = json.loads(text) return json.dumps(data).encode('utf-8') await retry_with_transform(client, "myqueue", fix_json_encoding) ``` ## Monitoring Dead-Letter Queues ### Count Dead-Lettered Messages ```python from azure.servicebus.management import ServiceBusAdministrationClient async def get_dlq_count(namespace: str, queue_name: str) -> int: """Get count of messages in dead-letter queue.""" admin_client = ServiceBusAdministrationClient( fully_qualified_namespace=namespace, credential=DefaultAzureCredential() ) async with admin_client: runtime_props = await admin_client.get_queue_runtime_properties(queue_name) return runtime_props.dead_letter_message_count # Alert if DLQ has messages dlq_count = await get_dlq_count("myns.servicebus.windows.net", "myqueue") if dlq_count > 0: print(f"ALERT: {dlq_count} messages in DLQ!") ``` ### DLQ Analysis Report ```python async def analyze_dlq(client, queue_name: str) -> dict: """Analyze DLQ contents by reason.""" analysis = { "total": 0, "by_reason": {}, "oldest": None, "newest": None } dlq_receiver = client.get_queue_receiver( queue_name=queue_name, sub_queue=ServiceBusSubQueue.DEAD_LETTER, receive_mode=ServiceBusReceiveMode.PEEK_LOCK ) async with dlq_receiver: # Peek without removing messages = await dlq_receiver.peek_messages(max_message_count=1000) for msg in messages: analysis["total"] += 1 reason = msg.dead_letter_reason or "Unknown" analysis["by_reason"][reason] = analysis["by_reason"].get(reason, 0) + 1 if analysis["oldest"] is None or msg.enqueued_time_utc < analysis["oldest"]: analysis["oldest"] = msg.enqueued_time_utc if analysis["newest"] is None or msg.enqueued_time_utc > analysis["newest"]: analysis["newest"] = msg.enqueued_time_utc return analysis # Generate report report = await analyze_dlq(client, "myqueue") print(f"DLQ Analysis for 'myqueue':") print(f" Total messages: {report['total']}") print(f" Oldest: {report['oldest']}") print(f" Newest: {report['newest']}") print(f" By reason:") for reason, count in report['by_reason'].items(): print(f" {reason}: {count}") ``` ## Preventing Dead-Letters ### Increase Max Delivery Count ```python from azure.servicebus.management import ServiceBusAdministrationClient async def increase_max_delivery(namespace: str, queue_name: str, max_count: int): """Increase max delivery count for a queue.""" admin_client = ServiceBusAdministrationClient( fully_qualified_namespace=namespace, credential=DefaultAzureCredential() ) async with admin_client: queue = await admin_client.get_queue(queue_name) queue.max_delivery_count = max_count await admin_client.update_queue(queue) print(f"Updated max delivery count to {max_count}") ``` ### Implement Proper Error Handling ```python async def resilient_processor(receiver, msg): """Process with proper error categorization.""" try: await process_message(msg) await receiver.complete_message(msg) except (ConnectionError, TimeoutError): # Transient - safe to retry await receiver.abandon_message(msg) except ValidationError: # Bad data - don't retry, dead-letter await receiver.dead_letter_message(msg, reason="ValidationError") except Exception as e: # Unknown - check delivery count if msg.delivery_count >= 3: # Enough retries, dead-letter with context await receiver.dead_letter_message( msg, reason="ProcessingFailed", error_description=f"Failed after {msg.delivery_count} attempts: {e}" ) else: # Still have retries left await receiver.abandon_message(msg) ``` ## Best Practices | Practice | Description | |----------|-------------| | Monitor DLQ counts | Alert when messages appear in DLQ | | Set appropriate max delivery | Balance between retries and DLQ accumulation | | Include context | Always provide reason and description when dead-lettering | | Regular cleanup | Process or archive old DLQ messages | | Categorize errors | Distinguish retryable vs permanent failures | | Implement retry handlers | Automate reprocessing where safe | | Archive for analysis | Keep DLQ data for debugging patterns | -
patterns.md 13.6 KB
# Messaging Patterns Reference Advanced messaging patterns for Azure Service Bus. ## Competing Consumers Multiple receivers processing messages from the same queue in parallel: ```python import asyncio from azure.servicebus.aio import ServiceBusClient from azure.identity.aio import DefaultAzureCredential async def worker(worker_id: int, namespace: str, queue_name: str): """Worker that processes messages from a shared queue.""" credential = DefaultAzureCredential() async with ServiceBusClient( fully_qualified_namespace=namespace, credential=credential ) as client: receiver = client.get_queue_receiver(queue_name=queue_name) async with receiver: while True: messages = await receiver.receive_messages( max_message_count=10, max_wait_time=5 ) if not messages: continue for msg in messages: try: print(f"Worker {worker_id}: Processing {str(msg)}") await process_message(msg) await receiver.complete_message(msg) except Exception as e: print(f"Worker {worker_id}: Error - {e}") await receiver.abandon_message(msg) async def run_workers(num_workers: int, namespace: str, queue_name: str): """Run multiple workers concurrently.""" workers = [ worker(i, namespace, queue_name) for i in range(num_workers) ] await asyncio.gather(*workers) # Run 5 competing consumers asyncio.run(run_workers(5, "myns.servicebus.windows.net", "work-queue")) ``` ## Message Sessions (Ordered Processing) Sessions ensure FIFO ordering for related messages: ```python from azure.servicebus import ServiceBusMessage, NEXT_AVAILABLE_SESSION from azure.servicebus.aio import ServiceBusClient async def send_order_messages(sender, order_id: str, items: list[str]): """Send order items as session messages (processed in order).""" for i, item in enumerate(items): message = ServiceBusMessage( body=item, session_id=order_id, message_id=f"{order_id}-{i}" ) # Messages with same session_id processed in send order await sender.send_messages(message) async def process_order_session(client, queue_name: str): """Process one session's messages in order.""" # Get next available session receiver = client.get_queue_receiver( queue_name=queue_name, session_id=NEXT_AVAILABLE_SESSION, max_wait_time=30 ) async with receiver: session = receiver.session print(f"Processing session: {session.session_id}") # Set session state (for checkpointing) await session.set_state(b"processing") async for msg in receiver: print(f" Item: {str(msg)}") await receiver.complete_message(msg) # Mark session complete await session.set_state(b"completed") print(f"Session {session.session_id} completed") async def session_worker(client, queue_name: str): """Continuously process sessions.""" while True: try: await process_order_session(client, queue_name) except Exception as e: if "No session available" in str(e): await asyncio.sleep(1) # Wait for new sessions else: raise ``` ## Retry Patterns ### Automatic Retry with Max Delivery Count Service Bus automatically retries abandoned messages up to `maxDeliveryCount`: ```python async def process_with_retry_tracking(receiver, msg): """Process message with delivery count awareness.""" delivery_count = msg.delivery_count max_retries = 5 # Should match queue's maxDeliveryCount print(f"Processing message (attempt {delivery_count}/{max_retries})") try: await process_message(msg) await receiver.complete_message(msg) except TransientError: if delivery_count >= max_retries - 1: # Last retry - dead-letter with context await receiver.dead_letter_message( msg, reason="MaxRetriesExceeded", error_description=f"Failed after {delivery_count} attempts" ) else: # Abandon for automatic retry await receiver.abandon_message(msg) except PermanentError as e: # Immediate dead-letter for non-retryable errors await receiver.dead_letter_message( msg, reason="PermanentFailure", error_description=str(e) ) ``` ### Exponential Backoff with Scheduled Retry ```python from datetime import datetime, timedelta, timezone async def process_with_backoff(client, receiver, msg): """Retry with exponential backoff using scheduled messages.""" retry_count = int(msg.application_properties.get("retry_count", 0)) max_retries = 5 try: await process_message(msg) await receiver.complete_message(msg) except TransientError: if retry_count >= max_retries: await receiver.dead_letter_message( msg, reason="MaxRetriesExceeded", error_description=f"Failed after {retry_count} retries" ) return # Calculate backoff: 2^retry seconds (1, 2, 4, 8, 16 seconds) backoff_seconds = 2 ** retry_count retry_time = datetime.now(timezone.utc) + timedelta(seconds=backoff_seconds) # Create retry message with incremented count retry_message = ServiceBusMessage( body=msg.body, application_properties={ **dict(msg.application_properties or {}), "retry_count": retry_count + 1, "original_enqueue_time": str(msg.enqueued_time_utc) } ) # Complete original and schedule retry await receiver.complete_message(msg) sender = client.get_queue_sender(queue_name=receiver.entity_path) async with sender: await sender.schedule_messages(retry_message, retry_time) print(f"Scheduled retry {retry_count + 1} for {retry_time}") ``` ## Request-Response Pattern Using reply queues for synchronous-style communication: ```python import uuid from asyncio import Event, wait_for class RequestResponseClient: """Send requests and wait for correlated responses.""" def __init__(self, client, request_queue: str, response_queue: str): self.client = client self.request_queue = request_queue self.response_queue = response_queue self.pending_requests: dict[str, tuple[Event, dict]] = {} async def start_response_listener(self): """Background task to receive responses.""" receiver = self.client.get_queue_receiver( queue_name=self.response_queue ) async with receiver: async for msg in receiver: correlation_id = msg.correlation_id if correlation_id in self.pending_requests: event, result = self.pending_requests[correlation_id] result["response"] = msg.body event.set() await receiver.complete_message(msg) async def send_request(self, body: str, timeout: float = 30.0) -> bytes: """Send request and wait for response.""" message_id = str(uuid.uuid4()) event = Event() result = {} self.pending_requests[message_id] = (event, result) try: # Send request with reply-to message = ServiceBusMessage( body=body, message_id=message_id, reply_to=self.response_queue ) sender = self.client.get_queue_sender(self.request_queue) async with sender: await sender.send_messages(message) # Wait for response await wait_for(event.wait(), timeout=timeout) return result["response"] finally: del self.pending_requests[message_id] # Request processor (server side) async def process_requests(client, request_queue: str): """Process requests and send responses.""" receiver = client.get_queue_receiver(queue_name=request_queue) async with receiver: async for msg in receiver: # Process request response_body = f"Processed: {str(msg.body)}" # Send response to reply_to queue if msg.reply_to: response = ServiceBusMessage( body=response_body, correlation_id=msg.message_id ) sender = client.get_queue_sender(queue_name=msg.reply_to) async with sender: await sender.send_messages(response) await receiver.complete_message(msg) ``` ## Pub/Sub with Filters Topics with filtered subscriptions: ```python # Publisher sends messages with properties async def publish_events(sender, events: list[dict]): """Publish events with filterable properties.""" for event in events: message = ServiceBusMessage( body=json.dumps(event["data"]), application_properties={ "event_type": event["type"], "priority": event["priority"], "region": event["region"] } ) await sender.send_messages(message) # Subscribers receive filtered messages # (Filters configured via Azure Portal or Management SDK) # # Subscription "high-priority": SqlFilter("priority = 'high'") # Subscription "us-region": SqlFilter("region = 'us'") # Subscription "orders": CorrelationFilter(label='order') async def subscribe_high_priority(client, topic: str): """Receive only high-priority messages.""" receiver = client.get_subscription_receiver( topic_name=topic, subscription_name="high-priority" ) async with receiver: async for msg in receiver: print(f"High priority: {str(msg)}") await receiver.complete_message(msg) ``` ## Transaction Support Atomic operations within a single entity: ```python from azure.servicebus import ServiceBusMessage async def transactional_receive_and_forward(client, source_queue: str, dest_queue: str): """Receive from one queue and send to another atomically.""" receiver = client.get_queue_receiver(queue_name=source_queue) sender = client.get_queue_sender(queue_name=dest_queue) async with receiver, sender: messages = await receiver.receive_messages(max_message_count=1) if messages: msg = messages[0] # Start transaction async with receiver.transaction_scope() as txn: # Transform message new_msg = ServiceBusMessage( body=f"Forwarded: {str(msg.body)}", application_properties={"original_id": msg.message_id} ) # Both operations in same transaction await sender.send_messages(new_msg, transaction=txn) await receiver.complete_message(msg, transaction=txn) # Transaction commits when context exits without error ``` ## Batch Processing with Prefetch Optimize throughput with prefetching: ```python async def batch_processor(client, queue_name: str, batch_size: int = 100): """Process messages in batches with prefetch.""" receiver = client.get_queue_receiver( queue_name=queue_name, prefetch_count=batch_size * 2 # Prefetch 2x batch size ) async with receiver: while True: messages = await receiver.receive_messages( max_message_count=batch_size, max_wait_time=5 ) if not messages: continue # Process batch results = await process_batch([str(m.body) for m in messages]) # Complete successful, abandon failed for msg, success in zip(messages, results): if success: await receiver.complete_message(msg) else: await receiver.abandon_message(msg) ``` ## Message Properties Reference | Property | Type | Description | |----------|------|-------------| | `message_id` | str | Unique message identifier | | `session_id` | str | Session identifier for FIFO | | `correlation_id` | str | For request-response correlation | | `reply_to` | str | Queue for responses | | `reply_to_session_id` | str | Session ID for responses | | `subject` | str | Message label/type | | `to` | str | Destination hint | | `content_type` | str | MIME type (e.g., "application/json") | | `time_to_live` | timedelta | Message expiration | | `scheduled_enqueue_time_utc` | datetime | Scheduled delivery time | | `application_properties` | dict | Custom key-value pairs | ## Best Practices Summary | Pattern | When to Use | |---------|-------------| | Competing Consumers | Scale message processing horizontally | | Sessions | Ordered processing for related messages | | Scheduled Messages | Delayed delivery, exponential backoff | | Request-Response | Synchronous-style communication | | Topics/Filters | Pub/sub with selective routing | | Transactions | Atomic multi-operation guarantees | | Prefetch | High-throughput batch processing |
-
-
scripts
-
setup_servicebus.py 13 KB
#!/usr/bin/env python3 """ Service Bus Administration CLI Tool Create and manage Azure Service Bus queues, topics, and subscriptions. Usage: python setup_servicebus.py queue create myqueue --max-delivery 10 --ttl 3600 python setup_servicebus.py queue info myqueue python setup_servicebus.py topic create mytopic python setup_servicebus.py subscription create mytopic mysub --filter "priority='high'" python setup_servicebus.py dlq count myqueue Environment Variables: SERVICEBUS_FULLY_QUALIFIED_NAMESPACE - Service Bus namespace (e.g., myns.servicebus.windows.net) SERVICEBUS_CONNECTION_STRING - Alternative: full connection string """ import argparse import json import os import sys from datetime import timedelta from typing import Any from azure.identity import DefaultAzureCredential from azure.servicebus.management import ServiceBusAdministrationClient from azure.servicebus.management import ( QueueProperties, TopicProperties, SubscriptionProperties, SqlRuleFilter, CorrelationRuleFilter, ) def get_admin_client() -> ServiceBusAdministrationClient: """Create Service Bus administration client.""" namespace = os.environ.get("SERVICEBUS_FULLY_QUALIFIED_NAMESPACE") conn_str = os.environ.get("SERVICEBUS_CONNECTION_STRING") if conn_str: return ServiceBusAdministrationClient.from_connection_string(conn_str) elif namespace: return ServiceBusAdministrationClient( fully_qualified_namespace=namespace, credential=DefaultAzureCredential() ) else: raise ValueError( "Set SERVICEBUS_FULLY_QUALIFIED_NAMESPACE or SERVICEBUS_CONNECTION_STRING" ) def create_queue( client: ServiceBusAdministrationClient, name: str, max_delivery_count: int = 10, ttl_seconds: int | None = None, lock_duration_seconds: int = 60, enable_sessions: bool = False, enable_partitioning: bool = False, ) -> dict[str, Any]: """Create a Service Bus queue.""" kwargs = { "max_delivery_count": max_delivery_count, "lock_duration": timedelta(seconds=lock_duration_seconds), "requires_session": enable_sessions, "enable_partitioning": enable_partitioning, } if ttl_seconds: kwargs["default_message_time_to_live"] = timedelta(seconds=ttl_seconds) queue = client.create_queue(name, **kwargs) return { "name": queue.name, "max_delivery_count": queue.max_delivery_count, "lock_duration": str(queue.lock_duration), "requires_session": queue.requires_session, "enable_partitioning": queue.enable_partitioning, } def get_queue_info(client: ServiceBusAdministrationClient, name: str) -> dict[str, Any]: """Get queue properties and runtime info.""" queue = client.get_queue(name) runtime = client.get_queue_runtime_properties(name) return { "name": queue.name, "max_delivery_count": queue.max_delivery_count, "lock_duration": str(queue.lock_duration), "default_ttl": str(queue.default_message_time_to_live), "requires_session": queue.requires_session, "enable_partitioning": queue.enable_partitioning, "runtime": { "active_message_count": runtime.active_message_count, "dead_letter_message_count": runtime.dead_letter_message_count, "scheduled_message_count": runtime.scheduled_message_count, "total_message_count": runtime.total_message_count, }, } def create_topic( client: ServiceBusAdministrationClient, name: str, ttl_seconds: int | None = None, enable_partitioning: bool = False, ) -> dict[str, Any]: """Create a Service Bus topic.""" kwargs = {"enable_partitioning": enable_partitioning} if ttl_seconds: kwargs["default_message_time_to_live"] = timedelta(seconds=ttl_seconds) topic = client.create_topic(name, **kwargs) return {"name": topic.name, "enable_partitioning": topic.enable_partitioning} def create_subscription( client: ServiceBusAdministrationClient, topic_name: str, subscription_name: str, sql_filter: str | None = None, max_delivery_count: int = 10, lock_duration_seconds: int = 60, enable_sessions: bool = False, ) -> dict[str, Any]: """Create a subscription with optional filter.""" subscription = client.create_subscription( topic_name=topic_name, subscription_name=subscription_name, max_delivery_count=max_delivery_count, lock_duration=timedelta(seconds=lock_duration_seconds), requires_session=enable_sessions, ) result = { "topic": topic_name, "subscription": subscription.name, "max_delivery_count": subscription.max_delivery_count, "requires_session": subscription.requires_session, } # Add SQL filter if provided if sql_filter: # Delete default rule and create filtered rule client.delete_rule(topic_name, subscription_name, "$Default") client.create_rule( topic_name=topic_name, subscription_name=subscription_name, rule_name="CustomFilter", filter=SqlRuleFilter(sql_filter), ) result["filter"] = sql_filter return result def get_dlq_count( client: ServiceBusAdministrationClient, name: str, is_subscription: bool = False, topic_name: str | None = None, ) -> dict[str, Any]: """Get dead-letter queue message count.""" if is_subscription: runtime = client.get_subscription_runtime_properties(topic_name, name) else: runtime = client.get_queue_runtime_properties(name) return { "entity": f"{topic_name}/{name}" if is_subscription else name, "dead_letter_message_count": runtime.dead_letter_message_count, "active_message_count": runtime.active_message_count, } def list_entities( client: ServiceBusAdministrationClient, entity_type: str, topic_name: str | None = None, ) -> list[str]: """List queues, topics, or subscriptions.""" if entity_type == "queues": return [q.name for q in client.list_queues()] elif entity_type == "topics": return [t.name for t in client.list_topics()] elif entity_type == "subscriptions" and topic_name: return [s.name for s in client.list_subscriptions(topic_name)] else: return [] def main(): parser = argparse.ArgumentParser( description="Manage Azure Service Bus entities", formatter_class=argparse.RawDescriptionHelpFormatter, ) subparsers = parser.add_subparsers(dest="entity", required=True) # Queue commands queue_parser = subparsers.add_parser("queue", help="Queue operations") queue_subparsers = queue_parser.add_subparsers(dest="action", required=True) queue_create = queue_subparsers.add_parser("create", help="Create queue") queue_create.add_argument("name", help="Queue name") queue_create.add_argument( "--max-delivery", type=int, default=10, help="Max delivery count" ) queue_create.add_argument("--ttl", type=int, help="Default message TTL in seconds") queue_create.add_argument( "--lock-duration", type=int, default=60, help="Lock duration in seconds" ) queue_create.add_argument("--sessions", action="store_true", help="Enable sessions") queue_create.add_argument( "--partitioned", action="store_true", help="Enable partitioning" ) queue_info = queue_subparsers.add_parser("info", help="Get queue info") queue_info.add_argument("name", help="Queue name") queue_list = queue_subparsers.add_parser("list", help="List queues") queue_delete = queue_subparsers.add_parser("delete", help="Delete queue") queue_delete.add_argument("name", help="Queue name") # Topic commands topic_parser = subparsers.add_parser("topic", help="Topic operations") topic_subparsers = topic_parser.add_subparsers(dest="action", required=True) topic_create = topic_subparsers.add_parser("create", help="Create topic") topic_create.add_argument("name", help="Topic name") topic_create.add_argument("--ttl", type=int, help="Default message TTL in seconds") topic_create.add_argument( "--partitioned", action="store_true", help="Enable partitioning" ) topic_list = topic_subparsers.add_parser("list", help="List topics") topic_delete = topic_subparsers.add_parser("delete", help="Delete topic") topic_delete.add_argument("name", help="Topic name") # Subscription commands sub_parser = subparsers.add_parser("subscription", help="Subscription operations") sub_subparsers = sub_parser.add_subparsers(dest="action", required=True) sub_create = sub_subparsers.add_parser("create", help="Create subscription") sub_create.add_argument("topic", help="Topic name") sub_create.add_argument("name", help="Subscription name") sub_create.add_argument("--filter", help="SQL filter expression") sub_create.add_argument( "--max-delivery", type=int, default=10, help="Max delivery count" ) sub_create.add_argument("--sessions", action="store_true", help="Enable sessions") sub_list = sub_subparsers.add_parser("list", help="List subscriptions") sub_list.add_argument("topic", help="Topic name") sub_delete = sub_subparsers.add_parser("delete", help="Delete subscription") sub_delete.add_argument("topic", help="Topic name") sub_delete.add_argument("name", help="Subscription name") # DLQ commands dlq_parser = subparsers.add_parser("dlq", help="Dead-letter queue operations") dlq_subparsers = dlq_parser.add_subparsers(dest="action", required=True) dlq_count = dlq_subparsers.add_parser("count", help="Get DLQ message count") dlq_count.add_argument("name", help="Queue or subscription name") dlq_count.add_argument("--topic", help="Topic name (for subscriptions)") parser.add_argument("--output", "-o", choices=["json", "text"], default="text") args = parser.parse_args() try: client = get_admin_client() except ValueError as e: print(f"Error: {e}") sys.exit(1) result = None try: if args.entity == "queue": if args.action == "create": result = create_queue( client, args.name, max_delivery_count=args.max_delivery, ttl_seconds=args.ttl, lock_duration_seconds=args.lock_duration, enable_sessions=args.sessions, enable_partitioning=args.partitioned, ) elif args.action == "info": result = get_queue_info(client, args.name) elif args.action == "list": result = list_entities(client, "queues") elif args.action == "delete": client.delete_queue(args.name) result = {"deleted": args.name} elif args.entity == "topic": if args.action == "create": result = create_topic( client, args.name, ttl_seconds=args.ttl, enable_partitioning=args.partitioned, ) elif args.action == "list": result = list_entities(client, "topics") elif args.action == "delete": client.delete_topic(args.name) result = {"deleted": args.name} elif args.entity == "subscription": if args.action == "create": result = create_subscription( client, args.topic, args.name, sql_filter=args.filter, max_delivery_count=args.max_delivery, enable_sessions=args.sessions, ) elif args.action == "list": result = list_entities(client, "subscriptions", args.topic) elif args.action == "delete": client.delete_subscription(args.topic, args.name) result = {"deleted": f"{args.topic}/{args.name}"} elif args.entity == "dlq": if args.action == "count": result = get_dlq_count( client, args.name, is_subscription=bool(args.topic), topic_name=args.topic, ) except Exception as e: print(f"Error: {e}") sys.exit(1) # Output result if result: if args.output == "json": print(json.dumps(result, indent=2, default=str)) else: if isinstance(result, list): for item in result: print(f" - {item}") elif isinstance(result, dict): for key, value in result.items(): if isinstance(value, dict): print(f"{key}:") for k, v in value.items(): print(f" {k}: {v}") else: print(f"{key}: {value}") if __name__ == "__main__": main()
-
-
SKILL.md 10.1 KB
--- name: azure-servicebus-py description: | Azure Service Bus SDK for Python messaging. Use for queues, topics, subscriptions, and enterprise messaging patterns. Triggers: "service bus", "ServiceBusClient", "queue", "topic", "subscription", "message broker". license: MIT metadata: author: Microsoft version: "1.0.0" package: azure-servicebus --- # 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 |
Comments (0)
Sign in to join the conversation.
Reviews (0)
No reviews yet.
No comments yet.