Manikandan — Manikandan
Microservices

Day 12: Transactional Outbox

ManikandanManikandan
16 min read·Updated Aug 21, 2022

The Transactional Outbox pattern solves the "dual-write" problem: a service must update its database and publish an event, but a database and a message broker cannot share one ACID transaction.

Intro

The Transactional Outbox pattern solves the “dual-write” problem: a service must update its database and publish an event, but a database and a message broker cannot share one ACID transaction. Instead of publishing directly, the service writes the event into an Outbox table in the same local database transaction as the business change. A separate relay process later reads the outbox and publishes the messages to the broker. Either both the business change and the event are saved, or neither is. Delivery to the broker becomes at-least-once, so consumers must be idempotent.

Why we need this

  • Microservices communicate through events (ClaimSubmitted, ClaimApproved). Other services (Payments, Notifications, Fraud) depend on those events to do their work.
  • A relational database (SQL Server, PostgreSQL) gives us atomic transactions, but only for its own tables. A broker (Azure Service Bus, Kafka, RabbitMQ) is a different system with its own commit. There is no cheap, reliable transaction spanning both (2PC/MSDTC is not supported by Azure Service Bus in a practical, cloud-native way, and it hurts availability).
  • Business need: “if a claim is saved, the world must eventually hear about it” and “if the world heard about it, the claim must really be saved.” Losing an event means a claim is approved but never paid. Publishing a phantom event means a customer is told about a claim that does not exist.
  • Technical need: turn an unreliable two-system operation into a reliable one-system operation (local ACID) plus a retryable background delivery.

What problem it solves

Problem (from the topic list): Dual-write of DB update plus event publish risks partial failure.

Consider the Claims service’s SubmitClaim handler without an outbox:

// PSEUDO-CODE showing the broken dual-write
await db.SaveChangesAsync(); // step 1: claim committed
await bus.PublishAsync(new ClaimSubmitted(id)); // step 2: publish

What goes wrong:

Failure pointResult
Crash / network error after step 1, before step 2Claim exists, event lost. Payments never starts. Silent data inconsistency.
Reverse order (publish, then save) and save failsEvent published for a claim that does not exist (phantom event).
Broker temporarily downUser gets HTTP 500 even though the DB is fine, or you swallow the error and lose the event.
Retry on the HTTP level after partial successDuplicate claims or duplicate events.

No try/catch fixes this, because the process itself can die between the two statements (deploy, OOM, pod eviction, power loss).

When it is needed (and when it is NOT)

Needed when:

  • A state change must be followed by an event, and losing the event causes business harm (payments, claim status, audit, inventory).
  • You use a database per service and integrate via async messaging.
  • You use Saga orchestration or choreography: every saga step emits events/commands and must not lose them.
  • You publish Domain Events (Day 11) to other bounded contexts.

Not needed / overkill when:

  • The event is purely informational and losing it is acceptable (for example a best-effort “page viewed” metric). Fire-and-forget is fine.
  • The service has no database state change tied to the event (for example a pure gateway).
  • The whole workload is one monolith with one database and in-process handlers; use a normal transaction.
  • The broker is the only store (for example a Kafka Streams processor with transactional producers/exactly-once inside Kafka).
  • You need synchronous, immediate consistency between two systems; outbox gives eventual consistency and adds latency (typically milliseconds to seconds).

How to identify the problem (key signals)

  1. Support tickets like “the claim shows Approved but the customer never got the payment/email”, with no error in logs.
  2. Reconciliation jobs regularly find rows in Service A with no matching effect in Service B.
  3. Code smell: SaveChanges() followed by Publish()/Send() in the same method, with no compensating logic.
  4. Errors such as ServiceBusException/TimeoutException on publish that occur after the database commit succeeded.
  5. Events observed for entities that do not exist (phantom events) after failed transactions or rolled-back deployments.
  6. Message counts in the broker are lower than the number of state changes in the DB for the same period (compare row counts to message metrics).
  7. Incidents cluster around deployments, restarts, or broker maintenance windows.

Flow Diagram

The business change and the event are committed together; a relay publishes later.

