Manikandan — Manikandan
Microservices

Day 9: Event Sourcing

ManikandanManikandan
18 min read·Updated Aug 18, 2022

Event Sourcing stores every change to a business object as an immutable event in an append-only log, instead of overwriting the current state in a row.

Intro

Event Sourcing stores every change to a business object as an immutable event in an append-only log, instead of overwriting the current state in a row. The current state of a claim is not stored; it is derived by replaying its events (ClaimFiled, ClaimAssigned, ClaimApproved, …). You get a complete, trustworthy history for free, and the events double as the messages that other services subscribe to.

Why we need this

Business reasons

  • Insurance is a regulated domain. Auditors and regulators ask “who changed this claim’s reserve from 5,000 to 12,000, when, and why?” A traditional UPDATE destroys that answer.
  • Disputes and fraud investigation need the exact sequence of what happened, not only the end result.
  • Product teams want to ask new questions of old data (“how long do claims sit in UnderReview on average?”) without having recorded that in advance.

Technical reasons

  • In a microservice, changing state and telling others about it are two operations (the dual-write problem from Day 12). When the state is the events, saving the events is the publish.
  • Appending is cheap and contention-free compared with updating hot rows, and the log is a natural fit for replay, rebuilding read models (Day 8, CQRS) and debugging.
  • Time travel: you can rebuild the state as it was at any point in the past.

What problem it solves

Problem (from the topic list): Update/delete discards history, complicating audits, reconciliation, and event publishing.

What goes wrong without it (claims example):

  • Claims.Status = 'Approved' overwrites 'UnderReview'. Nobody can tell who approved it, when, or how long it took.
  • A team adds an AuditTrail table with triggers or hand-written inserts. It drifts from reality: a developer forgets one code path, and the audit table lies.
  • The service updates SQL and then publishes ClaimApproved to a broker. The process crashes between the two. Billing never hears about the approval.
  • Reconciliation with finance requires “what was the reserve at month end?” and the answer is gone because the row was updated in place.

Event sourcing makes the history the source of truth, so audit, publishing and reconciliation stop being separate, fragile add-ons.

When it is needed (and when it is NOT)

Good fit

  • The history is the business value: claims, payments, ledgers, orders, inventory movements, policy endorsements.
  • Strong audit or regulatory requirements (non-repudiation, “as-of” reporting).
  • Complex domain with behaviour-rich aggregates where the sequence of state changes matters.
  • You already need reliable event publishing and multiple read models (paired with CQRS).
  • You need temporal queries or the ability to fix a bug and rebuild derived data from the log.

Wrong choice / overkill

  • Simple CRUD reference data (a list of branch offices, claim categories, user preferences).
  • Teams with no DDD or messaging experience and tight delivery deadlines: the learning curve is real.
  • Data that must be truly erased (GDPR “right to erasure”) and cannot be handled by crypto-shredding or by keeping personal data outside the event payload.
  • Reporting-heavy systems where every screen needs ad-hoc joins and no one wants to build projections.
  • Applying it to the whole system. Use it per aggregate or per bounded context where it pays off. A Claims context can use it while Customer Master Data stays plain SQL.

How to identify the problem (key signals)

  1. Auditors or support ask “what was the value on date X?” or “who changed this?” and the honest answer is “we can’t know.”
  2. The team maintains hand-written history/audit tables or SQL triggers, and they are known to be incomplete.
  3. Production bugs of the form “the event was never published” or “the DB says Approved but downstream never got it.”
  4. Support staff manually reconstruct a claim’s timeline from application logs.
  5. Columns such as LastModifiedBy, PreviousStatus, StatusChangedOn keep multiplying on a single table, a smell that you are hand-rolling history.
  6. Domain experts describe the business in past-tense facts (“the claim was filed, then assigned, then approved”) but your schema stores only “current status”.
  7. Finance reconciliation numbers cannot be reproduced for prior periods.

Flow Diagram

Commands append events to a stream; state and read models are derived from the log.

