Messaging means services talk to each other by putting messages on a durable broker (Azure Service Bus, Kafka, RabbitMQ) instead of calling each other directly and waiting.
Intro
Messaging means services talk to each other by putting messages on a durable broker (Azure Service Bus, Kafka, RabbitMQ) instead of calling each other directly and waiting. The sender hands off a message and moves on; one or more receivers process it when they are ready. This removes temporal coupling (both sides need not be up at the same time), absorbs traffic spikes, and lets each service scale independently. The price is eventual consistency, duplicate delivery, and harder debugging.
Running example: an insurance claims system where Claims.Api accepts a First Notice of Loss (FNOL), and Fraud, Policy, Payments and Notifications services react to it.
Why we need this
- Availability: if
Claims.Apicalls Fraud, Policy and Notifications synchronously, its availability is roughly the product of theirs (three services at 99.9% each is about 99.7% combined, before retries and latency). - Elasticity: after a storm, claim submissions spike 20x. A queue lets the intake tier accept work at spike rate while workers drain it at a sustainable rate (queue-based load levelling).
- Team autonomy: the publisher does not know who consumes its events; new consumers (e.g. an analytics service) are added without touching the publisher.
- Resilience to deploys: a consumer can be down for a release; messages wait on the broker.
What problem it solves
Problem (from the topic list): synchronous chains degrade availability, throughput and elasticity.
Without messaging, a claim submission is Api -> Fraud -> Policy -> Notifications, all blocking. Symptoms: one slow downstream (Fraud scoring at 8 s) makes the whole submission 8+ s; a Notifications outage returns HTTP 500 to a customer whose claim was actually valid; thread pools exhaust during peaks; retries by clients create duplicate claims.
With messaging the API validates, stores the claim, publishes ClaimSubmitted, and returns 202 Accepted. Downstream work happens asynchronously.
When it is needed (and when it is NOT)
Fits:
- Work that does not need to finish before the user gets a response (notifications, scoring, document generation).
- Fan-out: one event, many independent reactions.
- Spiky or bursty load; long-running work; integration with unreliable or slow partners.
- Cross-service workflows (sagas) that must survive restarts.
Does not fit:
- The caller needs the answer now (e.g. “is this policy active?” during quote screen) - use synchronous RPC (Day 15) or a local replica of the data.
- Small system, one team, one deployable: an in-process call or a background job is simpler.
- Strict read-your-own-writes UX with no tolerance for delay, unless you design for it (optimistic UI, polling on status).
- Team has no capacity to operate/monitor a broker, dead-letter queues and idempotency.
How to identify the problem (key signals)
- p95/p99 latency of an endpoint equals the sum of several downstream calls.
- A non-critical dependency outage (email, SMS) causes user-facing 5xx errors.
- Thread-pool starvation / HTTP client
TaskCanceledExceptionstorms during peaks. - Deployment of one service requires coordinating restarts of others (“deploy order” wiki pages).
- Nightly batch jobs exist purely to push data between services, and are always late.
- Publisher code contains a growing list of
await xClient.NotifyAsync(...)calls, one per new consumer. - Clients retry timed-out requests, creating duplicate records.
Flow Diagram
The API stores the claim, publishes once, and returns 202; consumers work asynchronously.
flowchart LR UI["Angular"] -- "POST /claims" --> API["Claims.Api"] API --> DB[("Claims + outbox")] API -- "202 Accepted" --> UI DB --> T{{"Topic: claim-submitted"}} T --> S1["fraud subscription"] --> FW["Fraud worker"] T --> S2["policy subscription"] --> PW["Policy worker"] T --> S3["notifications subscription"] --> NW["Notification worker"] FW -. "fails 10 times" .-> DLQ[("Dead-letter queue")] UI -- "poll GET /claims/id" --> APILevel 1: Beginner
Analogy: a post office box. You drop a letter and leave; the recipient collects it when they can. Direct calling is a phone call - both must be present.
Vocabulary: producer sends; queue holds messages for one competing consumer group; topic + subscription delivers a copy to every subscriber; consumer processes; acknowledge (complete) removes the message; dead-letter queue (DLQ) holds messages that repeatedly fail.
Minimal .NET example (Azure.Messaging.ServiceBus, .NET 10 LTS):
using Azure.Messaging.ServiceBus;using System.Text.Json;
public record ClaimSubmitted(Guid ClaimId, string PolicyNumber, decimal Amount);
await using var client = new ServiceBusClient("<namespace>.servicebus.windows.net", new Azure.Identity.DefaultAzureCredential());
// Producerawait using var sender = client.CreateSender("claim-submitted");var evt = new ClaimSubmitted(Guid.NewGuid(), "POL-1001", 2500m);await sender.SendMessageAsync(new ServiceBusMessage(JsonSerializer.Serialize(evt)){ MessageId = evt.ClaimId.ToString(), ContentType = "application/json", Subject = nameof(ClaimSubmitted)});
// Consumerawait using var processor = client.CreateProcessor("claim-submitted", new ServiceBusProcessorOptions { MaxConcurrentCalls = 4, AutoCompleteMessages = false });processor.ProcessMessageAsync += async args =>{ var claim = JsonSerializer.Deserialize<ClaimSubmitted>(args.Message.Body)!; Console.WriteLine($"Scoring claim {claim.ClaimId}"); await args.CompleteMessageAsync(args.Message);};processor.ProcessErrorAsync += args => { Console.WriteLine(args.Exception); return Task.CompletedTask; };await processor.StartProcessingAsync();Level 2: Intermediate
Real application shape: Angular -> Claims.Api (ASP.NET Core) -> SQL Server/PostgreSQL + Service Bus -> worker services.
API (returns 202, status via polling):
app.MapPost("/claims", async (SubmitClaim cmd, ClaimsDb db, ServiceBusSender sender, CancellationToken ct) =>{ var claim = new Claim(Guid.NewGuid(), cmd.PolicyNumber, cmd.Amount, ClaimStatus.Received); db.Claims.Add(claim); await db.SaveChangesAsync(ct); // NOTE: dual write - see Day 12 (Transactional Outbox) for the safe version. await sender.SendMessageAsync(new ServiceBusMessage(JsonSerializer.SerializeToUtf8Bytes( new ClaimSubmitted(claim.Id, claim.PolicyNumber, claim.Amount))) { MessageId = claim.Id.ToString() }, ct); return Results.Accepted($"/claims/{claim.Id}", new { claim.Id, status = "Received" });});
app.MapGet("/claims/{id:guid}", async (Guid id, ClaimsDb db) => await db.Claims.FindAsync(id) is { } c ? Results.Ok(new { c.Id, c.Status }) : Results.NotFound());Worker (BackgroundService registered via DI):
builder.Services.AddAzureClients(b =>{ b.AddServiceBusClientWithNamespace("<ns>.servicebus.windows.net"); b.AddClient<ServiceBusProcessor, ServiceBusClientOptions>((_, _, sp) => sp.GetRequiredService<ServiceBusClient>().CreateProcessor("claim-submitted", "fraud", new ServiceBusProcessorOptions { MaxConcurrentCalls = 8, AutoCompleteMessages = false, PrefetchCount = 16 }));});builder.Services.AddHostedService<FraudWorker>();Angular (current major) polling the status:
submit(cmd: SubmitClaim) { return this.http.post<{ id: string }>('/api/claims', cmd).pipe( switchMap(r => timer(0, 2000).pipe( switchMap(() => this.http.get<{ status: string }>(`/api/claims/${r.id}`)), takeWhile(s => s.status === 'Received' || s.status === 'Scoring', true) )) );}Alternatives to polling: SignalR / Azure SignalR Service push, or Web PubSub.
Level 3: Advanced
Delivery semantics. Brokers give at-least-once delivery in practice. Design consumers to be idempotent (Day 17). Exactly-once end to end is achieved by idempotent processing, not by the broker alone.
Ordering. Ordering is only guaranteed within a unit: Service Bus sessions (SessionId = ClaimId) or Kafka partition key. Global ordering kills throughput. Choose the key by aggregate.
Poison messages. After MaxDeliveryCount (default 10 on Service Bus) the message goes to the DLQ. Monitor DLQ depth, alert on it, and build a replay tool. Do not blindly retry forever.
Backpressure and concurrency. MaxConcurrentCalls x instances must not exceed what the database or downstream API can handle. Tune PrefetchCount carefully - prefetched messages hold their lock and can expire.
Lock duration. Peek-lock default is 1 minute (max 5 minutes on Service Bus); the processor auto-renews locks (MaxAutoLockRenewalDuration). Long jobs should checkpoint or hand off to a durable workflow.
Security. Use Microsoft Entra ID + managed identity (Azure Service Bus Data Sender/Receiver roles) instead of connection strings; use private endpoints (Premium). Never put PII you do not need into messages; consider claim-check pattern for large payloads/documents (store in Blob, send the URI).
Schema evolution. Version events (ClaimSubmitted.v2), add fields only, tolerate unknown fields, and validate with contract tests (Day 51).
Common mistakes:
- Dual write (DB save + publish) without an outbox.
- Non-idempotent consumers.
- Using a queue where a topic is needed (or the reverse).
- Treating messages as RPC (request/reply with tight timeouts everywhere).
- No DLQ monitoring; giant messages; no correlation IDs.
- Auto-completing messages before processing succeeds.
Level 4: Expert and Architect view
| Option | Best for | Trade-offs |
|---|---|---|
| Azure Service Bus | Enterprise messaging: queues, topics, sessions, DLQ, scheduled delivery, transactions | Per-message model, not a replayable log; throughput bounded by tier/units |
| Azure Event Hubs (Kafka-compatible endpoint) | High-volume telemetry/event streams, replay, partition-ordered processing | Consumer manages offsets/checkpoints; no per-message DLQ, no sessions |
| Azure Event Grid | Reactive event notification (resource events, webhooks) | Lightweight push, not for heavy work queues |
| Azure Storage Queues | Very cheap simple queues | Few features, no topics/ordering guarantees |
| RabbitMQ / Kafka self-managed | Portability, existing ops skills | You operate it |
| Synchronous RPC | Immediate answer needed | Temporal coupling, cascading failure |
Combines with: Transactional Outbox (12), Idempotent Consumer (17), Saga (7), Domain Event (11), CQRS (8), Circuit Breaker/Retry for the sync edges, Distributed Tracing (31) via W3C traceparent in message properties.
ADR (example)
- Title: ADR-016 Use asynchronous messaging for claim intake downstream processing
- Status: Proposed
- Context: Claim submission currently makes three blocking calls (Fraud, Policy, Notifications); p95 is 6.8 s and Notifications outages cause failed submissions. Peak season load is ~15x normal.
- Decision:
Claims.Apipersists the claim and publishesClaimSubmittedto an Azure Service Bus topic via a transactional outbox; Fraud, Policy and Notifications each consume through their own subscription. API returns 202 with a status resource. - Consequences: (+) intake latency drops to DB write time; outages isolated; new consumers need no publisher change. (-) eventual consistency, need idempotent consumers, DLQ operations, and UI changes for async status. Requires on-call runbooks for DLQ.
- Alternatives considered: keep sync with Polly retries/circuit breakers (does not fix load spikes); Event Hubs (replay not required, per-message DLQ needed).
Azure implementation
Services: Azure Service Bus (primary), Event Hubs (streams), Event Grid (notifications), Azure Functions Service Bus trigger or AKS/Container Apps workers (consumers), Application Insights + Azure Monitor (observability), Key Vault (only if secrets are unavoidable), Entra ID managed identities.
Configuration (Azure CLI, sketch):
az servicebus namespace create -g rg-claims -n sb-claims-prod --sku Standard --location westeuropeaz servicebus topic create -g rg-claims --namespace-name sb-claims-prod -n claim-submitted \ --enable-duplicate-detection true --duplicate-detection-history-time-window PT10Maz servicebus topic subscription create -g rg-claims --namespace-name sb-claims-prod \ --topic-name claim-submitted -n fraud --max-delivery-count 10 --lock-duration PT1Maz role assignment create --assignee <fraud-worker-managed-identity> \ --role "Azure Service Bus Data Receiver" --scope <subscription-resource-id>Tiers (verify current numbers on the Azure pricing page before budgeting):
- Basic: queues only, no topics.
- Standard: queues + topics, sessions, duplicate detection, pay per operation plus a base charge; shared infrastructure, throttling possible.
- Premium: dedicated messaging units, predictable latency, VNet/private endpoints, larger messages (up to 100 MB), geo-disaster recovery; priced per messaging unit-hour. Choose Premium for production workloads needing network isolation or steady high throughput.
Reference architecture (text): Angular SPA (Static Web Apps) -> Azure Front Door/API Management -> Claims.Api on Container Apps -> Azure SQL / PostgreSQL Flexible Server (claims + outbox table) -> outbox relay publishes to Service Bus topic claim-submitted -> subscriptions fraud, policy, notifications -> each worker (Container Apps with KEDA scaling on subscription message count, or Functions) -> DLQ alerts in Azure Monitor -> traces and logs in Application Insights with correlation ID from the message.
Teaching guide for my team
2-minute beginner explanation: “Instead of phoning the fraud team and waiting on hold, we put the claim in their inbox and tell the customer we received it. The inbox (broker) keeps it safe until they read it. If they are busy or off sick, nothing is lost.” Show the sender/processor code from Level 1.
5-minute intermediate explanation: Sync chain vs async flow diagram; queue vs topic; peek-lock and complete; why messages can arrive twice; why the API returns 202 and the UI polls; DLQ as the safety net; and the dual-write problem that leads to the outbox.
Hands-on exercise: Create a Service Bus namespace (or use the local emulator), a topic claim-submitted with two subscriptions. Build Claims.Api publishing on POST and two console consumers. Kill one consumer, submit 10 claims, restart it. Expected: the other consumer processed immediately, the restarted one catches up on all 10 with no loss. Then throw an exception for one claim and observe it land in the DLQ after 10 attempts.
Interview questions:
- Queue vs topic? A queue delivers each message to one consumer (competing consumers); a topic copies each message to every subscription.
- Why must consumers be idempotent? Brokers deliver at least once; lock expiry, retries or crashes after processing but before completing cause redelivery.
- How do you preserve order? Only per key: use sessions (Service Bus) or partition keys (Kafka/Event Hubs), keyed by aggregate ID.
Mastery checklist
- Can explain queue vs topic vs stream and pick between Service Bus, Event Hubs and Event Grid.
- Implements a producer and processor with peek-lock, manual completion and error handler.
- Designs idempotent consumers and DLQ handling with replay.
- Explains and avoids the dual-write problem.
- Chooses ordering key and understands its throughput impact.
- Uses managed identity and private networking rather than connection strings.
- Propagates correlation/trace IDs through messages.
- Writes an ADR justifying async vs sync for a given flow.
Key takeaway
Use messaging when the caller does not need the answer now: you trade immediate consistency for availability and scalability, and you pay for it with idempotency, DLQ operations and asynchronous UX.
