Introduction
In a world where every click, every transfer, and every pollination can be recorded, the way we store change matters as much as the data itself. Traditional databases keep the current state of a system—a bank account balance, a hive’s health index, or the latest configuration of an autonomous AI agent. When a problem surfaces, engineers often have to dig through logs, reconstruct missing steps, or hope that a backup contains the truth. Event sourcing flips that model on its head: the source of truth is the immutable series of events that led to the present. By persisting each state transition as an indivisible fact, we gain a built‑in audit trail, deterministic replay, and a natural fit for distributed, self‑governing agents.
Why does this matter for a platform like Apiary, which intertwines bee conservation with self‑organizing AI? Bees themselves operate on a simple set of immutable actions— foraging, dancing, and queen‑laying—yet the colony’s health emerges from the sequence of those actions. Similarly, an AI agent that manages a financial ledger can be built from the same principle: each decision is an event, and the collective behavior emerges from replaying them. In the banking domain, where regulatory compliance demands traceability down to the cent, event sourcing offers a concrete, auditable, and flexible architecture that can also inspire how we think about ecological data and autonomous agents.
In this pillar article we’ll walk through the mechanics of event sourcing, illustrate it with a detailed banking example, and sprinkle in connections to bee data and AI governance where they naturally fit. By the end you’ll have a working mental model, concrete code‑level patterns, and a sense of why immutable events are becoming a cornerstone of modern, trustworthy systems.
What Is Event Sourcing?
Event sourcing is a persistence pattern where state changes are stored as a chronologically ordered log of events rather than as overwrites of a row in a table. Each event captures a fact that has already happened—e.g., “AccountCreated”, “MoneyDeposited”, “WithdrawalAttempted”. The system’s current state is reconstructed by replaying these events from the beginning of the stream. This is not merely a fancy audit log; it is the source of truth for the entire application.
The idea dates back to the early 2000s, popularized by Martin Fowler and later refined in the context of Domain‑Driven Design (DDD) domain-driven-design. A 2021 survey of 3,200 software engineers reported that 27 % of large‑scale fintech teams had adopted event sourcing for at least one microservice, citing regulatory compliance and debugging speed as top reasons. The pattern shines when you need auditability, temporal queries (what was the balance on 2023‑05‑12?), or system replay (re‑process a failed batch by replaying only the affected events).
Event sourcing also dovetails with Command‑Query Responsibility Segregation (CQRS), which separates the write model (commands that generate events) from the read model (projections built for queries). This separation enables each side to scale independently and to be optimized for its own workload, a crucial advantage for high‑throughput banking APIs that handle thousands of transactions per second.
Core Building Blocks
Events
An event is an immutable data structure that records a single state transition. It must be self‑describing (include a type identifier, timestamp, and payload) and idempotent (replaying it twice yields the same result). A typical MoneyDeposited event might look like:
{
"eventId": "e3f9c7a2-5b4d-4c9e-9f12-2a3d5e7b9c01",
"type": "MoneyDeposited",
"timestamp": "2024-03-15T14:23:05Z",
"accountId": "ACC-00123",
"amount": 2500,
"currency": "USD",
"metadata": { "source": "mobile-app" }
}
Because events never change, they can be safely replicated across data centers, archived for years, and even signed cryptographically to prove integrity—an appealing feature for regulators.
Aggregates
An aggregate is the consistency boundary that owns a set of events and enforces business invariants. In a banking domain, the BankAccount aggregate ensures that a withdrawal never pushes the balance below zero (unless an overdraft is explicitly allowed). The aggregate’s apply method takes an event and mutates the in‑memory state; the handle method validates a command, decides which events to emit, and returns them to the caller.
Streams
A stream (sometimes called an event stream or event log) is the ordered sequence of events for a given aggregate identifier. In most implementations the stream is stored in an event store, a purpose‑built database that guarantees append‑only writes and total ordering per stream. Streams enable optimistic concurrency: when you try to append a new event, you supply the expected stream version; if the version has changed, the store rejects the write, forcing you to reload and re‑apply the latest events.
These three concepts—event, aggregate, stream—form the backbone of any event‑sourced system, whether you’re tracking a bank balance or the foraging trips of a honeybee colony.
Modeling a Banking Account with Immutable Events
Let’s build a concrete example: a simple checking account that supports deposits, withdrawals, and transfers. We’ll model each action as an event and see how the aggregate enforces rules.
1. Defining the Event Types
| Event | Payload | Business Meaning |
|---|---|---|
AccountCreated | {accountId, ownerId, currency, createdAt} | Initializes a new aggregate. |
MoneyDeposited | {accountId, amount, currency, source} | Increases balance. |
MoneyWithdrawn | {accountId, amount, currency, destination} | Decreases balance, must respect overdraft policy. |
TransferInitiated | {fromAccountId, toAccountId, amount, correlationId} | Starts a two‑phase transfer. |
TransferCompleted | {correlationId, status} | Marks the transfer as succeeded or failed. |
Each event is stored with a global sequence number (often called a position or offset) assigned by the event store. This number enables temporal queries: “what was the balance after event #1,234,567?”
2. The Aggregate Logic
public class BankAccount {
private decimal _balance;
private readonly string _accountId;
private readonly string _currency;
private readonly List<Event> _uncommitted = new();
public BankAccount(string accountId, string currency) {
Apply(new AccountCreated(accountId, currency));
}
// Command handlers
public void Deposit(decimal amount, string source) {
if (amount <= 0) throw new ArgumentException("Amount must be positive");
var e = new MoneyDeposited(_accountId, amount, _currency, source);
ApplyAndStage(e);
}
public void Withdraw(decimal amount, string destination) {
if (amount <= 0) throw new ArgumentException("Amount must be positive");
if (_balance - amount < 0) throw new InvalidOperationException("Insufficient funds");
var e = new MoneyWithdrawn(_accountId, amount, _currency, destination);
ApplyAndStage(e);
}
// Event appliers
private void Apply(AccountCreated e) {
_accountId = e.AccountId;
_currency = e.Currency;
_balance = 0;
}
private void Apply(MoneyDeposited e) => _balance += e.Amount;
private void Apply(MoneyWithdrawn e) => _balance -= e.Amount;
// Helper
private void ApplyAndStage(Event e) {
Apply(e);
_uncommitted.Add(e);
}
public IReadOnlyCollection<Event> GetUncommitted() => _uncommitted;
}
The code demonstrates the single source of truth: the _balance field is never set directly; it is always the result of applying events. If we later need to rebuild the account after a crash, we simply load all events from the stream and call the corresponding Apply methods in order.
3. Replaying the Stream
public static BankAccount Rehydrate(string accountId, IEventStore store) {
var events = store.LoadStream(accountId);
var account = new BankAccount(); // empty ctor
foreach (var e in events) {
account.Apply(e); // internal, no staging
}
return account;
}
Notice that rehydration is deterministic: given the same ordered events, the aggregate always ends up in the same state. This property is essential for regulatory replay—if a central bank asks to see the ledger as of a specific date, you can simply replay up to that point.
4. Handling Transfers
Transfers are a classic case where eventual consistency shines. A TransferInitiated event is written to the source account’s stream, and a corresponding MoneyDeposited event is written to the destination account’s stream only after the source has successfully emitted a MoneyWithdrawn. The TransferCompleted event ties the two together, allowing auditors to trace the whole transaction across multiple aggregates.
5. Numbers in Practice
A mid‑size credit union with 150,000 active accounts generated 3.2 million events per month in 2023—roughly 110,000 events per day, or 1,300 per minute on average. By storing these events in an append‑only log, the organization reduced its nightly reconciliation window from 6 hours to under 10 minutes, because the nightly job now only needed to compute projections rather than scan raw tables.
The Event Store: Persistence, Ordering, and Retrieval
An event store is more than a table that holds JSON blobs. It must guarantee:
- Append‑only semantics – no updates or deletes of existing events.
- Total ordering per stream – each event gets a monotonically increasing version.
- Optimistic concurrency control – writes specify an expected version; mismatches trigger a retry.
- Scalable read access – queries for a stream’s events must be fast even when the stream contains millions of entries.
1. Storage Technologies
- Relational databases (e.g., PostgreSQL) can act as an event store using a table with a
BIGSERIALprimary key, aVARCHARfor stream ID, and aJSONBpayload. This is simple to set up and works well for low‑to‑moderate throughput (up to ~10 k writes/sec). - Purpose‑built log systems like EventStoreDB, Apache Kafka, or Azure Event Hubs provide higher throughput (hundreds of thousands of events per second) and built‑in replication. Kafka, for example, stores events in partitions; each partition is an ordered log that maps cleanly to a stream.
- NoSQL document stores (e.g., DynamoDB) can be used with a composite key (
streamId#version) to guarantee ordering and atomic writes.
Choosing the right backend depends on latency requirements, regulatory retention periods, and budget. For a banking microservice that must commit a transaction within 150 ms, an in‑memory write‑ahead cache layered on top of an SSD‑backed event store is common.
2. Snapshotting
Replaying a stream that contains 10 years of daily events can be costly. Snapshotting stores a serialized aggregate state after a certain number of events (e.g., every 5,000 events). When rehydrating, the system loads the latest snapshot and then replays only the events that occurred after that point. In practice, a snapshot size of 2 KB for a bank account is negligible, yet it can cut rehydration time from minutes to milliseconds.
3. Global vs. Per‑Stream Ordering
Regulators often require a global order of all financial events to detect race conditions like double‑spending. EventStoreDB provides a global position that increments across all streams, while Kafka’s offset is per‑partition. To achieve a global view, you can either:
- Use a single partition (limited throughput).
- Combine per‑stream positions with a global sequence generated by a dedicated service (e.g., a Snowflake‑style ID generator).
4. Querying the Store
Most queries are projection‑based (see cqrs). However, occasional ad‑hoc queries—like “find all accounts that deposited more than $10 k in the last 30 days”—can be served by materialized views built from the event stream. Tools such as Kafka Streams, Apache Flink, or EventStoreDB’s projection engine let you define continuous queries that update in near‑real time.
Queries and Projections: The CQRS Pattern
Event sourcing stores the write side of the system; the read side is built from projections—denormalized views that answer specific questions efficiently. This separation is the essence of CQRS (Command‑Query Responsibility Segregation).
1. Types of Projections
| Projection | Typical Store | Use‑Case |
|---|---|---|
| Account Balance View | In‑memory cache (Redis) or relational table | Real‑time balance display on mobile apps. |
| Transaction Ledger | Append‑only table or Elasticsearch | Full‑text search of deposits, compliance reporting. |
| Overdraft Alerts | Stream processing job (Kafka Streams) | Push notifications when balance < -$100. |
| Hive Health Index (Bee example) | Time‑series DB (InfluxDB) | Track pollen collection trends per day. |
Each projection consumes events from the event store, applies them, and writes the result to a storage technology optimized for that query pattern.
2. Building a Balance Projection
public class BalanceProjection {
private readonly IDictionary<string, decimal> _balances = new ConcurrentDictionary<string, decimal>();
public void When(MoneyDeposited e) => _balances[e.AccountId] += e.Amount;
public void When(MoneyWithdrawn e) => _balances[e.AccountId] -= e.Amount;
public decimal GetBalance(string accountId) => _balances.TryGetValue(accountId, out var bal) ? bal : 0m;
}
A background worker subscribes to the event store (via a subscription in EventStoreDB or a consumer group in Kafka) and invokes the appropriate handler for each incoming event. Because the projection is idempotent, if the same event is delivered twice (possible in at‑least‑once delivery semantics), the balance remains correct.
3. Handling Eventual Consistency
Since the read model is built asynchronously, there is a short window where a newly deposited amount may not yet appear in the balance view. In banking, this latency is usually acceptable as long as critical operations (e.g., another withdrawal) consult the write model for validation. For user‑facing UI, you can show a pending state until the projection catches up.
4. Scaling Projections
Projections can be sharded by stream identifier. For example, a balance projection can be partitioned by the first two characters of the account ID, allowing multiple worker instances to process disjoint subsets of events in parallel. This approach scales linearly with the number of accounts, which is essential for systems handling millions of customers.
Evolving Schemas: Versioning and Compatibility
In a live banking environment, you cannot freeze the data model forever. Event schemas must evolve without breaking the ability to replay historic streams.
1. Upcasting
Upcasting transforms an old event version into the latest shape before it reaches the aggregate. Suppose the MoneyDeposited event originally stored currency as a three‑letter code, but later you add a exchangeRate field for multi‑currency support. An upcaster reads the old JSON, adds a default exchangeRate = 1.0, and passes the enriched event downstream. This transformation is performed once, at load time, keeping the stored events untouched.
2. Downcasting (Rare)
If a new version removes a field, you may need a downcaster for compatibility with older projection services that still expect the removed field. However, the preferred strategy is to keep old fields optional and ignore them in newer code.
3. Event Version Numbers
Each event type carries a schema version (e.g., MoneyDeposited.v2). When a new version is introduced, you add a new handler for that version while keeping the old handler for backward compatibility. Over time, you can migrate old streams by re‑writing them into a new compact stream, but this is usually done during a major system upgrade, not during normal operation.
4. Real‑World Numbers
A large European bank migrated 12 years of event history (≈ 250 million events) to a new schema over a weekend by running parallel upcasters on a 48‑node Spark cluster. The operation took 7 hours, and the resulting event store size grew by only 3 % due to added metadata.
Consistency, Transactions, and Idempotency
Event sourcing does not eliminate the need for transactional guarantees; it re‑frames them.
1. Atomic Append
When a command generates multiple events (e.g., a transfer that withdraws from one account and deposits into another), the system must guarantee that all events are appended atomically. This is typically achieved by a transactional outbox pattern: the command handler writes the events to a local relational database within a transaction, then a background process publishes them to the event store. If the publish fails, the transaction is rolled back, preserving atomicity.
2. Idempotent Handlers
Because events may be delivered more than once, both aggregate apply methods and projection handlers must be idempotent. A common technique is to store the eventId in a deduplication table; before processing, check if the ID has already been seen. In high‑throughput systems, a Bloom filter can provide a probabilistic, low‑memory deduplication layer.
3. Handling Concurrency
Optimistic concurrency uses the expected version pattern. When appending a new event, the client supplies the last known version (expectedVersion). If another writer has already added an event, the store returns a ConcurrencyException, prompting the client to reload the stream, reapply business logic, and retry. In banking, this prevents two simultaneous withdrawals from both seeing the same pre‑withdrawal balance.
4. ACID vs. BASE
Event sourcing leans toward BASE (Basically Available, Soft state, Eventual consistency) for the read side, while the write side remains ACID thanks to atomic appends and strict version checks. This hybrid model satisfies both regulatory precision and modern scalability.
Testing, Debugging, and Auditing
A major selling point of event sourcing is observability. Every state change is a first‑class artifact, which simplifies testing and forensic analysis.
1. Unit‑Testing Aggregates
Because aggregates are pure functions of events and commands, you can write given‑when‑then tests:
// Given
var account = new BankAccount();
account.Apply(new AccountCreated("ACC-1", "USD"));
// When
account.Deposit(100, "online");
// Then
account.GetUncommitted().ShouldContain(e => e is MoneyDeposited);
No database is needed; the test runs in memory, executes instantly, and guarantees that business rules are enforced.
2. Integration Tests with Event Store
For end‑to‑end verification, spin up a lightweight in‑memory event store (e.g., EventStoreDB in Docker) and run a scenario that writes commands, reads projections, and asserts final balances. Because the store persists events, you can pause a test, inspect the raw JSON files, and manually replay them to understand failures.
3. Auditing Tools
Many event stores expose a REST endpoint that streams events in order. Auditors can pull a JSON dump for a given period and run compliance scripts that check, for example, “no withdrawal exceeds 10 % of the account’s average monthly deposit”. Because the data is immutable, the audit trail cannot be tampered with without detection.
4. Debugging with Time Travel
Developers can time‑travel by loading a snapshot at a