flowchart LR
CMD["Command: ApproveClaim"] --> LOAD["Replay stream to rebuild Claim"]
LOAD --> CHK{"Valid in current state?"}
CHK -- "no" --> ERR["Reject command"]
CHK -- "yes" --> APP["Append ClaimApproved with expected version"]
APP --> ES[("Event store - append only")]
ES --> INL["Inline projection - snapshot"]
ES --> ASY["Async projections - worklists, reports"]
ES --> PUB["Integration events to Service Bus"]
ES --> TL["Audit timeline / time travel"]

Level 1: Beginner

Analogy: A bank statement versus a balance. A traditional table stores only “Balance: 1,200.” Event sourcing stores the statement: deposit 1,000, withdrawal 300, deposit 500. The balance is just the sum, and you can recompute it, or the balance as of last Tuesday, at any time. You never erase a line on the statement; to fix a mistake you add a correcting line.

Minimal working example (plain C#, no libraries). State is derived by folding events.

using System;
using System.Collections.Generic;
// Events: immutable facts, named in the past tense.
public abstract record ClaimEvent(Guid ClaimId, DateTimeOffset At);
public record ClaimFiled(Guid ClaimId, DateTimeOffset At, string PolicyNumber, decimal Amount)
: ClaimEvent(ClaimId, At);
public record ClaimApproved(Guid ClaimId, DateTimeOffset At, string ApprovedBy)
: ClaimEvent(ClaimId, At);
public enum ClaimStatus { None, Filed, Approved }
public class Claim
{
public Guid Id { get; private set; }
public ClaimStatus Status { get; private set; } = ClaimStatus.None;
public decimal Amount { get; private set; }
// Rebuild current state by replaying history.
public static Claim Replay(IEnumerable<ClaimEvent> history)
{
var claim = new Claim();
foreach (var e in history) claim.Apply(e);
return claim;
}
private void Apply(ClaimEvent e)
{
switch (e)
{
case ClaimFiled f: Id = f.ClaimId; Amount = f.Amount; Status = ClaimStatus.Filed; break;
case ClaimApproved: Status = ClaimStatus.Approved; break;
}
}
}
public static class Demo
{
public static void Main()
{
var id = Guid.NewGuid();
var log = new List<ClaimEvent> // the "event store"
{
new ClaimFiled(id, DateTimeOffset.UtcNow, "POL-1001", 4500m),
new ClaimApproved(id, DateTimeOffset.UtcNow, "adjuster.rao")
};
var claim = Claim.Replay(log);
Console.WriteLine($"{claim.Status}, {claim.Amount}"); // Approved, 4500
}
}

Key rules to remember: events are immutable, past-tense, and only ever appended.

Level 2: Intermediate

In a real application you use an event store (a database that provides per-stream append with optimistic concurrency). On .NET with PostgreSQL, Marten is a common choice: it turns PostgreSQL into a document DB plus event store. SQL Server has no equivalent first-class library, so teams either use a purpose-built store (KurrentDB, formerly EventStoreDB), or Cosmos DB, or build a simple Events table (see the schema below).

Rules of the aggregate

  1. A command is validated against current state; if valid it produces new events.
  2. Events are appended to the aggregate’s stream with an expected version (optimistic concurrency).
  3. State is rebuilt from the stream (optionally from a snapshot).
  4. Read screens use projections (read models), never replay in a query.

.NET side: aggregate + Marten (PostgreSQL)

// NuGet: Marten (check the current 8.x release; requires a supported .NET runtime)
using Marten;
public record ClaimFiled(Guid ClaimId, string PolicyNumber, decimal Amount);
public record ClaimApproved(Guid ClaimId, string ApprovedBy);
// Marten aggregates use Create/Apply conventions.
public class Claim
{
public Guid Id { get; set; }
public string Status { get; set; } = "None";
public decimal Amount { get; set; }
public string? ApprovedBy { get; set; }
public static Claim Create(ClaimFiled e) =>
new() { Id = e.ClaimId, Status = "Filed", Amount = e.Amount };
public void Apply(ClaimApproved e) { Status = "Approved"; ApprovedBy = e.ApprovedBy; }
}
public class ClaimService(IDocumentStore store)
{
public async Task<Guid> FileAsync(string policyNumber, decimal amount)
{
var id = Guid.NewGuid();
await using var session = store.LightweightSession();
session.Events.StartStream<Claim>(id, new ClaimFiled(id, policyNumber, amount));
await session.SaveChangesAsync(); // one ACID transaction
return id;
}
public async Task ApproveAsync(Guid id, string adjuster)
{
await using var session = store.LightweightSession();
var claim = await session.Events.AggregateStreamAsync<Claim>(id)
?? throw new KeyNotFoundException();
if (claim.Status != "Filed")
throw new InvalidOperationException("Only a filed claim can be approved.");
session.Events.Append(id, new ClaimApproved(id, adjuster));
await session.SaveChangesAsync();
}
}

Registration in Program.cs (ASP.NET Core, current .NET LTS):

builder.Services.AddMarten(opts =>
{
opts.Connection(builder.Configuration.GetConnectionString("Claims")!);
// Live read model kept up to date in the same transaction as the events:
opts.Projections.Snapshot<Claim>(SnapshotLifecycle.Inline);
});

For concurrency-safe commands, use Marten’s FetchForWriting<Claim>(id) so that two adjusters approving/rejecting the same claim at once produce a concurrency exception instead of a corrupt stream (verify the exact API in the docs for your Marten version).

If you must use SQL Server: a minimal event table (pseudo-schema; no library)

CREATE TABLE dbo.ClaimEvents (
GlobalSeq BIGINT IDENTITY PRIMARY KEY, -- total order for projections
StreamId UNIQUEIDENTIFIER NOT NULL, -- the claim id
StreamVersion INT NOT NULL, -- 1,2,3... per claim
EventType NVARCHAR(200) NOT NULL,
Payload NVARCHAR(MAX) NOT NULL, -- JSON
Metadata NVARCHAR(MAX) NULL, -- user, correlationId, causationId
OccurredAt DATETIMEOFFSET NOT NULL DEFAULT SYSDATETIMEOFFSET(),
CONSTRAINT UQ_Stream_Version UNIQUE (StreamId, StreamVersion) -- optimistic concurrency
);

The unique constraint on (StreamId, StreamVersion) is what gives you optimistic concurrency: if two writers try to append version 4, one fails and retries.

Angular side. The UI never reads events. It calls a query endpoint that returns a projection (e.g. GET /api/claims/{id}) and, optionally, a timeline endpoint that returns the event history for an “audit” tab.

// claim-timeline.component.ts (Angular, standalone component)
import { Component, inject, input } from '@angular/core';
import { HttpClient } from '@angular/common/http';
import { rxResource } from '@angular/core/rxjs-interop';
import { DatePipe } from '@angular/common';
interface TimelineItem { type: string; at: string; by?: string; }
@Component({
selector: 'app-claim-timeline',
imports: [DatePipe],
template: `
<ul>
@for (e of timeline.value() ?? []; track e.at) {
<li>{{ e.at | date:'medium' }} - {{ e.type }} @if (e.by) { by {{ e.by }} }</li>
}
</ul>`
})
export class ClaimTimelineComponent {
private http = inject(HttpClient);
claimId = input.required<string>();
timeline = rxResource({
params: () => this.claimId(),
stream: ({ params }) => this.http.get<TimelineItem[]>(`/api/claims/${params}/timeline`)
});
}

(rxResource with params/stream is the naming used from Angular 20 onward; on Angular 19 the option was request/loader. Check your version.)

Important consequence for the UI: because projections may be updated asynchronously, a screen can briefly show stale data right after a command. Return the new version number from the command and let the UI wait or refresh (eventual consistency).

Level 3: Advanced

Performance

  • Long streams: replaying 50,000 events per command is slow. Keep streams short by design (one stream per aggregate instance, not one per “Claims”), and use snapshots every N events (e.g. every 100).
  • Projections: inline projections cost write latency but are consistent; async projections (a background daemon reading the global sequence) scale better but are eventually consistent. Choose per read model.
  • Read path: never replay in a query. Serve reads from projections in PostgreSQL/SQL/Cosmos/Redis.
  • Indexes: the stream lookup (StreamId, StreamVersion) and a global ordered sequence are the two access paths that must be fast.

Scalability

  • Streams are the unit of concurrency and of partitioning. Use the aggregate id as the partition key (e.g. Cosmos DB partition key /streamId).
  • Consumers must be idempotent (Day 17): projections and subscribers can see the same event twice after a restart.

Security and privacy

  • Events are forever, so never put more personal data than needed in the payload. Store a CustomerId, not the name and national ID.
  • For unavoidable personal data, use crypto-shredding: encrypt the personal fields with a per-customer key (kept in Azure Key Vault or a key table) and delete the key to make the data unreadable. This is how teams reconcile immutability with erasure requests. Take legal advice on whether this satisfies your regulator.
  • Restrict who can read raw streams; they contain everything.

Failure modes and common mistakes

MistakeConsequenceFix
Events named as commands (ApproveClaim) or CRUD (ClaimUpdated)Loses business meaning; unusable for analyticsPast-tense business facts: ClaimApproved, ReserveIncreased
Changing or deleting old eventsBreaks replay and trustNever edit; append a compensating event (ClaimApprovalRevoked)
No event versioning strategyOld events cannot be deserialized after code changesUpcasting old shapes to new, additive changes, explicit version in metadata
Huge fat aggregates / one stream for everythingContention and slow replaySmall aggregates, one stream per instance
Using the event store as a query databaseSlow, complex queriesBuild projections
Publishing every internal event as an integration contractConsumers couple to your internalsMap to a smaller set of public integration events
Treating it as a system-wide defaultMassive accidental complexityApply only to contexts that benefit
Non-deterministic Apply methods (calling DateTime.Now, external APIs)Replay gives different state each timeApply only mutates state from event data; timestamps live in the event
Ignoring projection rebuild timeMulti-hour outages during schema changeBlue/green projections, rebuild in the background, then switch

Level 4: Expert and Architect view

Trade-offs versus alternatives

ApproachHistory / auditComplexityQuery flexibilityFits when
Plain CRUD (state only)None unless addedLowHigh (SQL)Simple reference data
CRUD + audit table / triggersPartial, can driftLow-MediumHighBasic audit need, no replay need
SQL Server temporal tables / PostgreSQL history extensionsRow versions, no intentLowHigh for as-of queriesYou need “as-of” but not business intent
CDC on a CRUD tableTechnical row changesMediumMediumFeed downstream, not domain history
Event SourcingComplete, with business intentHighNeeds projectionsHistory is the business value
Event Sourcing + CQRSComplete + optimized readsHighestHigh via read modelsComplex domains, scale, audit

Patterns it combines with

  • CQRS (Day 8): events feed the read models. Event sourcing without CQRS is painful for queries.
  • Domain Event (Day 11): in event sourcing, the persisted events are the domain events.
  • Saga (Day 7): saga state and steps are often event sourced too; events trigger the next step.
  • Transactional Outbox (Day 12): largely replaced, because the event store append is atomic; you still need a reliable subscription/publisher (or CDC/change feed) to push events out.
  • Idempotent Consumer (Day 17): mandatory on projections and subscribers.
  • Database per Service (Day 5): each service owns its own event store.

ADR (architecture review)

ADR-009: Use Event Sourcing for the Claims aggregate Status: Proposed Context: Claims must be fully auditable, reserve changes must be reproducible for month-end, and downstream services (Billing, Fraud) must be notified reliably of every state change. Today the Claims service updates rows in SQL Server and publishes to Service Bus in a second step; we have had lost-event incidents and incomplete audit history. Decision: The Claims bounded context will persist claim state as an append-only event stream per claim (PostgreSQL + Marten). Read models (claim list, adjuster worklist, reserve reports) are projections. Customer and Policy contexts remain state-based. Consequences (+): Complete audit trail by construction; atomic state-change + event; replayable projections; temporal queries. Consequences (-): Team learning curve; eventual consistency in read models; need for event versioning/upcasting and a projection rebuild process; GDPR handled with crypto-shredding and minimal PII in events. Alternatives considered: CRUD + audit table (drift risk), temporal tables (no business intent), CDC (technical, not domain events). Review trigger: Reassess if stream lengths or projection lag exceed agreed SLOs.

Azure implementation

Services that implement or support it

NeedAzure serviceNotes
Event store (relational)Azure Database for PostgreSQL – Flexible ServerRuns Marten as the event store. Good default for the .NET team.
Event store (NoSQL)Azure Cosmos DB for NoSQLPartition key = stream id; use a transactional batch for atomic multi-event append within one partition; the change feed streams new events to processors.
Event store (dedicated)KurrentDB (formerly EventStoreDB) via Kurrent Cloud on Microsoft Marketplace, or self-host on AKS/VMsPurpose-built event store; a third-party service billed through the marketplace.
Event distributionAzure Service Bus (commands, reliable integration events), Azure Event Hubs (high-volume streams)Publish integration events from the store.
Projection hostAzure Container Apps / AKS / App Service WebJob running the Marten async daemon or a Cosmos change feed processor; Azure Functions Cosmos DB trigger for simple cases
Read modelsAzure SQL Database, PostgreSQL, Cosmos DB, Azure Cache for Redis, Azure AI Search
Secrets/keysAzure Key VaultConnection strings, and per-customer keys for crypto-shredding.
Long-term archiveAzure Blob Storage (cool/archive tiers)Archive old closed-claim streams.
ObservabilityApplication Insights / Azure MonitorTrack projection lag, append latency, concurrency conflicts.

How to configure

  • PostgreSQL Flexible Server + Marten: create a Flexible Server, enable zone-redundant HA for production, enforce TLS, use Microsoft Entra authentication or Key Vault-referenced secrets, and let Marten create its schema on deploy (or apply migrations from CI). Turn on PITR backups and set retention according to your audit requirements.
  • Cosmos DB: container with partition key /streamId; store one document per event with streamId, version, type, data, metadata. Enforce ordering by version and use a unique key policy on (streamId, version) (unique keys are scoped per partition) for optimistic concurrency. Consume via the change feed processor (the default latest version mode is enough for event consumers, because events are only inserted). Keep single-stream appends in one partition so transactional batches work.
  • Service Bus: a topic per integration-event family; subscriptions per consumer; enable duplicate detection and sessions if you need per-claim ordering.
  • Identity: use managed identities from Container Apps/AKS to reach Key Vault, PostgreSQL and Service Bus; avoid keys in config.

Pricing and tier considerations (always confirm on the Azure pricing pages and the Azure pricing calculator; prices vary by region and change)

  • PostgreSQL Flexible Server: billed by compute tier (Burstable / General Purpose / Memory Optimized), storage and backup. Start with Burstable or a small General Purpose SKU for dev/test; use General Purpose with zone-redundant HA for production. Event tables grow forever, so budget storage growth and plan archiving.
  • Cosmos DB: provisioned throughput (manual or autoscale) or serverless. Serverless suits dev/test and spiky, low-volume workloads; autoscale suits production with variable load. Cost is driven by Request Units and storage, and writes are more expensive than point reads, so batch appends sensibly.
  • Service Bus: Standard tier for topics and duplicate detection; Premium when you need isolation, larger messages, or predictable performance and private networking.
  • Kurrent Cloud: subscription pricing through the Microsoft Marketplace; compare against operating PostgreSQL yourself.
  • Cost trap: projections rebuilt frequently on Cosmos DB consume RUs; run rebuilds off-peak or against a temporary higher throughput.

Reference architecture (text)

  1. Angular SPA calls the Claims API through Azure API Management / Application Gateway.
  2. The Claims API (ASP.NET Core on Azure Container Apps) handles commands: loads the claim stream from PostgreSQL (Marten), validates business rules, and appends new events with an expected version.
  3. Inline projections keep the Claim snapshot consistent; an async projection daemon (a separate Container App replica, one active daemon coordinated by Marten) builds the ClaimList and AdjusterWorklist read tables.
  4. A subscriber picks new events, maps them to integration events and publishes them to a Service Bus topic; Billing and Fraud services consume them idempotently.
  5. Queries hit the read models (PostgreSQL read tables or Azure SQL); the timeline endpoint reads the raw stream for the audit tab.
  6. Key Vault holds secrets and customer encryption keys; Application Insights collects traces (correlation id stored in event metadata); Blob Storage receives archived closed streams; PITR backups protect the store.

Teaching guide for my team

Explain to a beginner in 2 minutes

“Normally we store the latest value and overwrite it. With event sourcing we store what happened, like a bank statement, and we never erase lines. A claim isn’t a row that says ‘Approved’; it’s a list: Filed, Assigned, Approved. If you want the current state, you add the list up. If you want last Tuesday’s state, you add up the list until Tuesday. Because the list can’t be edited, the audit trail is always right.”

Explain to an intermediate developer in 5 minutes

Cover, in order: (1) events are immutable past-tense facts; (2) a command checks the current state (rebuilt from the stream) and appends events with an expected version, which gives optimistic concurrency; (3) Apply methods are pure and deterministic; (4) reads never replay: projections build read models, inline (consistent) or async (eventually consistent); (5) snapshots keep long streams fast; (6) events change over time, so you need versioning/upcasting; (7) personal data belongs outside events or must be crypto-shredded; (8) it pairs naturally with CQRS and needs idempotent consumers. Finish with the cost: more moving parts, so use it only where history is the business value.

Hands-on exercise

Task: Using the plain C# sample from Level 1, add ReserveChanged(ClaimId, At, decimal NewReserve, string Reason) and ClaimClosed. Then:

  1. Add a Claim.ReplayUntil(history, DateTimeOffset asOf) method that ignores events after asOf.
  2. Write a command method Approve that throws if the claim is not in Filed status, and returns the new event instead of mutating state.
  3. Add an expected-version check when appending to the in-memory list; simulate two concurrent approvals.

Expected outcome: Replaying the full log gives the current reserve and status; ReplayUntil with a past timestamp returns the older reserve; the second concurrent append throws a concurrency error; no event is ever modified or removed.

Interview-style questions

  1. How do you handle a mistake in an already-stored event? You never edit it. You append a compensating or correcting event (e.g. ReserveCorrected), so the history shows both the mistake and the fix.
  2. How do you query current state efficiently? With projections/read models updated from the event stream (inline or async), optionally with snapshots for aggregate loading. You do not replay on every query.
  3. How do you reconcile event sourcing with GDPR erasure? Keep personal data out of events where possible (use ids), and encrypt unavoidable personal fields with per-subject keys so deleting the key crypto-shreds the data; confirm the approach with your legal/compliance team.

Mastery checklist

  • I can explain the difference between state-based persistence and an event log, and give a business case where the history matters.
  • I can design events as past-tense business facts and reject CRUD-style event names.
  • I can implement an aggregate with Apply methods that are pure and deterministic.
  • I can implement optimistic concurrency on a stream with an expected version or a unique (StreamId, StreamVersion) constraint.
  • I can choose between inline and async projections and explain the consistency consequences to a UI developer.
  • I can describe snapshots, upcasting/versioning, and projection rebuilds.
  • I can explain how to handle personal data and erasure requests (minimal PII, crypto-shredding).
  • I can decide when NOT to use event sourcing and defend that decision in an architecture review.

Key takeaway

Event Sourcing makes the sequence of business facts your source of truth, giving you audit, replay and reliable event publishing by construction, at the price of projections, versioning and eventual consistency. Use it where history is the business value, not everywhere.

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 11: Domain Event

A Domain Event is an immutable record that something meaningful happened in the business domain, named in the past tense (for example ClaimApproved).

Manikandan
Manikandan·16 min read
Microservices

Day 10: API Composition

API Composition solves the "no more SQL JOIN" problem that appears once each microservice owns its own database.

Manikandan
Manikandan·13 min read
Microservices

Day 8: CQRS

CQRS (Command Query Responsibility Segregation) splits an application's model into two sides.

Manikandan
Manikandan·18 min read
Microservices

Day 7: Saga Pattern

A saga is a way to complete one business process that spans several services, each with its own database, without a distributed lock or a two-phase commit.

Manikandan
Manikandan·20 min read