| name | Event Sourcing |
| description | Architectural pattern storing state changes as immutable events for complete audit trail and temporal queries |
| category | software-development |
Event Sourcing
What I do
I provide an architectural pattern where application state is stored as a sequence of events rather than as current state. Every state change is captured as an immutable event, creating a complete audit trail and enabling temporal queries. Event sourcing allows reconstructing any past state, provides full audit capability, and enables sophisticated event processing patterns. The event log becomes the source of truth.
When to use me
Use event sourcing when you need a complete audit trail, when temporal queries are important, or when business rules depend on event history. It's valuable for domain models where state transitions are complex. Event sourcing excels in systems requiring eventual consistency, audit compliance, or the ability to replay events. Avoid for simple CRUD applications or when you need only current state efficiently.
Core Concepts
- Event: Immutable record of something that happened
- Event Store: Specialized database for events
- Aggregate: Entity reconstructed from events
- Snapshot: Optimization to avoid replaying all events
- Projection: Building read models from events
- Event Versioning: Handling evolving event schemas
- Replay: Rebuilding state from event history
- Append-Only: Events are never modified or deleted
- Domain Events: Business events from domain model
- Integration Events: Events for external systems
Code Examples
Event Definitions
from abc import ABC
from dataclasses import dataclass, field
from datetime import datetime
from typing import Generic, TypeVar, Protocol
from uuid import UUID, uuid4
T = TypeVar("T", bound='Event')
@dataclass
class Event:
event_id: UUID
aggregate_id: UUID
event_type: str
occurred_at: datetime
version: int
def to_dict(self) -> dict:
return {
"event_id": str(self.event_id),
"aggregate_id": str(self.aggregate_id),
"event_type": self.event_type,
"occurred_at": self.occurred_at.isoformat(),
"version": self.version
}
@dataclass
class UserCreated(Event):
email: str
name: str
def __init__(
self,
aggregate_id: UUID,
email: str,
name: str,
version: int
):
().__init__(
event_id=uuid4(),
aggregate_id=aggregate_id,
event_type=,
occurred_at=datetime.utcnow(),
version=version
)
.email = email
.name = name
():
old_email:
new_email:
():
().__init__(
event_id=uuid4(),
aggregate_id=aggregate_id,
event_type=,
occurred_at=datetime.utcnow(),
version=version
)
.old_email = old_email
.new_email = new_email
():
reason:
():
().__init__(
event_id=uuid4(),
aggregate_id=aggregate_id,
event_type=,
occurred_at=datetime.utcnow(),
version=version
)
.reason = reason
():
reason:
():
().__init__(
event_id=uuid4(),
aggregate,
event_type=,
occurred_at=datetime.utcnow(),
version=version
)
.reason = reason
Aggregate Root with Event Sourcing
from abc import ABC, abstractmethod
from typing import TypeVar, Generic
E = TypeVar("E", bound=Event)
class AggregateRoot(ABC, Generic[E]):
def __init__(self, aggregate_id: UUID):
self._id = aggregate_id
self._version = 0
self._events: list[E] = []
@property
def id(self) -> UUID:
return self._id
@property
def version(self) -> int:
return self._version
@property
def uncommitted_events(self) -> list[E]:
return self._events.copy()
def clear_uncommitted_events(self) -> None:
self._events.clear()
def _apply(self, event: E) -> None:
self._version +=
._events.append(event)
._apply_event(event)
() -> :
(AggregateRoot[Event]):
():
().__init__(user_id)
._email: | =
._name: | =
._is_active: =
._email_history: [] = []
() -> :
aggregate = cls(uuid4())
event = UserCreated(
aggregate_id=aggregate._,
email=email,
name=name,
version=
)
aggregate._apply(event)
aggregate
() -> :
._is_active:
ValueError()
new_email == ._email:
event = UserEmailChanged(
aggregate_id=._,
old_email=._email ,
new_email=new_email,
version=._version +
)
._apply(event)
() -> :
._is_active:
ValueError()
event = UserDeactivated(
aggregate_id=._,
reason=reason,
version=._version +
)
._apply(event)
() -> :
._is_active:
ValueError()
event = UserReactivated(
aggregate_id=._,
reason=reason,
version=._version +
)
._apply(event)
() -> :
(event, UserCreated):
._email = event.email
._name = event.name
._is_active =
._email_history = [event.email]
(event, UserEmailChanged):
._email = event.new_email
._email_history.append(event.new_email)
(event, UserDeactivated):
._is_active =
(event, UserReactivated):
._is_active =
() -> | :
._email
() -> | :
._name
() -> :
._is_active
() -> []:
._email_history.copy()
Event Store
from abc import ABC, abstractmethod
from typing import Protocol, TypeVar
E = TypeVar("E", bound=Event)
class EventStore(Protocol[E]):
@abstractmethod
def append(self, event: E) -> None:
pass
@abstractmethod
def get_events(self, aggregate_id: UUID) -> list[E]:
pass
@abstractmethod
def get_all_events(self, from_version: int = 0) -> list[E]:
pass
class InMemoryEventStore(EventStore[Event]):
def __init__(self):
self._events: list[Event] = []
self._by_aggregate: dict[UUID, list[Event]] = {}
def append(self, event: Event) -> None:
self._events.append(event)
if event.aggregate_id not in self._by_aggregate:
self._by_aggregate[event.aggregate_id] = []
self._by_aggregate[event.aggregate_id].append(event)
() -> [Event]:
._by_aggregate.get(aggregate_id, []).copy()
() -> [Event]:
[e e ._events e.version > from_version]
:
():
._snapshots: [UUID, ] = {}
() -> :
._snapshots[aggregate_id] = {
: version,
: state
}
() -> | :
._snapshots.get(aggregate_id)
Repository with Snapshots
class UserRepository:
def __init__(
self,
event_store: EventStore[Event],
snapshot_store: SnapshotStore | None = None,
snapshot_threshold: int = 10
):
self._event_store = event_store
self._snapshot_store = snapshot_store
self._snapshot_threshold = snapshot_threshold
def save(self, aggregate: UserAggregate) -> None:
for event in aggregate.uncommitted_events:
self._event_store.append(event)
aggregate.clear_uncommitted_events()
if self._snapshot_store and self._should_snapshot(aggregate):
self._save_snapshot(aggregate)
def load(self, aggregate_id: UUID) -> UserAggregate:
snapshot = None
if self._snapshot_store:
snapshot = self._snapshot_store.get_snapshot(aggregate_id)
if snapshot:
events = self._event_store.get_events(aggregate_id)
events_after_snapshot = [
e for e in events
if e.version > snapshot["version"]
]
aggregate = self._reconstitute(snapshot["state"])
event events_after_snapshot:
aggregate._apply(event)
aggregate
events = ._event_store.get_events(aggregate_id)
._reconstitute_from_events(events)
() -> :
(
._snapshot_store
aggregate.version % ._snapshot_threshold ==
)
() -> :
state = {
: aggregate.email,
: aggregate.name,
: aggregate.is_active,
: aggregate.email_history
}
._snapshot_store.save_snapshot(
aggregate.,
aggregate.version,
state
)
() -> UserAggregate:
aggregate = UserAggregate(uuid4())
aggregate._email = state[]
aggregate._name = state[]
aggregate._is_active = state[]
aggregate._email_history = state[]
aggregate
() -> UserAggregate:
aggregate = UserAggregate(uuid4())
event events:
aggregate._apply(event)
aggregate
Projections for Read Models
from abc import ABC, abstractmethod
from dataclasses import dataclass, field
from datetime import datetime
from typing import Generic, TypeVar, Protocol
P = TypeVar("P", bound='Projection')
@dataclass
class UserView:
user_id: UUID
email: str
name: str
is_active: bool
created_at: datetime
updated_at: datetime
version: int
class Projection(ABC):
@abstractmethod
def apply(self, event: Event) -> None:
pass
class UserProjection(Projection):
def __init__(self):
self._users: dict[UUID, UserView] = {}
@property
def users(self) -> dict[UUID, UserView]:
return self._users.copy()
def get_user(self, user_id: UUID) -> UserView | None:
return self._users.get(user_id)
() -> [UserView]:
[u u ._users.values() u.is_active]
() -> :
(event, UserCreated):
._users[event.aggregate_id] = UserView(
user_id=event.aggregate_id,
email=event.email,
name=event.name,
is_active=,
created_at=event.occurred_at,
updated_at=event.occurred_at,
version=event.version
)
(event, UserEmailChanged):
user := ._users.get(event.aggregate_id):
user.email = event.new_email
user.updated_at = event.occurred_at
user.version = event.version
(event, UserDeactivated):
user := ._users.get(event.aggregate_id):
user.is_active =
user.updated_at = event.occurred_at
user.version = event.version
:
():
._projections: [Projection] = []
() -> :
._projections.append(projection)
() -> :
projection ._projections:
:
projection.apply(event)
Exception e:
()
:
():
._event_store = event_store
._projection_manager = ProjectionManager()
._last_processed_version: [, ] = {}
() -> :
._projection_manager.add(projection)
() -> :
threading
():
:
event ._event_store.get_all_events():
._projection_manager.apply(event)
threading.Thread(target=process, daemon=).start()
Event Versioning
from abc import ABC, abstractmethod
from datetime import datetime
from typing import dict
class EventUpgrader(ABC):
@abstractmethod
def can_upgrade(self, event_type: str, version: int) -> bool:
pass
@abstractmethod
def upgrade(self, event: dict) -> dict:
pass
class UserEventUpgrader(EventUpgrader):
V1_TO_V2_MAPPING = {
"email": "contact_email"
}
def can_upgrade(self, event_type: str, version: int) -> bool:
return event_type == "UserCreated" and version == 1
def upgrade(self, event: dict) -> dict:
upgraded = event.copy()
upgraded["version"] = 2
for old_field, new_field in self.V1_TO_V2_MAPPING.items():
old_field upgraded:
upgraded[new_field] = upgraded.pop(old_field)
upgraded[] = {
: ,
: datetime.utcnow().isoformat()
}
upgraded
:
():
._upgraders: [EventUpgrader] = []
() -> :
._upgraders.append(upgrader)
() -> :
current = event
upgrader ._upgraders:
upgrader.can_upgrade(
current[],
current[]
):
current = upgrader.upgrade(current)
current
:
():
._store = base_store
._upgrader_chain = EventUpgraderChain()
._upgrader_chain.add_upgrader(UserEventUpgrader())
() -> :
._store.append(event)
() -> []:
events = ._store.get_events(aggregate_id)
[._upgrader_chain.upgrade(e.to_dict()) e events]
Best Practices
- Immutable Events: Never modify or delete events
- Idempotent Appends: Handle duplicate event writes safely
- Version Events: Plan for event schema evolution
- Use Snapshots: Optimize loading large event histories
- Atomic Writes: Save events in transactions
- Compensating Events: For undo actions, emit inverse events
- Projection Async: Projections should not block writes
- Event Ordering: Preserve order within aggregates
- Backwards Compatibility: New code should read old events
- Testing: Test aggregates through their event history
- Governance: Document event schemas thoroughly
- Performance: Index events by aggregate ID and type