GitHub Copilot
ChatGPT
Claude
Codex CLI
Cursor
opencode
Skill
Text
azure-servicebus-ts
Build messaging applications using Azure Service Bus SDK for JavaScript (@azure/service-bus). Use when implementing queues, topics/subscriptions, message sessions, dead-letter handling, or enterprise messaging patterns.
Virus-scanned
Reviewed automatically before listing.
Download
microsoft-skills-.github_plugins_azure-sdk-typescript_skills_azure-servicebus-ts-e58528d.zip · 9 KB
Install
skills CLI
npx skills add https://github.com/microsoft/skills/tree/main/.github/plugins/azure-sdk-typescript/skills/azure-servicebus-ts
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 TypeScript
Enterprise messaging with queues, topics, and subscriptions.
Installation
npm install @azure/service-bus @azure/identity
Environment Variables
SERVICEBUS_NAMESPACE=<namespace>.servicebus.windows.net
SERVICEBUS_QUEUE_NAME=my-queue
SERVICEBUS_TOPIC_NAME=my-topic
SERVICEBUS_SUBSCRIPTION_NAME=my-subscription
AZURE_TOKEN_CREDENTIALS=prod # Required only if DefaultAzureCredential is used in production
Authentication
import { ServiceBusClient } from "@azure/service-bus";
import { DefaultAzureCredential, ManagedIdentityCredential } from "@azure/identity";
// Local dev: DefaultAzureCredential. Production: set AZURE_TOKEN_CREDENTIALS=prod or AZURE_TOKEN_CREDENTIALS=<specific_credential>
const credential = new DefaultAzureCredential({requiredEnvVars: ["AZURE_TOKEN_CREDENTIALS"]});
// Or use a specific credential directly in production:
// See https://learn.microsoft.com/javascript/api/overview/azure/identity-readme?view=azure-node-latest#credential-classes
// const credential = new ManagedIdentityCredential();
const fullyQualifiedNamespace = process.env.SERVICEBUS_NAMESPACE!;
const client = new ServiceBusClient(fullyQualifiedNamespace, credential);
Core Workflow
Send Messages to Queue
const sender = client.createSender("my-queue");
// Single message
await sender.sendMessages({
body: { orderId: "12345", amount: 99.99 },
contentType: "application/json",
});
// Batch messages
const batch = await sender.createMessageBatch();
batch.tryAddMessage({ body: "Message 1" });
batch.tryAddMessage({ body: "Message 2" });
await sender.sendMessages(batch);
await sender.close();
Receive Messages from Queue
const receiver = client.createReceiver("my-queue");
// Receive batch
const messages = await receiver.receiveMessages(10, { maxWaitTimeInMs: 5000 });
for (const message of messages) {
console.log(`Received: ${message.body}`);
await receiver.completeMessage(message);
}
await receiver.close();
Subscribe to Messages (Event-Driven)
const receiver = client.createReceiver("my-queue");
const subscription = receiver.subscribe({
processMessage: async (message) => {
console.log(`Processing: ${message.body}`);
// Message auto-completed on success
},
processError: async (args) => {
console.error(`Error: ${args.error}`);
},
});
// Stop after some time
setTimeout(async () => {
await subscription.close();
await receiver.close();
}, 60000);
Topics and Subscriptions
// Send to topic
const topicSender = client.createSender("my-topic");
await topicSender.sendMessages({
body: { event: "order.created", data: { orderId: "123" } },
applicationProperties: { eventType: "order.created" },
});
// Receive from subscription
const subscriptionReceiver = client.createReceiver("my-topic", "my-subscription");
const messages = await subscriptionReceiver.receiveMessages(10);
Message Sessions
// Send session message
const sender = client.createSender("session-queue");
await sender.sendMessages({
body: { step: 1, data: "First step" },
sessionId: "workflow-123",
});
// Receive session messages
const sessionReceiver = await client.acceptSession("session-queue", "workflow-123");
const messages = await sessionReceiver.receiveMessages(10);
// Get/set session state
const state = await sessionReceiver.getSessionState();
await sessionReceiver.setSessionState(Buffer.from(JSON.stringify({ progress: 50 })));
await sessionReceiver.close();
Dead-Letter Handling
// Move to dead-letter
await receiver.deadLetterMessage(message, {
deadLetterReason: "Validation failed",
deadLetterErrorDescription: "Missing required field: orderId",
});
// Process dead-letter queue
const dlqReceiver = client.createReceiver("my-queue", { subQueueType: "deadLetter" });
const dlqMessages = await dlqReceiver.receiveMessages(10);
for (const msg of dlqMessages) {
console.log(`DLQ Reason: ${msg.deadLetterReason}`);
// Reprocess or log
await dlqReceiver.completeMessage(msg);
}
Scheduled Messages
const sender = client.createSender("my-queue");
// Schedule for future delivery
const scheduledTime = new Date(Date.now() + 60000); // 1 minute from now
const sequenceNumber = await sender.scheduleMessages(
{ body: "Delayed message" },
scheduledTime
);
// Cancel scheduled message
await sender.cancelScheduledMessages(sequenceNumber);
Message Deferral
// Defer message for later
await receiver.deferMessage(message);
// Receive deferred message by sequence number
const deferredMessage = await receiver.receiveDeferredMessages(message.sequenceNumber!);
await receiver.completeMessage(deferredMessage[0]);
Peek Messages (Non-Destructive)
const receiver = client.createReceiver("my-queue");
// Peek without removing
const peekedMessages = await receiver.peekMessages(10);
for (const msg of peekedMessages) {
console.log(`Peeked: ${msg.body}`);
}
Key Types
import {
ServiceBusClient,
ServiceBusSender,
ServiceBusReceiver,
ServiceBusSessionReceiver,
ServiceBusMessage,
ServiceBusReceivedMessage,
ProcessMessageCallback,
ProcessErrorCallback,
} from "@azure/service-bus";
Receive Modes
// Peek-Lock (default) - message locked until completed/abandoned
const receiver = client.createReceiver("my-queue", { receiveMode: "peekLock" });
await receiver.completeMessage(message); // Remove from queue
await receiver.abandonMessage(message); // Return to queue
await receiver.deferMessage(message); // Defer for later
await receiver.deadLetterMessage(message); // Move to DLQ
// Receive-and-Delete - message removed immediately
const receiver = client.createReceiver("my-queue", { receiveMode: "receiveAndDelete" });
Best Practices
- Use Microsoft Entra Token Credential - Use
DefaultAzureCredentialfor local development; useManagedIdentityCredentialorWorkloadIdentityCredentialfor production - Reuse clients - Create
ServiceBusClientonce, share across senders/receivers - Close resources - Always close senders/receivers when done
- Handle errors - Implement
processErrorcallback for subscription receivers - Use sessions for ordering - When message order matters within a group
- Configure dead-letter - Always handle DLQ messages
- Batch sends - Use
createMessageBatch()for multiple messages
Reference Documentation
For detailed patterns, see:
- Queues vs Topics Patterns - Queue/topic patterns, sessions, receive modes, message settlement
- Error Handling and Reliability - ServiceBusError codes, DLQ handling, lock renewal, graceful shutdown
Files (skills)
-
references
-
error-handling.md 11.3 KB
# Error Handling and Reliability Comprehensive error handling patterns for @azure/service-bus. ## ServiceBusError All Service Bus errors extend `ServiceBusError` with a `code` property: ```typescript import { ServiceBusError } from "@azure/service-bus"; try { await sender.sendMessages(message); } catch (error) { if (error instanceof ServiceBusError) { console.log(`Code: ${error.code}`); console.log(`Message: ${error.message}`); console.log(`Retryable: ${error.retryable}`); handleServiceBusError(error); } } ``` ## Error Codes Reference | Code | Description | Retryable | Action | |------|-------------|-----------|--------| | `GeneralError` | Unspecified error | Maybe | Log and investigate | | `MessagingEntityNotFound` | Queue/topic/subscription doesn't exist | No | Check entity name, create entity | | `MessageLockLost` | Lock expired before settlement | No | Message will be redelivered | | `MessageNotFound` | Message no longer available | No | Already processed or expired | | `MessageSizeExceeded` | Message too large (>256KB standard, >100MB premium) | No | Reduce message size or use claim check | | `MessagingEntityAlreadyExists` | Entity already exists | No | Use existing entity | | `MessagingEntityDisabled` | Entity is disabled | No | Enable entity in portal | | `QuotaExceeded` | Namespace quota exceeded | No | Delete messages, upgrade tier | | `ServiceBusy` | Service temporarily overloaded | Yes | Retry with backoff | | `ServiceTimeout` | Operation timed out | Yes | Retry with backoff | | `ServiceCommunicationProblem` | Network/connection issue | Yes | Retry with backoff | | `SessionCannotBeLocked` | Session locked by another receiver | Yes | Retry or use different session | | `SessionLockLost` | Session lock expired | No | Re-accept session | | `UnauthorizedAccess` | Authentication/authorization failed | No | Check credentials/permissions | ## Error Handling by Code ```typescript function handleServiceBusError(error: ServiceBusError): void { switch (error.code) { case "MessagingEntityNotFound": console.error(`Entity not found: ${error.message}`); // Create entity or fix configuration break; case "MessageLockLost": console.warn("Message lock lost - will be redelivered"); // No action needed, message returns to queue break; case "MessageSizeExceeded": console.error("Message too large - use claim check pattern"); // Store payload in blob, send reference break; case "QuotaExceeded": console.error("Quota exceeded - namespace is full"); // Alert ops team, consider cleanup or upgrade break; case "ServiceBusy": case "ServiceTimeout": case "ServiceCommunicationProblem": if (error.retryable) { console.warn(`Transient error: ${error.code} - will retry`); // SDK handles retry automatically } break; case "SessionCannotBeLocked": console.warn("Session busy - trying another session"); // Use acceptNextSession() instead break; case "SessionLockLost": console.warn("Session lock lost - re-accepting session"); // Re-accept the session break; case "UnauthorizedAccess": console.error("Authorization failed - check credentials"); // Verify RBAC roles or connection string break; default: console.error(`Unexpected error: ${error.code} - ${error.message}`); } } ``` ## ProcessErrorArgs in Subscribe When using `receiver.subscribe()`, errors are delivered via `processError`: ```typescript import { ProcessErrorArgs } from "@azure/service-bus"; receiver.subscribe({ processMessage: async (message) => { // Process message }, processError: async (args: ProcessErrorArgs) => { console.error(`Error source: ${args.errorSource}`); console.error(`Entity path: ${args.entityPath}`); console.error(`Namespace: ${args.fullyQualifiedNamespace}`); if (args.error instanceof ServiceBusError) { console.error(`Code: ${args.error.code}`); console.error(`Retryable: ${args.error.retryable}`); } // Error sources: // - "receive" - Error receiving messages // - "processMessageCallback" - Error in your processMessage handler // - "renewLock" - Error renewing message lock // - "complete" / "abandon" / "deadLetter" - Settlement errors switch (args.errorSource) { case "receive": console.log("Connection issue - SDK will reconnect"); break; case "processMessageCallback": console.log("Bug in message handler - fix code"); break; case "renewLock": console.log("Lock renewal failed - message may be redelivered"); break; } }, }); ``` ## Dead Letter Queue Handling Messages that can't be processed go to the dead letter queue (DLQ): ```typescript // Create DLQ receiver const dlqReceiver = client.createReceiver("my-queue", { subQueueType: "deadLetter", }); // Process dead letters const deadLetters = await dlqReceiver.receiveMessages(10); for (const message of deadLetters) { console.log(`Dead letter reason: ${message.deadLetterReason}`); console.log(`Error description: ${message.deadLetterErrorDescription}`); console.log(`Original body: ${JSON.stringify(message.body)}`); console.log(`Delivery count: ${message.deliveryCount}`); console.log(`Enqueued time: ${message.enqueuedTimeUtc}`); // Analyze and fix the issue if (canReprocess(message)) { // Resend to main queue const sender = client.createSender("my-queue"); await sender.sendMessages({ body: message.body }); await sender.close(); } // Remove from DLQ await dlqReceiver.completeMessage(message); } await dlqReceiver.close(); ``` ### DLQ for Topics ```typescript // DLQ receiver for topic subscription const dlqReceiver = client.createReceiver( "my-topic", "my-subscription", { subQueueType: "deadLetter" } ); ``` ### Automatic Dead Lettering Messages are automatically dead-lettered when: - `deliveryCount` exceeds `maxDeliveryCount` (default: 10) - Message TTL expires (if `deadLetteringOnMessageExpiration` is enabled) - Subscription filter evaluation fails ## Graceful Shutdown Always close resources in the correct order: ```typescript const client = new ServiceBusClient(namespace, credential); const sender = client.createSender("my-queue"); const receiver = client.createReceiver("my-queue"); // Subscribe to messages const subscription = receiver.subscribe({ processMessage: async (msg) => { /* ... */ }, processError: async (args) => { /* ... */ }, }); // Graceful shutdown handler async function shutdown(): Promise<void> { console.log("Shutting down..."); // 1. Stop receiving new messages await subscription.close(); // 2. Close receiver (waits for in-flight messages) await receiver.close(); // 3. Close sender (waits for pending sends) await sender.close(); // 4. Close client last await client.close(); console.log("Shutdown complete"); } process.on("SIGTERM", shutdown); process.on("SIGINT", shutdown); ``` ## Lock Renewal ### Automatic Lock Renewal (Subscribe) ```typescript // subscribe() automatically renews locks receiver.subscribe({ processMessage: async (message) => { // Lock is auto-renewed while processing await longRunningOperation(message.body); }, processError: async (args) => { /* ... */ }, }, { // Max time to auto-renew (default: 5 minutes) maxAutoLockRenewalDurationInMs: 10 * 60 * 1000, // 10 minutes }); ``` ### Manual Lock Renewal (receiveMessages) ```typescript const [message] = await receiver.receiveMessages(1); // For long processing, manually renew lock const renewalInterval = setInterval(async () => { try { await receiver.renewMessageLock(message); console.log("Lock renewed"); } catch (error) { console.error("Lock renewal failed:", error); clearInterval(renewalInterval); } }, 30000); // Renew every 30 seconds try { await longRunningOperation(message.body); await receiver.completeMessage(message); } finally { clearInterval(renewalInterval); } ``` ### Session Lock Renewal ```typescript const sessionReceiver = await client.acceptSession("my-queue", "session-123"); // Auto-renewal for sessions const sessionReceiver = await client.acceptSession("my-queue", "session-123", { maxAutoLockRenewalDurationInMs: 10 * 60 * 1000, }); // Manual renewal await sessionReceiver.renewSessionLock(); ``` ## Connection Recovery The SDK automatically handles connection recovery: ```typescript // SDK reconnects automatically on transient failures // No manual reconnection code needed receiver.subscribe({ processMessage: async (message) => { // Processing continues after reconnection }, processError: async (args) => { if (args.errorSource === "receive") { // Connection issue - SDK is reconnecting console.log("Connection lost, SDK reconnecting..."); } }, }); ``` ## Retry Configuration Configure retry behavior at client level: ```typescript import { ServiceBusClient } from "@azure/service-bus"; const client = new ServiceBusClient(namespace, credential, { retryOptions: { maxRetries: 3, retryDelayInMs: 1000, maxRetryDelayInMs: 30000, mode: "Exponential", // or "Fixed" }, }); ``` ## Idempotent Processing Design for at-least-once delivery: ```typescript receiver.subscribe({ processMessage: async (message) => { const messageId = message.messageId; // Check if already processed (use database, Redis, etc.) if (await isAlreadyProcessed(messageId)) { console.log(`Duplicate message: ${messageId}`); return; // Message will be completed } // Process message await processOrder(message.body); // Mark as processed await markAsProcessed(messageId); }, processError: async (args) => { /* ... */ }, }); ``` ## Poison Message Handling Handle messages that repeatedly fail: ```typescript receiver.subscribe({ processMessage: async (message) => { // Check delivery count if (message.deliveryCount > 5) { console.warn(`Message ${message.messageId} failed ${message.deliveryCount} times`); // Dead letter with reason await receiver.deadLetterMessage(message, { deadLetterReason: "MaxRetriesExceeded", deadLetterErrorDescription: `Failed after ${message.deliveryCount} attempts`, }); return; } try { await processMessage(message.body); } catch (error) { // Abandon to retry await receiver.abandonMessage(message, { propertiesToModify: { lastError: error.message, lastAttempt: new Date().toISOString(), }, }); } }, processError: async (args) => { /* ... */ }, }, { autoCompleteMessages: false, // Manual settlement }); ``` ## Best Practices Summary 1. **Always handle errors** - Implement `processError` callback 2. **Use peek-lock mode** - Ensures at-least-once delivery 3. **Design for idempotency** - Messages may be delivered multiple times 4. **Monitor dead letter queues** - Set up alerts for DLQ messages 5. **Configure appropriate lock duration** - Match to processing time 6. **Use auto-lock renewal** - For long-running operations 7. **Graceful shutdown** - Close resources in correct order 8. **Log error codes** - Helps diagnose issues 9. **Set maxDeliveryCount** - Prevent infinite retry loops 10. **Handle poison messages** - Dead letter after max retries -
queues-topics.md 9.5 KB
# Queues vs Topics Patterns Detailed patterns for Azure Service Bus queues and topics with @azure/service-bus. ## Queue Patterns (Point-to-Point) Queues deliver each message to exactly one consumer. Use for work distribution. ### Basic Queue Sender ```typescript import { ServiceBusClient, ServiceBusMessage } from "@azure/service-bus"; import { DefaultAzureCredential } from "@azure/identity"; const client = new ServiceBusClient( process.env.SERVICEBUS_NAMESPACE!, new DefaultAzureCredential() ); const sender = client.createSender("order-queue"); // Send single message const message: ServiceBusMessage = { body: { orderId: "12345", amount: 99.99 }, contentType: "application/json", messageId: "unique-id-12345", correlationId: "request-abc", applicationProperties: { priority: "high", source: "web-app", }, }; await sender.sendMessages(message); await sender.close(); ``` ### Batch Sending (Recommended for Multiple Messages) ```typescript const sender = client.createSender("order-queue"); // Create batch with size limits const batch = await sender.createMessageBatch(); const orders = [ { orderId: "001", amount: 50 }, { orderId: "002", amount: 75 }, { orderId: "003", amount: 100 }, ]; for (const order of orders) { const added = batch.tryAddMessage({ body: order }); if (!added) { // Batch is full - send current batch and create new one await sender.sendMessages(batch); const newBatch = await sender.createMessageBatch(); newBatch.tryAddMessage({ body: order }); } } // Send remaining messages if (batch.count > 0) { await sender.sendMessages(batch); } await sender.close(); ``` ### Queue Receiver (Pull Model) ```typescript const receiver = client.createReceiver("order-queue", { receiveMode: "peekLock", // Default - message locked until settled }); // Receive batch of messages const messages = await receiver.receiveMessages(10, { maxWaitTimeInMs: 5000, }); for (const message of messages) { try { const order = message.body; console.log(`Processing order: ${order.orderId}`); // Process successfully - remove from queue await receiver.completeMessage(message); } catch (error) { // Processing failed - return to queue for retry await receiver.abandonMessage(message); } } await receiver.close(); ``` ### Queue Subscriber (Push Model - Event-Driven) ```typescript const receiver = client.createReceiver("order-queue"); const subscription = receiver.subscribe({ processMessage: async (message) => { console.log(`Received: ${JSON.stringify(message.body)}`); // Message auto-completed on success when using subscribe() }, processError: async (args) => { console.error(`Error source: ${args.errorSource}`); console.error(`Error: ${args.error.message}`); if (args.error.code === "MessageLockLost") { console.log("Message lock expired - will be redelivered"); } }, }, { autoCompleteMessages: true, // Default maxConcurrentCalls: 5, }); // Graceful shutdown process.on("SIGTERM", async () => { await subscription.close(); await receiver.close(); await client.close(); }); ``` ## Topic Patterns (Publish-Subscribe) Topics deliver messages to multiple subscriptions. Use for fan-out scenarios. ### Topic Publisher ```typescript // Sending to a topic is identical to sending to a queue const sender = client.createSender("order-events"); await sender.sendMessages({ body: { event: "order.created", orderId: "12345" }, applicationProperties: { eventType: "order.created", region: "us-west", }, subject: "orders/created", // Can be used for filtering }); await sender.close(); ``` ### Subscription Receiver ```typescript // Receive from a specific subscription const receiver = client.createReceiver( "order-events", // Topic name "billing-service" // Subscription name ); const messages = await receiver.receiveMessages(10); for (const message of messages) { console.log(`Billing received: ${message.body.event}`); await receiver.completeMessage(message); } ``` ### Multiple Subscriptions (Fan-Out) ```typescript // Each subscription gets a copy of every message // Create receivers for different services const billingReceiver = client.createReceiver("order-events", "billing"); const inventoryReceiver = client.createReceiver("order-events", "inventory"); const notificationReceiver = client.createReceiver("order-events", "notifications"); // Each processes independently billingReceiver.subscribe({ processMessage: async (msg) => { await processBilling(msg.body); }, processError: async (args) => console.error(args.error), }); inventoryReceiver.subscribe({ processMessage: async (msg) => { await updateInventory(msg.body); }, processError: async (args) => console.error(args.error), }); ``` ## Session-Enabled Queues (Ordered Processing) Sessions guarantee FIFO ordering within a session and enable stateful processing. ### Send Session Messages ```typescript const sender = client.createSender("workflow-queue"); // All messages with same sessionId processed in order by same receiver await sender.sendMessages([ { body: { step: 1, action: "validate" }, sessionId: "order-123" }, { body: { step: 2, action: "charge" }, sessionId: "order-123" }, { body: { step: 3, action: "fulfill" }, sessionId: "order-123" }, ]); ``` ### Accept Specific Session ```typescript // Lock a specific session const sessionReceiver = await client.acceptSession( "workflow-queue", "order-123" ); // Process all messages for this session in order const messages = await sessionReceiver.receiveMessages(10); for (const message of messages) { console.log(`Step ${message.body.step}: ${message.body.action}`); await sessionReceiver.completeMessage(message); } await sessionReceiver.close(); ``` ### Accept Next Available Session ```typescript // Accept any available session (load balancing) const sessionReceiver = await client.acceptNextSession("workflow-queue", { maxAutoLockRenewalDurationInMs: 300000, // 5 minutes }); console.log(`Processing session: ${sessionReceiver.sessionId}`); // Get/set session state for checkpointing const state = await sessionReceiver.getSessionState(); if (state) { const checkpoint = JSON.parse(state.toString()); console.log(`Resuming from step: ${checkpoint.lastStep}`); } // Process messages... // Save checkpoint await sessionReceiver.setSessionState( Buffer.from(JSON.stringify({ lastStep: 3, status: "complete" })) ); await sessionReceiver.close(); ``` ## Receive Modes ### Peek-Lock (Default - At-Least-Once) ```typescript const receiver = client.createReceiver("my-queue", { receiveMode: "peekLock", }); const [message] = await receiver.receiveMessages(1); // Message is locked - other receivers can't see it // You MUST settle the message: await receiver.completeMessage(message); // Success - remove from queue await receiver.abandonMessage(message); // Failure - return to queue await receiver.deferMessage(message); // Defer - retrieve by sequence number await receiver.deadLetterMessage(message); // Move to DLQ ``` ### Receive-and-Delete (At-Most-Once) ```typescript const receiver = client.createReceiver("my-queue", { receiveMode: "receiveAndDelete", }); const [message] = await receiver.receiveMessages(1); // Message is immediately deleted - no settlement needed // If processing fails, message is lost console.log(`Received and deleted: ${message.body}`); ``` ## Message Settlement Actions ```typescript const receiver = client.createReceiver("my-queue"); const [message] = await receiver.receiveMessages(1); // COMPLETE: Processing succeeded, remove message await receiver.completeMessage(message); // ABANDON: Processing failed, return to queue for retry // Increments deliveryCount await receiver.abandonMessage(message, { propertiesToModify: { retryReason: "timeout" }, }); // DEFER: Can't process now, retrieve later by sequence number await receiver.deferMessage(message); // Later: await receiver.receiveDeferredMessages([message.sequenceNumber!]); // DEAD LETTER: Poison message, move to DLQ await receiver.deadLetterMessage(message, { deadLetterReason: "ValidationFailed", deadLetterErrorDescription: "Missing required field: customerId", }); ``` ## Scheduled Messages ```typescript const sender = client.createSender("my-queue"); // Schedule for future delivery const scheduledTime = new Date(Date.now() + 60 * 60 * 1000); // 1 hour from now const sequenceNumbers = await sender.scheduleMessages( [ { body: "Reminder: Complete your order" }, { body: "Your cart is waiting" }, ], scheduledTime ); console.log(`Scheduled messages: ${sequenceNumbers}`); // Cancel scheduled messages if needed await sender.cancelScheduledMessages(sequenceNumbers); ``` ## Peek Messages (Non-Destructive) ```typescript const receiver = client.createReceiver("my-queue"); // Peek without removing or locking const peekedMessages = await receiver.peekMessages(10); for (const msg of peekedMessages) { console.log(`Peeked: ${msg.body}`); console.log(`Sequence: ${msg.sequenceNumber}`); console.log(`Enqueued: ${msg.enqueuedTimeUtc}`); } // Peek from specific sequence number const moreMessages = await receiver.peekMessages(10, { fromSequenceNumber: 100n, }); ``` ## When to Use Queues vs Topics | Scenario | Use | |----------|-----| | Work distribution (one consumer per message) | Queue | | Fan-out (multiple consumers per message) | Topic | | Competing consumers (load balancing) | Queue | | Event broadcasting | Topic | | Request-reply pattern | Queue with reply-to | | Ordered processing | Session-enabled Queue | | Filtered delivery | Topic with subscription filters |
-
-
SKILL.md 7.1 KB
--- name: azure-servicebus-ts description: Build messaging applications using Azure Service Bus SDK for JavaScript (@azure/service-bus). Use when implementing queues, topics/subscriptions, message sessions, dead-letter handling, or enterprise messaging patterns. license: MIT metadata: author: Microsoft version: "1.0.0" package: '@azure/service-bus' --- # Azure Service Bus SDK for TypeScript Enterprise messaging with queues, topics, and subscriptions. ## Installation ```bash npm install @azure/service-bus @azure/identity ``` ## Environment Variables ```bash SERVICEBUS_NAMESPACE=<namespace>.servicebus.windows.net SERVICEBUS_QUEUE_NAME=my-queue SERVICEBUS_TOPIC_NAME=my-topic SERVICEBUS_SUBSCRIPTION_NAME=my-subscription AZURE_TOKEN_CREDENTIALS=prod # Required only if DefaultAzureCredential is used in production ``` ## Authentication ```typescript import { ServiceBusClient } from "@azure/service-bus"; import { DefaultAzureCredential, ManagedIdentityCredential } from "@azure/identity"; // Local dev: DefaultAzureCredential. Production: set AZURE_TOKEN_CREDENTIALS=prod or AZURE_TOKEN_CREDENTIALS=<specific_credential> const credential = new DefaultAzureCredential({requiredEnvVars: ["AZURE_TOKEN_CREDENTIALS"]}); // Or use a specific credential directly in production: // See https://learn.microsoft.com/javascript/api/overview/azure/identity-readme?view=azure-node-latest#credential-classes // const credential = new ManagedIdentityCredential(); const fullyQualifiedNamespace = process.env.SERVICEBUS_NAMESPACE!; const client = new ServiceBusClient(fullyQualifiedNamespace, credential); ``` ## Core Workflow ### Send Messages to Queue ```typescript const sender = client.createSender("my-queue"); // Single message await sender.sendMessages({ body: { orderId: "12345", amount: 99.99 }, contentType: "application/json", }); // Batch messages const batch = await sender.createMessageBatch(); batch.tryAddMessage({ body: "Message 1" }); batch.tryAddMessage({ body: "Message 2" }); await sender.sendMessages(batch); await sender.close(); ``` ### Receive Messages from Queue ```typescript const receiver = client.createReceiver("my-queue"); // Receive batch const messages = await receiver.receiveMessages(10, { maxWaitTimeInMs: 5000 }); for (const message of messages) { console.log(`Received: ${message.body}`); await receiver.completeMessage(message); } await receiver.close(); ``` ### Subscribe to Messages (Event-Driven) ```typescript const receiver = client.createReceiver("my-queue"); const subscription = receiver.subscribe({ processMessage: async (message) => { console.log(`Processing: ${message.body}`); // Message auto-completed on success }, processError: async (args) => { console.error(`Error: ${args.error}`); }, }); // Stop after some time setTimeout(async () => { await subscription.close(); await receiver.close(); }, 60000); ``` ### Topics and Subscriptions ```typescript // Send to topic const topicSender = client.createSender("my-topic"); await topicSender.sendMessages({ body: { event: "order.created", data: { orderId: "123" } }, applicationProperties: { eventType: "order.created" }, }); // Receive from subscription const subscriptionReceiver = client.createReceiver("my-topic", "my-subscription"); const messages = await subscriptionReceiver.receiveMessages(10); ``` ## Message Sessions ```typescript // Send session message const sender = client.createSender("session-queue"); await sender.sendMessages({ body: { step: 1, data: "First step" }, sessionId: "workflow-123", }); // Receive session messages const sessionReceiver = await client.acceptSession("session-queue", "workflow-123"); const messages = await sessionReceiver.receiveMessages(10); // Get/set session state const state = await sessionReceiver.getSessionState(); await sessionReceiver.setSessionState(Buffer.from(JSON.stringify({ progress: 50 }))); await sessionReceiver.close(); ``` ## Dead-Letter Handling ```typescript // Move to dead-letter await receiver.deadLetterMessage(message, { deadLetterReason: "Validation failed", deadLetterErrorDescription: "Missing required field: orderId", }); // Process dead-letter queue const dlqReceiver = client.createReceiver("my-queue", { subQueueType: "deadLetter" }); const dlqMessages = await dlqReceiver.receiveMessages(10); for (const msg of dlqMessages) { console.log(`DLQ Reason: ${msg.deadLetterReason}`); // Reprocess or log await dlqReceiver.completeMessage(msg); } ``` ## Scheduled Messages ```typescript const sender = client.createSender("my-queue"); // Schedule for future delivery const scheduledTime = new Date(Date.now() + 60000); // 1 minute from now const sequenceNumber = await sender.scheduleMessages( { body: "Delayed message" }, scheduledTime ); // Cancel scheduled message await sender.cancelScheduledMessages(sequenceNumber); ``` ## Message Deferral ```typescript // Defer message for later await receiver.deferMessage(message); // Receive deferred message by sequence number const deferredMessage = await receiver.receiveDeferredMessages(message.sequenceNumber!); await receiver.completeMessage(deferredMessage[0]); ``` ## Peek Messages (Non-Destructive) ```typescript const receiver = client.createReceiver("my-queue"); // Peek without removing const peekedMessages = await receiver.peekMessages(10); for (const msg of peekedMessages) { console.log(`Peeked: ${msg.body}`); } ``` ## Key Types ```typescript import { ServiceBusClient, ServiceBusSender, ServiceBusReceiver, ServiceBusSessionReceiver, ServiceBusMessage, ServiceBusReceivedMessage, ProcessMessageCallback, ProcessErrorCallback, } from "@azure/service-bus"; ``` ## Receive Modes ```typescript // Peek-Lock (default) - message locked until completed/abandoned const receiver = client.createReceiver("my-queue", { receiveMode: "peekLock" }); await receiver.completeMessage(message); // Remove from queue await receiver.abandonMessage(message); // Return to queue await receiver.deferMessage(message); // Defer for later await receiver.deadLetterMessage(message); // Move to DLQ // Receive-and-Delete - message removed immediately const receiver = client.createReceiver("my-queue", { receiveMode: "receiveAndDelete" }); ``` ## Best Practices 1. **Use Microsoft Entra Token Credential** - Use `DefaultAzureCredential` for local development; use `ManagedIdentityCredential` or `WorkloadIdentityCredential` for production 2. **Reuse clients** - Create `ServiceBusClient` once, share across senders/receivers 3. **Close resources** - Always close senders/receivers when done 4. **Handle errors** - Implement `processError` callback for subscription receivers 5. **Use sessions for ordering** - When message order matters within a group 6. **Configure dead-letter** - Always handle DLQ messages 7. **Batch sends** - Use `createMessageBatch()` for multiple messages ## Reference Documentation For detailed patterns, see: - [Queues vs Topics Patterns](references/queues-topics.md) - Queue/topic patterns, sessions, receive modes, message settlement - [Error Handling and Reliability](references/error-handling.md) - ServiceBusError codes, DLQ handling, lock renewal, graceful shutdown
Comments (0)
Sign in to join the conversation.
Reviews (0)
No reviews yet.
No comments yet.