flowchart LR
REQ["Submit claim request"] --> TX["One local DB transaction"]
TX --> C[("Claims table")]
TX --> O[("OutboxMessages table")]
O --> RELAY["Outbox relay - background"]
RELAY -- "MessageId = outbox Id" --> SB{{"Service Bus topic"}}
RELAY -- "mark ProcessedUtc" --> O
SB --> CON["Idempotent consumers"]

Level 1: Beginner

Analogy: A shop clerk writes the sale and a “please mail the receipt” note on the same page of the same ledger, then signs the page once. Later a mail clerk walks through the ledger, mails each pending receipt, and ticks it off. If the mail clerk faints, the note is still in the ledger; nothing is lost.

Minimal working example (SQL + C#, console-level):

CREATE TABLE Claims (
Id UNIQUEIDENTIFIER PRIMARY KEY,
PolicyNumber NVARCHAR(30) NOT NULL,
Amount DECIMAL(18,2) NOT NULL,
Status NVARCHAR(20) NOT NULL
);
CREATE TABLE OutboxMessages (
Id UNIQUEIDENTIFIER PRIMARY KEY,
Type NVARCHAR(200) NOT NULL,
Payload NVARCHAR(MAX) NOT NULL,
OccurredUtc DATETIME2 NOT NULL,
ProcessedUtc DATETIME2 NULL
);
using System.Data.SqlClient; // Microsoft.Data.SqlClient in new projects
using System.Text.Json;
async Task SubmitClaimAsync(SqlConnection conn, Guid id, string policy, decimal amount)
{
await using var tx = (SqlTransaction)await conn.BeginTransactionAsync();
await using (var cmd = new SqlCommand(
"INSERT INTO Claims (Id, PolicyNumber, Amount, Status) VALUES (@i,@p,@a,'Submitted')", conn, tx))
{
cmd.Parameters.AddWithValue("@i", id);
cmd.Parameters.AddWithValue("@p", policy);
cmd.Parameters.AddWithValue("@a", amount);
await cmd.ExecuteNonQueryAsync();
}
var evt = new { ClaimId = id, PolicyNumber = policy, Amount = amount };
await using (var cmd = new SqlCommand(
"INSERT INTO OutboxMessages (Id, Type, Payload, OccurredUtc) VALUES (@i,@t,@p,SYSUTCDATETIME())", conn, tx))
{
cmd.Parameters.AddWithValue("@i", Guid.NewGuid());
cmd.Parameters.AddWithValue("@t", "ClaimSubmitted");
cmd.Parameters.AddWithValue("@p", JsonSerializer.Serialize(evt));
await cmd.ExecuteNonQueryAsync();
}
await tx.CommitAsync(); // both rows or neither
}

Key idea: one transaction, two inserts. Publishing is someone else’s job (Days 13 and 14 cover the two ways to relay it).

Level 2: Intermediate

Real .NET (current LTS, .NET 10) + EF Core + Angular + SQL Server/PostgreSQL.

Design: the outbox rows are written by SaveChangesAsync through an EF Core interceptor that converts domain events collected on aggregates into outbox rows. The same DbContext transaction commits both.

public sealed class OutboxMessage
{
public Guid Id { get; init; } = Guid.NewGuid();
public string Type { get; init; } = default!;
public string Payload { get; init; } = default!;
public DateTime OccurredUtc { get; init; } = DateTime.UtcNow;
public DateTime? ProcessedUtc { get; set; }
public int Attempts { get; set; }
public string? Error { get; set; }
}
public interface IDomainEvent { }
public sealed record ClaimSubmitted(Guid ClaimId, string PolicyNumber, decimal Amount) : IDomainEvent;
public abstract class AggregateRoot
{
private readonly List<IDomainEvent> _events = new();
public IReadOnlyList<IDomainEvent> DomainEvents => _events;
protected void Raise(IDomainEvent e) => _events.Add(e);
public void ClearEvents() => _events.Clear();
}
public sealed class Claim : AggregateRoot
{
public Guid Id { get; private set; }
public string PolicyNumber { get; private set; } = default!;
public decimal Amount { get; private set; }
public string Status { get; private set; } = default!;
private Claim() { }
public static Claim Submit(string policy, decimal amount)
{
var c = new Claim { Id = Guid.NewGuid(), PolicyNumber = policy, Amount = amount, Status = "Submitted" };
c.Raise(new ClaimSubmitted(c.Id, policy, amount));
return c;
}
}
using System.Text.Json;
using Microsoft.EntityFrameworkCore;
using Microsoft.EntityFrameworkCore.Diagnostics;
public sealed class OutboxInterceptor : SaveChangesInterceptor
{
public override ValueTask<InterceptionResult<int>> SavingChangesAsync(
DbContextEventData eventData, InterceptionResult<int> result, CancellationToken ct = default)
{
var ctx = eventData.Context!;
var aggregates = ctx.ChangeTracker.Entries<AggregateRoot>().Select(e => e.Entity).ToList();
foreach (var agg in aggregates)
{
foreach (var e in agg.DomainEvents)
{
ctx.Set<OutboxMessage>().Add(new OutboxMessage
{
Type = e.GetType().AssemblyQualifiedName!,
Payload = JsonSerializer.Serialize(e, e.GetType())
});
}
agg.ClearEvents();
}
return base.SavingChangesAsync(eventData, result, ct);
}
}
public sealed class ClaimsDbContext(DbContextOptions<ClaimsDbContext> o) : DbContext(o)
{
public DbSet<Claim> Claims => Set<Claim>();
public DbSet<OutboxMessage> Outbox => Set<OutboxMessage>();
protected override void OnModelCreating(ModelBuilder b)
{
b.Entity<OutboxMessage>(m =>
{
m.ToTable("OutboxMessages");
m.HasKey(x => x.Id);
m.Property(x => x.Type).HasMaxLength(500);
// Partial/filtered index so the relay only scans pending rows
m.HasIndex(x => x.OccurredUtc).HasFilter("[ProcessedUtc] IS NULL");
});
}
}

Registration and the API endpoint (minimal API):

builder.Services.AddSingleton<OutboxInterceptor>();
builder.Services.AddDbContext<ClaimsDbContext>((sp, o) =>
o.UseSqlServer(builder.Configuration.GetConnectionString("Claims"))
.AddInterceptors(sp.GetRequiredService<OutboxInterceptor>()));
app.MapPost("/api/claims", async (SubmitClaimRequest r, ClaimsDbContext db, CancellationToken ct) =>
{
var claim = Claim.Submit(r.PolicyNumber, r.Amount);
db.Claims.Add(claim);
await db.SaveChangesAsync(ct); // claim + outbox row: one transaction
return Results.Accepted($"/api/claims/{claim.Id}", new { claim.Id });
});
public sealed record SubmitClaimRequest(string PolicyNumber, decimal Amount);

Angular side: because the flow is asynchronous, the API returns 202 Accepted; the UI shows “Submitted, processing” and observes progress (polling, SignalR, or a status endpoint).

// claim.service.ts (Angular, standalone, HttpClient)
import { inject, Injectable } from '@angular/core';
import { HttpClient } from '@angular/common/http';
@Injectable({ providedIn: 'root' })
export class ClaimService {
private http = inject(HttpClient);
submit(policyNumber: string, amount: number) {
return this.http.post<{ id: string }>('/api/claims', { policyNumber, amount });
}
status(id: string) {
return this.http.get<{ status: string }>(`/api/claims/${id}/status`);
}
}

Database note: PostgreSQL supports the same design. Use timestamptz, a partial index (WHERE processed_utc IS NULL), and FOR UPDATE SKIP LOCKED in the relay (Day 14).

Level 3: Advanced

Delivery guarantee. Outbox gives at-least-once, never exactly-once. A relay can publish a message and crash before marking it processed, so it will publish again. Consumers must be idempotent (Day 17). Give every message a stable MessageId (the outbox row Id) so consumers and the broker can deduplicate.

Ordering. Ordering is guaranteed only if you enforce it. Options: publish per aggregate in OccurredUtc/sequence order, use Service Bus sessions (SessionId = ClaimId) or Kafka partition keys (key = ClaimId), and run a single relay per partition of work. Global ordering is neither needed nor affordable.

Performance.

  • Index only pending rows (filtered/partial index); otherwise the table becomes a scan magnet.
  • Batch the relay (for example 100 rows per cycle) and batch sends to the broker.
  • Purge or archive processed rows (retention job, for example delete after 7 days) or use table partitioning; an ever-growing outbox slows inserts and the relay.
  • Keep payloads small; put large blobs in storage and send a reference (Claim Check).

Scalability of the relay. Several relay instances need row claiming to avoid duplicate work: SQL Server WITH (UPDLOCK, READPAST, ROWLOCK), PostgreSQL FOR UPDATE SKIP LOCKED. Duplicates are still possible, so idempotent consumers remain mandatory.

Security.

  • The outbox payload is data at rest: apply the same encryption, access control, and PII rules as business tables (Always Encrypted / TDE, column masking where needed). Do not put secrets or unnecessary PII in events.
  • The relay’s broker credential should be a managed identity with only “send” rights on specific topics.

Failure modes.

FailureEffectMitigation
Relay crashes after publish, before mark-processedDuplicate publishIdempotent consumers, broker duplicate detection
Poison payload (cannot serialize/deserialize)Relay blocked on the same rowAttempts counter, dead-letter/park after N failures, alert
Outbox table grows without cleanupSlow inserts, big backupsRetention job, partitioning
Relay is downEvents delayed but not lostAlert on “oldest unprocessed age”
Transaction rolled backOutbox row also rolled back: no phantom eventThis is the point of the pattern
Schema change of eventConsumers breakVersion event types, additive changes, contract tests (Day 51)

Common mistakes.

  1. Writing the outbox row in a different transaction/DbContext than the business change (defeats the pattern).
  2. Publishing inside the request and also via outbox (double publishing).
  3. Assuming exactly-once and skipping consumer idempotency.
  4. No monitoring of outbox lag.
  5. Deleting rows immediately after publish with no audit or retry visibility.
  6. Using DateTime.Now (local time) for ordering; use UTC or a sequence.

Level 4: Expert and Architect view

Alternatives compared:

ApproachAtomicity DB + eventDeliveryComplexityLatencyMain drawback
Dual write (naive)NoLossyLowLowestLost/phantom events
Transactional Outbox + relayYes (local ACID)At-least-onceMediumms to secondsExtra table and relay to operate
Outbox + CDC / log tailing (Day 13)YesAt-least-onceHigher infraLowNeeds CDC tooling and DB permissions
Outbox + polling publisher (Day 14)YesAt-least-onceLow to mediumPoll intervalQuery load, latency
2PC / distributed transaction (MSDTC/XA)YesExactly-once (theoretical)HighHighPoor availability, limited broker support
Event Sourcing as the source (Day 9)Yes (event store is the DB)At-least-once via projection publisherHighLowDifferent persistence model altogether
Listen-to-yourselfPartialAt-least-onceMediumMediumNeeds an idempotent local consumer

Combines with: Domain Event (Day 11: source of the rows), Polling Publisher / Transaction Log Tailing (relay strategies), Idempotent Consumer (Day 17: mandatory partner), Saga (Day 7: reliable step messages), Messaging (Day 16), Distributed Tracing (propagate traceparent inside the outbox row so traces continue across the async gap).

ADR-style justification

  • Title: ADR-012 Use Transactional Outbox for publishing Claims domain events.
  • Context: Claims service stores state in SQL Server and must notify Payments, Fraud, and Notifications through Azure Service Bus. Dual writes caused lost ClaimApproved events during deployments.
  • Decision: Persist events into an OutboxMessages table in the same EF Core transaction as the aggregate change; a background relay publishes to Service Bus using a managed identity; consumers are idempotent.
  • Alternatives considered: naive dual write (rejected: data loss), distributed transactions (rejected: unsupported/availability), CDC-based relay (deferred: revisit if relay latency or DB load becomes an issue).
  • Consequences: (+) no lost or phantom events, survives broker outages; (-) eventual consistency, at-least-once delivery, operational burden of relay and retention, outbox lag must be monitored.
  • Status: Accepted. Review when volume exceeds roughly 1,000 events per second per service or when CDC becomes available.

Azure implementation

Services involved

  • Azure SQL Database / SQL Server on Azure or Azure Database for PostgreSQL - Flexible Server: holds business tables and the outbox in one transaction.
  • Azure Service Bus (topics/subscriptions) as the broker. Alternatives: Event Hubs (streaming, Kafka protocol), Event Grid (event notifications).
  • Azure Container Apps, AKS, or App Service hosts the API and the relay (BackgroundService, or a separate worker/Container Apps job).
  • Azure Functions can host a timer-triggered relay. Note that timer triggers are at-least-once and may overlap unless singleton/locking is handled.
  • Managed Identity + Azure RBAC: role Azure Service Bus Data Sender for the relay; the consumers get Data Receiver.
  • Application Insights / Azure Monitor for metrics and traces; Key Vault for any remaining secrets.

Configuration highlights

  1. Create a Service Bus namespace and a topic claims-events. Enable duplicate detection on the topic (window, for example 10 minutes to 1 day) and set the outgoing MessageId to the outbox row Id. This removes many relay-retry duplicates, but it is a helper, not a replacement for idempotent consumers.
  2. Set SessionId (or use subscriptions with sessions enabled) to ClaimId when per-claim ordering is required.
  3. Relay authentication with DefaultAzureCredential:
using Azure.Identity;
using Azure.Messaging.ServiceBus;
using Microsoft.EntityFrameworkCore;
public sealed class OutboxRelay(IServiceScopeFactory scopes, ServiceBusClient client, ILogger<OutboxRelay> log)
: BackgroundService
{
private readonly ServiceBusSender _sender = client.CreateSender("claims-events");
protected override async Task ExecuteAsync(CancellationToken ct)
{
using var timer = new PeriodicTimer(TimeSpan.FromSeconds(2));
while (await timer.WaitForNextTickAsync(ct))
{
try { await PublishBatchAsync(ct); }
catch (Exception ex) when (ex is not OperationCanceledException)
{ log.LogError(ex, "Outbox relay cycle failed"); }
}
}
private async Task PublishBatchAsync(CancellationToken ct)
{
using var scope = scopes.CreateScope();
var db = scope.ServiceProvider.GetRequiredService<ClaimsDbContext>();
var batch = await db.Outbox
.Where(m => m.ProcessedUtc == null && m.Attempts < 10)
.OrderBy(m => m.OccurredUtc)
.Take(100)
.ToListAsync(ct);
foreach (var m in batch)
{
try
{
var msg = new ServiceBusMessage(m.Payload)
{
MessageId = m.Id.ToString(), // enables duplicate detection
Subject = m.Type,
ContentType = "application/json"
};
await _sender.SendMessageAsync(msg, ct);
m.ProcessedUtc = DateTime.UtcNow;
}
catch (Exception ex)
{
m.Attempts++;
m.Error = ex.Message;
}
}
await db.SaveChangesAsync(ct);
}
}
// Program.cs
builder.Services.AddSingleton(_ => new ServiceBusClient(
"<namespace>.servicebus.windows.net", new DefaultAzureCredential()));
builder.Services.AddHostedService<OutboxRelay>();

The simple version above uses a single relay instance. With several replicas, claim rows using locking hints as described in Level 3 (details in Day 14).

Pricing and tier considerations (verify current numbers on the Azure pricing page before budgeting; prices vary by region)

  • Service Bus Basic has no topics, so it cannot be used for this fan-out design. Standard supports topics/subscriptions and is billed as a base charge plus per-operation charges; Premium uses dedicated Messaging Units, gives predictable performance, VNet/private endpoint integration, and larger message sizes. Duplicate detection is available on Standard and Premium.
  • Azure SQL: the outbox adds a small write per event and a little storage; keep the retention job on to avoid paying for dead data. Choose the service tier by the total workload, not by the outbox.
  • Cost driver in practice: relay polling frequency (DB reads) and Service Bus operations. Batch sends reduce both.

Reference architecture (text)

Angular SPA -> Azure Front Door / API Management -> Claims API (Container Apps) -> Azure SQL (Claims + OutboxMessages in one transaction). The Outbox Relay (worker in the same Container Apps environment, managed identity) reads pending rows and sends to a Service Bus topic claims-events. Payments, Fraud, and Notifications services each have a subscription and idempotent consumers (with their own inbox/dedup table). Application Insights collects traces (trace context stored in the outbox row), and an Azure Monitor alert fires when the oldest unprocessed outbox row is older than, say, 5 minutes.

Teaching guide for my team

Explain to a beginner in 2 minutes

“When we save a claim we also need to tell other services. Saving to the database and sending a message are two separate systems, and the app can crash between them. So we do not send the message directly. We write the message into a table called Outbox in the same database transaction as the claim. A small background worker reads that table and sends the messages, marking each one as sent. If the worker or the broker fails, the message is still in the table and will be sent later. The claim and its message are saved together or not at all.”

Explain to an intermediate developer in 5 minutes

Walk through: the dual-write failure table; aggregate raises a domain event; the SaveChanges interceptor writes an OutboxMessage row in the same transaction; the relay batches pending rows, publishes with MessageId = row.Id, marks processed; at-least-once means consumers need idempotency; ordering needs session/partition keys; operations need lag monitoring and cleanup. Finish with the relay choices: polling (Day 14) versus CDC/log tailing (Day 13).

Hands-on exercise

  1. Build the Claims API with the OutboxInterceptor above and a POST /api/claims endpoint against a local SQL Server or PostgreSQL container.
  2. Add the OutboxRelay with a fake sender that writes to the console.
  3. Test A: submit a claim and confirm one Claims row and one OutboxMessages row appear and the relay marks it processed.
  4. Test B: throw an exception after db.Claims.Add(...) but before commit; confirm neither row exists.
  5. Test C: kill the relay between “send” and SaveChangesAsync; restart it and observe the message being sent a second time.

Expected outcome: A shows atomic write and delivery; B proves no phantom events; C demonstrates at-least-once and motivates the Idempotent Consumer.

Interview-style questions

  1. Why can’t we just call SaveChanges() then Publish() in a try/catch? The process can die between the two calls, and a try/catch cannot undo a commit or recover an event that was never recorded. The outbox makes the event part of the same atomic commit.
  2. Does the outbox give exactly-once delivery? No. It gives at-least-once. Duplicates can occur when the relay publishes but fails before marking the row. Consumers must be idempotent, or use broker duplicate detection as a partial help.
  3. How do you keep the outbox table from growing forever, and how do you know the relay is healthy? Run a retention job or partition to remove processed rows after a retention period, and monitor “age of the oldest unprocessed message” and the count of rows with failed attempts.

Mastery checklist

  • I can draw the dual-write failure timeline and show each place where the event can be lost or phantom.
  • I can implement the outbox write in the same EF Core transaction as the business change.
  • I can build a relay that publishes in batches, marks rows processed, and handles poison rows.
  • I can explain why the guarantee is at-least-once and design an idempotent consumer for it.
  • I can explain how ordering is (and is not) preserved, using sessions or partition keys.
  • I can run multiple relay instances safely (UPDLOCK, READPAST / SKIP LOCKED).
  • I can set up the Azure side: Service Bus topic, duplicate detection, managed identity, and lag alerts.
  • I can compare outbox with 2PC, CDC, and event sourcing in an ADR.

Key takeaway

Never write to the database and the broker separately: write the event to an outbox table in the same transaction, publish it afterwards, and make consumers tolerate duplicates.

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 13: Transaction Log Tailing

Transaction Log Tailing is the way to get the events in your Transactional Outbox out to a message broker without writing a polling loop.

Manikandan
Manikandan·23 min read
Microservices

Day 53: Client-Side UI Composition (Micro-Frontends)

Client-Side UI Composition splits one large single-page application into several smaller front-end applications ("micro-frontends") that are built, tested, and deployed independently by different teams.

Manikandan
Manikandan·17 min read