Manikandan — Manikandan
Microservices

Day 17: Idempotent Consumer

ManikandanManikandan
15 min read·Updated Aug 26, 2022

An idempotent consumer is a message handler that produces the same end result whether a message is processed once or several times.

Intro

An idempotent consumer is a message handler that produces the same end result whether a message is processed once or several times. Brokers such as Azure Service Bus deliver at-least-once, so duplicates are normal, not exceptional. In our insurance claims system, the Payments service must pay a claim exactly once even if the ClaimApproved message is delivered three times; the consumer achieves this with a message ID, a deduplication store, or an operation that is naturally safe to repeat.

Why we need this

  • Business reason: Paying a claim twice, sending two “claim approved” letters, or creating two reserve entries is real money and real regulatory risk. Insurers reconcile ledgers; duplicates surface as audit findings and customer complaints.
  • Technical reason: Reliable brokers choose “never lose a message” over “never repeat a message”. Lock expiry, consumer crashes after doing the work but before acknowledging, network timeouts on the acknowledgement, producer retries and broker failover all cause redelivery.
  • Exactly-once is an end-to-end property: it comes from at-least-once delivery plus idempotent processing. The consumer is the only place that can guarantee it.
  • Enables safe retries: if handlers are idempotent, you can retry aggressively, replay a dead-letter queue, and re-drive a failed batch without fear.

What problem it solves

Problem (from the topic list): At-least-once delivery leads to duplicate processing (for example, double charge).

Failure timeline without idempotency:

  1. Payments receives ClaimApproved(claim 42, 2,500).
  2. It calls the bank API and pays 2,500. Success.
  3. The process crashes (or the Service Bus lock expires) before CompleteMessageAsync.
  4. The broker redelivers the message. Payments pays 2,500 again.

Result: the policyholder is paid 5,000. No exception was ever logged; every component behaved “correctly”. Other symptoms: duplicate emails, inflated counters and reports, double inventory-style reservations, and repeated audit entries.

When it is needed (and when it is NOT)

Needed

  • Any consumer of a broker with at-least-once delivery (Service Bus PeekLock, Kafka with manual commit, RabbitMQ with manual ack) whose handler has side effects: money movement, state transitions with counters, email/SMS, calls to external APIs, inserts without natural unique keys.
  • HTTP endpoints that clients or gateways may retry (POST /payments with an Idempotency-Key header) - the same idea.
  • Saga participants and Outbox consumers (Days 7 and 12): retries are part of their design.

Not needed / overkill

  • Handlers that are naturally idempotent: UPDATE Claims SET Status='Approved' WHERE Id=@id, upserts keyed by business ID, cache invalidation, search index refresh from the latest state.
  • Read-only handlers or pure computations with no side effects.
  • Loss-tolerant telemetry where an occasional duplicate is harmless and cost of dedup exceeds value (for example, page-view counters).
  • A dedup store is not a substitute for the producer-side Outbox: it protects against duplicates, not against lost messages.

How to identify the problem (key signals)

  1. Customers report two payments, two emails, or two SMS for one claim, usually clustered around a deployment, restart or broker incident.
  2. Business tables show duplicate rows for the same business key (two Payment rows with the same ClaimId).
  3. Service Bus metrics show a high delivery count, and messages with DeliveryCount > 1 in logs are not followed by any deduplication check.
  4. Ledger reconciliation totals differ from claim totals by an exact multiple of a known payment.
  5. Handler code performs “insert” or “call external API” with no guard, and completion happens after the side effect.
  6. Load tests that kill the consumer mid-run leave inconsistent results.
  7. A consumer cannot safely replay the dead-letter queue; the team is afraid to click “resubmit”.

Flow Diagram

The dedup record and the business change commit atomically, so redelivery is a no-op.

flowchart TD
M(["ClaimApproved message"]) --> TX["Begin transaction"]
TX --> INS["Insert ProcessedMessage - MessageId, Consumer"]
INS --> PAY["Insert Payment"]
PAY --> SAVE{"Commit succeeded?"}
SAVE -- "yes" --> DONE["Complete message"]
SAVE -- "duplicate key" --> SKIP["Log duplicate, treat as success"]
SKIP --> DONE
DONE -. "crash before complete" .-> M

Level 1: Beginner

