GitHub Copilot ChatGPT Claude Codex CLI Cursor opencode Skill Text

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".

Ciza · 0 points · 24 views 0 listing impressions 0 install-command copies
Virus-scanned Reviewed automatically before listing.

Full trust report

Download microsoft-skills-.github_plugins_azure-sdk-python_skills_azure-servicebus-py-e58528d.zip · 13 KB
Part of microsoft/skills — 195 skills

Install

skills CLI npx skills add https://github.com/microsoft/skills/tree/main/.github/plugins/azure-sdk-python/skills/azure-servicebus-py
Claude Code claude plugin marketplace add https://llmmart.ai/marketplace.json && claude plugin install microsoft-skills@llmmart
Git 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:

  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.

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

  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 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.

No comments yet.

Reviews (0)

No reviews yet.

Related