Analogy: A lift button. Pressing “3” five times doesn’t take you to floor 15; the lift remembers “go to 3”. A bad handler is a vending machine that dispenses a snack on every press of the same paid button.

Three ways to be idempotent

  1. Natural idempotency: design the operation as “set state to X” rather than “add Y”.
  2. Message ID + dedup store: remember IDs of processed messages; skip if already seen.
  3. Idempotency key to a downstream system: pass a stable key so the other side dedups.

Minimal example (in-memory, teaching only; production needs a durable store)

public record ClaimApproved(Guid MessageId, Guid ClaimId, decimal Amount);
public class PaymentHandler
{
private readonly HashSet<Guid> _processed = new(); // NOT for production: lost on restart
public void Handle(ClaimApproved msg)
{
if (!_processed.Add(msg.MessageId))
{
Console.WriteLine($"Duplicate {msg.MessageId} ignored");
return;
}
Console.WriteLine($"Paying {msg.Amount} for claim {msg.ClaimId}");
}
}

Running Handle twice with the same MessageId pays once. The flaw: memory disappears on restart, and marking and paying are not atomic. Level 2 fixes both.

Level 2: Intermediate

Scenario: Payments service (ASP.NET Core .NET 10, EF Core, SQL Server) consumes ClaimApproved from Service Bus and records a payment. The dedup record and the business change are saved in the same database transaction, so either both happen or neither.

Schema and entities

using Microsoft.EntityFrameworkCore;
public class ProcessedMessage
{
public string MessageId { get; set; } = ""; // Service Bus MessageId
public string Consumer { get; set; } = ""; // handler name, lets several handlers share the table
public DateTime ProcessedAtUtc { get; set; }
}
public class Payment
{
public Guid Id { get; set; }
public Guid ClaimId { get; set; }
public decimal Amount { get; set; }
public DateTime CreatedAtUtc { get; set; }
}
public class PaymentsDb(DbContextOptions<PaymentsDb> options) : DbContext(options)
{
public DbSet<ProcessedMessage> ProcessedMessages => Set<ProcessedMessage>();
public DbSet<Payment> Payments => Set<Payment>();
protected override void OnModelCreating(ModelBuilder b)
{
b.Entity<ProcessedMessage>().HasKey(x => new { x.MessageId, x.Consumer }); // unique constraint = the guard
b.Entity<Payment>().HasIndex(x => x.ClaimId).IsUnique(); // second line of defence
}
}

Handler with atomic dedup

using Azure.Messaging.ServiceBus;
using Microsoft.EntityFrameworkCore;
using System.Text.Json;
public record ClaimApproved(Guid ClaimId, decimal Amount);
public class ClaimApprovedHandler(PaymentsDb db, ILogger<ClaimApprovedHandler> log)
{
private const string Name = nameof(ClaimApprovedHandler);
public async Task HandleAsync(ProcessMessageEventArgs args)
{
var messageId = args.Message.MessageId;
var evt = JsonSerializer.Deserialize<ClaimApproved>(args.Message.Body)!;
await using var tx = await db.Database.BeginTransactionAsync(args.CancellationToken);
try
{
db.ProcessedMessages.Add(new ProcessedMessage
{
MessageId = messageId, Consumer = Name, ProcessedAtUtc = DateTime.UtcNow
});
db.Payments.Add(new Payment
{
Id = Guid.NewGuid(), ClaimId = evt.ClaimId,
Amount = evt.Amount, CreatedAtUtc = DateTime.UtcNow
});
await db.SaveChangesAsync(args.CancellationToken); // throws on duplicate key
await tx.CommitAsync(args.CancellationToken);
}
catch (DbUpdateException ex) when (IsDuplicateKey(ex))
{
log.LogInformation("Duplicate message {MessageId} skipped", messageId);
// fall through: it is safe to complete because the work was already done
}
await args.CompleteMessageAsync(args.Message, args.CancellationToken);
}
// SQL Server: 2627 = unique constraint violation, 2601 = unique index violation
private static bool IsDuplicateKey(DbUpdateException ex) =>
ex.InnerException is Microsoft.Data.SqlClient.SqlException { Number: 2627 or 2601 };
}

Key points:

  • The ProcessedMessage insert and the Payment insert commit together. A crash before commit leaves neither; after commit, a redelivery hits the primary key and is ignored.
  • CompleteMessageAsync happens after commit. If completion fails, redelivery is harmless.
  • The unique index on Payment.ClaimId is a second guard using a business key, useful when different messages (different MessageIds) could describe the same business fact.

When the side effect is an external call (bank API, email): a database transaction cannot cover it. Pass an idempotency key so the provider dedups:

using var req = new HttpRequestMessage(HttpMethod.Post, "https://bank.example/payments")
{
Content = JsonContent.Create(new { claimId = evt.ClaimId, amount = evt.Amount })
};
req.Headers.Add("Idempotency-Key", $"claim-payment-{evt.ClaimId}"); // stable across retries
var resp = await http.SendAsync(req, ct);
resp.EnsureSuccessStatusCode();

Angular relevance: the same principle protects the UI’s “Pay now” button. The client generates an idempotency key once per user intent and sends it on every retry.

import { HttpClient, HttpHeaders } from '@angular/common/http';
import { inject, Injectable } from '@angular/core';
@Injectable({ providedIn: 'root' })
export class PaymentApi {
private http = inject(HttpClient);
payClaim(claimId: string, amount: number, idempotencyKey: string) {
return this.http.post('/api/payments', { claimId, amount }, {
headers: new HttpHeaders({ 'Idempotency-Key': idempotencyKey }),
});
}
}
// In the component: create `crypto.randomUUID()` when the form is opened, reuse it on retry/double-click.

Level 3: Advanced

Choosing the dedup key

  • Use a stable ID assigned by the producer (for example the Outbox row ID or event ID), not one generated at receive time. Service Bus MessageId is preserved on redelivery; the SequenceNumber is assigned by the broker and is not stable across a producer re-send.
  • For business-level idempotency, key on the business identity (ClaimId + payment purpose), not the message ID.

Race conditions

  • Two consumer instances may receive the same message concurrently (lock expiry while the first is still running). The unique constraint arbitrates: one commit wins, the other gets a duplicate-key error. Never rely on “check then insert” alone (if (!exists) insert), which is racy.
  • Lock renewal: keep MaxAutoLockRenewalDuration above worst-case processing time to reduce concurrent duplicates.

In-progress vs completed states

For long or external side effects, store status: Started / Completed. On redelivery: Completed means skip; Started means resume or verify with the downstream system before repeating. A plain “processed” flag written at the end can double-run; written at the start can lose work on crash.

Storage options and trade-offs

OptionProsCons
Relational table, same transaction as business dataAtomic, simple, exactly-once for DB effectsTable grows; needs cleanup; DB-only effects
Redis SET NX with TTLFast, auto-expiryNot atomic with your DB; duplicates possible after TTL or Redis loss
Cosmos DB with unique key / TTLScales, TTL cleanupAtomicity only within a partition/transactional batch
Service Bus duplicate detectionNo codeOnly dedups sends within a window (max 7 days), does not protect against redelivery
Natural idempotency (upsert, set state)No extra storageNot always possible (money movement, emails)

Retention and cleanup

Delete ProcessedMessage rows older than the maximum redelivery horizon (max TTL of the message plus DLQ replay window, for example 14-30 days) using a scheduled job. Deleting too early reopens the duplicate window; index by ProcessedAtUtc for cheap purges.

Security

Validate that message IDs and payloads come from trusted publishers (RBAC on send). Do not trust client-supplied idempotency keys across tenants: scope keys by tenant/user (tenantId + key) to avoid one caller replaying or blocking another’s key.

Failure modes and common mistakes

  1. Acknowledging before the work is durable (loss instead of duplicates).
  2. Dedup write and business write in separate transactions (a crash between them causes either loss or duplicates).
  3. Using a random GUID generated on receive as the “message ID”.
  4. Dedup in memory only.
  5. Dedup table with no cleanup; a multi-hundred-million-row table hurts inserts.
  6. Returning an error on duplicates, which makes the broker retry forever and fills the DLQ. Duplicate means success: complete the message.
  7. Forgetting that side effects to third parties need their own idempotency keys.
  8. Same key, different payload: return a conflict (HTTP 409/422) rather than silently accepting.

Level 4: Expert and Architect view

Approaches compared

ApproachGuaranteeComplexityBest for
Natural idempotency (state-based commands)Strong, no storageLowState transitions, upserts
Inbox table (dedup in same DB transaction)Exactly-once effect on local DBMediumMoney, ledgers, entity creation
External idempotency key (provider-side)Depends on providerLow-mediumPayment gateways, email APIs
Redis SET NXBest effortLowThrowaway, high-volume, loss-tolerant work
Broker-side duplicate detectionSend-side onlyVery lowCutting producer retries; complements, never replaces, the consumer guard
Kafka transactions / exactly-once semanticsWithin Kafka read-process-writeHighKafka-to-Kafka pipelines only, not external side effects

Combines with

  • Transactional Outbox (Day 12): producer side guarantees the event is published at least once; the consumer’s inbox guarantees it is applied once. Together they give effectively-once processing.
  • Messaging (Day 16): the delivery model that makes this necessary.
  • Saga (Day 7): each step and each compensation must be idempotent.
  • Retry & Backoff, Circuit Breaker (Days 26 and 28): retries are only safe on idempotent operations.
  • Event Sourcing / CQRS: projections are rebuilt by replay, so handlers must tolerate reprocessing; store the last processed event position for a projection.

ADR (architecture review)

  • Title: ADR-017 All message consumers with side effects must be idempotent using an Inbox table.
  • Context: Service Bus delivers at-least-once. A duplicate ClaimApproved could cause a duplicate payout. Message loss is unacceptable, so at-most-once is rejected.
  • Decision: Producers set a stable MessageId (Outbox event ID). Consumers write a ProcessedMessage row in the same transaction as their business changes; a primary key on (MessageId, Consumer) enforces the guard. External calls carry a deterministic idempotency key. Duplicate detection at the broker is enabled as an optimization only. Cleanup job purges rows older than 30 days.
  • Consequences (positive): safe redelivery, DLQ replay, and retries; audit-friendly; simple to reason about.
  • Consequences (negative): extra table and write per message; cleanup job; a shared library is needed to keep the implementation uniform; external effects still rely on provider support.
  • Alternatives considered: Redis-only dedup (rejected: not atomic with DB, TTL gaps); relying on broker duplicate detection (rejected: does not cover redelivery); at-most-once delivery (rejected: loses claims).
  • Status: Accepted.

Azure implementation

Services

  • Azure Service Bus: PeekLock receive mode, MaxDeliveryCount, dead-letter queue, optional duplicate detection (window up to 7 days, keyed by MessageId) on queues/topics, sessions if per-claim ordering is required.
  • Azure SQL Database or Azure Database for PostgreSQL Flexible Server: hosts the inbox table beside business data (same transaction).
  • Azure Cache for Redis: optional fast pre-check for high-volume, loss-tolerant flows; never the only guard for money.
  • Azure Cosmos DB: alternative when the consumer’s data is already there; use a container with unique key policy and TTL.
  • Azure Functions with Service Bus trigger / Container Apps / AKS: consumer hosting. Functions retries messages on failure, so functions need the same idempotency handling.
  • Azure API Management: for HTTP idempotency keys, you can cache responses per key with policies, but the durable guarantee must live in the service.
  • Application Insights / Azure Monitor: track duplicates-skipped counter, DLQ depth, delivery count distribution; alert on spikes.

Configuration

Terminal window
# Enable broker-side duplicate detection (optimization, not a substitute)
az servicebus queue create -g rg-claims --namespace-name claims-demo -n claim-approved \
--enable-duplicate-detection true --duplicate-detection-history-time-window PT10M \
--max-delivery-count 10 --lock-duration PT1M
-- Purge job (run daily via Azure SQL Elastic Jobs, Azure Functions timer or Logic Apps)
DELETE TOP (5000) FROM ProcessedMessages WHERE ProcessedAtUtc < DATEADD(day, -30, SYSUTCDATETIME());

Note that duplicate detection cannot be changed after the queue or topic is created, and it is available on Standard and Premium tiers, not Basic; verify current limits on the Service Bus documentation and pricing pages before finalizing. Standard bills per operation plus a base charge; Premium bills per messaging unit per hour. Dedup itself adds no broker cost, but the inbox table adds database storage and write IOPS (size Azure SQL DTUs/vCores or Postgres storage for one extra insert per message).

Reference architecture (text)

Claims API writes claim state and an Outbox row in one Azure SQL transaction; the Outbox publisher sends ClaimApproved to a Service Bus topic with MessageId equal to the outbox event ID. The Payments consumer on Azure Container Apps (KEDA scaling on queue length) receives via PeekLock, opens a transaction against its own Azure SQL database, inserts an inbox row and the payment together, commits, then completes the message. The bank call uses a deterministic idempotency key. Poison messages go to the DLQ; an operator tool can safely resubmit them because handlers are idempotent. Application Insights tracks duplicates skipped and DLQ depth with alerts through Azure Monitor. Managed identities and RBAC secure Service Bus and SQL access.

Teaching guide for my team

2-minute beginner explanation

“Our message system promises never to lose a message, and the price is that it sometimes delivers one twice. If our payment handler just does what the message says each time, we pay twice. So the handler keeps a notebook: ‘I already handled message 123.’ Before doing anything it checks the notebook, and it writes in the notebook and does the work at the same instant, so a crash can’t leave one without the other. If the notebook says done, we simply say ‘OK, already done’ and move on.”

5-minute intermediate explanation

Cover: why duplicates occur (crash before ack, lock expiry, producer retry, failover); at-least-once plus idempotency equals effectively-once; the three techniques (natural, inbox table, external key); why the check-then-insert pattern races and unique constraints fix it; same-transaction requirement; duplicate means complete, not fail; cleanup retention; business-key guard as a second line of defence; HTTP Idempotency-Key from Angular. Whiteboard the crash timeline from Section 2 and show where the inbox row blocks the second payment.

Hands-on exercise

Task: Take the Level 2 handler. (1) Send the same ClaimApproved message five times (same MessageId) and confirm only one Payment row exists. (2) Add a temporary throw after SaveChangesAsync but before CompleteMessageAsync so the message is redelivered, and verify the second delivery logs “Duplicate skipped” and completes. (3) Remove the transaction and the unique key, repeat, and observe the duplicate payment. Expected outcome: With the guard: exactly one payment, no DLQ entries, duplicate logs visible. Without the guard: multiple payments, proving the guard is what provides safety.

Interview-style questions

  1. Why can’t the broker just guarantee exactly-once? The consumer can crash after doing its side effect but before acknowledging, and the broker cannot know the effect happened; so brokers choose at-least-once, and exactly-once effects come from idempotent consumers.
  2. Why must the dedup record and the business change be in the same transaction? Otherwise a crash between them either loses the work (marked processed but not done) or repeats it (done but not marked).
  3. What should a consumer do when it detects a duplicate? Treat it as success: skip the work and complete the message. Throwing causes endless redelivery and DLQ noise.

Mastery checklist

  • I can explain exactly which failures cause redelivery in Service Bus (crash, lock expiry, abandon, timeouts).
  • I can implement an inbox table with a composite primary key and atomic commit with business data.
  • I can explain why check-then-insert is racy and how unique constraints solve it.
  • I can choose between natural idempotency, inbox, external key, and Redis for a given side effect.
  • I can pass and scope idempotency keys for HTTP calls and Angular clients.
  • I can design retention and cleanup for the dedup store.
  • I can explain why broker-side duplicate detection is not enough.
  • I can test idempotency by forcing crashes and duplicate deliveries.

Key takeaway

Assume every message will arrive more than once: record “already handled” in the same transaction as the business change, and treat a duplicate as a successful no-op.

Interactive Architectural Roadmaps

Explore Complete Roadmaps & Pattern Checklists

Track your learning with interactive checklists for all 23 Gang of Four patterns and modern Microservice architecture patterns.

Share:
Back to Blog

Related Posts

View All Posts
Microservices

Day 18: Domain-specific Protocol

A domain-specific protocol means choosing a wire protocol that fits the *shape of the traffic* instead of defaulting to HTTP/JSON for everything.

Manikandan
Manikandan·22 min read
Microservices

Day 16: Messaging

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.

Manikandan
Manikandan·10 min read
Microservices

Day 15: Remote Procedure Invocation (RPI)

Remote Procedure Invocation (RPI) is the simplest way for one service to use another: the caller sends a request over the network (REST/HTTP, gRPC, or GraphQL), waits, and gets a response, as if it had called a local...

Manikandan
Manikandan·12 min read