Event Sourcing and CQRS in Production: Processing 50K Events/Second
Production patterns for event sourcing at scale — how we process 50K events per second with CQRS, handle schema evolution, and maintain sub-100ms read projections.

Event sourcing sounds elegant in conference talks. In production at 50,000 events per second, it becomes an exercise in managing complexity, storage growth, schema evolution, and projection rebuild times that can stretch into hours. After three years of operating an event-sourced platform processing financial transactions, I can tell you what actually works — and what the tutorials conveniently skip.
This article covers the production patterns that keep our event store healthy, our projections consistent, and our engineers sane. No toy examples. Real throughput numbers, real failure modes, and the operational tooling you need before going live.
Why Event Sourcing (And When Not To)
We adopted event sourcing for our payment processing domain because the audit trail is the product. Financial regulators require us to prove the exact sequence of state changes for any account over any time period. With traditional CRUD, we would need separate audit logging — essentially building event sourcing with extra steps and worse consistency guarantees.
Event sourcing is not appropriate for every domain. We use it selectively:
| Domain | Pattern | Reasoning |
|---|---|---|
| Payment processing | Event sourced | Regulatory audit trail, temporal queries required |
| User accounts | CRUD + audit log | Simple state, low write volume, no temporal queries |
| Notifications | Event-driven (not sourced) | Fire-and-forget, no replay value |
| Analytics | Event streaming | Append-only by nature, but no aggregate reconstruction |
| Inventory | Event sourced | Concurrency conflicts need event-level resolution |
Architecture Overview
Our event sourcing infrastructure consists of three layers: the event store, the command processing pipeline, and the projection engine.
Event Store Design
We use Apache Kafka as our event transport and PostgreSQL as our persistent event store. This might surprise purists who expect a purpose-built event store, but PostgreSQL gives us ACID guarantees on writes, mature operational tooling, and the ability to run complex queries against the event log when needed.
-- Event store schema with optimistic concurrency control
CREATE TABLE events (
event_id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
aggregate_id UUID NOT NULL,
aggregate_type VARCHAR(128) NOT NULL,
event_type VARCHAR(256) NOT NULL,
event_version INTEGER NOT NULL,
sequence_number BIGSERIAL,
payload JSONB NOT NULL,
metadata JSONB NOT NULL DEFAULT '{}',
schema_version INTEGER NOT NULL DEFAULT 1,
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
-- Optimistic concurrency: unique constraint on aggregate + version
CONSTRAINT unique_aggregate_version
UNIQUE (aggregate_id, event_version)
);
-- Partitioning by month for manageable table sizes
CREATE TABLE events_2025_12 PARTITION OF events
FOR VALUES FROM ('2025-12-01') TO ('2026-01-01');
-- Index for aggregate replay (most common read pattern)
CREATE INDEX idx_events_aggregate_replay
ON events (aggregate_id, event_version ASC);
-- Index for projection catch-up
CREATE INDEX idx_events_sequence
ON events (sequence_number ASC)
WHERE sequence_number > 0;
-- Index for event type queries (debugging, analytics)
CREATE INDEX idx_events_type_created
ON events (event_type, created_at DESC);
Command Processing Pipeline
Commands are validated, processed, and converted to events through a strict pipeline. Each aggregate has a command handler that enforces business invariants before emitting events.
// Command handler with optimistic concurrency and retry logic
import { EventStore, Aggregate, Command, Event } from './core';
interface ProcessPaymentCommand extends Command {
type: 'ProcessPayment';
accountId: string;
amount: number;
currency: string;
idempotencyKey: string;
merchantId: string;
}
interface PaymentProcessedEvent extends Event {
type: 'PaymentProcessed';
accountId: string;
amount: number;
currency: string;
balanceAfter: number;
merchantId: string;
processedAt: string;
}
class PaymentCommandHandler {
constructor(
private eventStore: EventStore,
private maxRetries: number = 3
) {}
async handle(command: ProcessPaymentCommand): Promise<PaymentProcessedEvent> {
// Idempotency check — critical for exactly-once processing
const existing = await this.eventStore.findByIdempotencyKey(
command.idempotencyKey
);
if (existing) {
return existing as PaymentProcessedEvent;
}
let attempt = 0;
while (attempt < this.maxRetries) {
try {
// Load current aggregate state from event history
const account = await this.loadAggregate(command.accountId);
// Validate business invariants
this.validatePayment(account, command);
// Generate event
const event: PaymentProcessedEvent = {
type: 'PaymentProcessed',
aggregateId: command.accountId,
version: account.version + 1,
accountId: command.accountId,
amount: command.amount,
currency: command.currency,
balanceAfter: account.balance - command.amount,
merchantId: command.merchantId,
processedAt: new Date().toISOString(),
metadata: {
idempotencyKey: command.idempotencyKey,
correlationId: command.correlationId,
causationId: command.id
}
};
// Persist with optimistic concurrency
await this.eventStore.append(event);
return event;
} catch (error) {
if (error.code === 'CONCURRENCY_CONFLICT') {
attempt++;
await this.backoff(attempt);
continue;
}
throw error;
}
}
throw new Error(`Command failed after ${this.maxRetries} retries`);
}
private validatePayment(account: AccountAggregate, cmd: ProcessPaymentCommand): void {
if (account.status !== 'active') {
throw new BusinessRuleViolation('Account is not active');
}
if (account.balance < cmd.amount) {
throw new BusinessRuleViolation('Insufficient funds');
}
if (cmd.amount <= 0) {
throw new BusinessRuleViolation('Amount must be positive');
}
}
private async backoff(attempt: number): Promise<void> {
const delay = Math.min(100 * Math.pow(2, attempt), 2000);
await new Promise(resolve => setTimeout(resolve, delay));
}
}
Projection Engine
Projections transform the event stream into query-optimized read models. We run multiple projections per domain, each optimized for specific read patterns.
Projection Performance
| Projection | Events/Second | Lag (p99) | Storage | Rebuild Time |
|---|---|---|---|---|
| Account Balance | 50,000 | 45ms | PostgreSQL | 4.2 hours |
| Transaction History | 50,000 | 62ms | Elasticsearch | 6.8 hours |
| Daily Aggregates | 50,000 | 180ms | ClickHouse | 2.1 hours |
| Fraud Signals | 50,000 | 23ms | Redis | 45 minutes |
| Regulatory Report | 50,000 | 340ms | PostgreSQL | 8.4 hours |
The rebuild times are our biggest operational concern. When a projection has a bug or schema change, we must rebuild from the full event history. For the Transaction History projection, that means replaying 4.3 billion events over 6.8 hours.
Schema Evolution
Event schemas evolve over time. We handle this with explicit versioning and upcasters — transformers that convert old event formats to current format at read time.
// Event upcaster chain for schema evolution
interface EventUpcaster {
eventType: string;
fromVersion: number;
toVersion: number;
upcast(event: RawEvent): RawEvent;
}
const paymentProcessedUpcasters: EventUpcaster[] = [
{
eventType: 'PaymentProcessed',
fromVersion: 1,
toVersion: 2,
upcast(event: RawEvent): RawEvent {
// V1 -> V2: Added currency field (default to USD for historical events)
return {
...event,
payload: {
...event.payload,
currency: event.payload.currency || 'USD'
},
schemaVersion: 2
};
}
},
{
eventType: 'PaymentProcessed',
fromVersion: 2,
toVersion: 3,
upcast(event: RawEvent): RawEvent {
// V2 -> V3: Renamed merchantId to merchant.id, added merchant.name
return {
...event,
payload: {
...event.payload,
merchant: {
id: event.payload.merchantId,
name: null // Will be enriched by projection
}
},
schemaVersion: 3
};
}
}
];
class UpcasterPipeline {
private upcasters: Map<string, EventUpcaster[]> = new Map();
register(upcaster: EventUpcaster): void {
const key = upcaster.eventType;
const chain = this.upcasters.get(key) || [];
chain.push(upcaster);
chain.sort((a, b) => a.fromVersion - b.fromVersion);
this.upcasters.set(key, chain);
}
upcast(event: RawEvent): RawEvent {
const chain = this.upcasters.get(event.eventType) || [];
let current = event;
for (const upcaster of chain) {
if (current.schemaVersion === upcaster.fromVersion) {
current = upcaster.upcast(current);
}
}
return current;
}
}
Snapshotting Strategy
At 50K events/second, aggregate event histories grow rapidly. Loading a full account history (potentially millions of events) on every command would be unacceptable. We snapshot aggregate state every 1,000 events.
| Snapshot Strategy | Load Time (avg) | Storage Cost | Consistency |
|---|---|---|---|
| No snapshots | 340ms | None | Perfect |
| Every 100 events | 4ms | High (10x snapshots) | Perfect |
| Every 1,000 events | 12ms | Moderate | Perfect |
| Every 10,000 events | 48ms | Low | Perfect |
We chose every 1,000 events as our sweet spot — 12ms average load time with manageable storage overhead.
Operational Concerns
Event Store Compaction
Our event store grows at approximately 4.3 billion events per month. We partition by time and archive partitions older than 90 days to cold storage (S3 with Parquet format). Archived events remain queryable through our data lake but are not part of the hot event store.
Monitoring and Alerting
Key metrics we monitor:
- Projection lag — Alert if any projection falls more than 5 seconds behind
- Command retry rate — High retry rate indicates concurrency hotspots
- Event store write latency — p99 above 50ms triggers investigation
- Snapshot freshness — Stale snapshots increase aggregate load times
- Dead letter queue depth — Events that failed processing need immediate attention
Conclusion
Event sourcing at scale is an operational commitment, not just an architectural pattern. The benefits — complete audit trail, temporal queries, event replay for debugging — are substantial but come with real costs: complex projections, schema evolution overhead, and lengthy rebuild times.
If you are considering event sourcing, start with a single bounded context where the audit trail provides genuine business value. Build the operational tooling (projection monitoring, snapshot management, schema evolution pipeline) before scaling out. The pattern works beautifully at 50K events/second, but only if you respect the complexity it introduces.
Recommended reading

Per-Team Cost Allocation in Shared Kubernetes Clusters: From Chaos to Clarity
Implementing accurate per-namespace cost allocation in multi-tenant Kubernetes clusters, covering request vs. usage attribution, shared resource amortization, and building showback dashboards that drive accountability.

Measuring and Eliminating Toil: From 40% to 12% of Engineering Time
A systematic approach to identifying, measuring, and automating toil—the repetitive operational work that scales linearly with service growth and prevents engineers from doing creative work.

Serverless Postgres in Production: Branching, Scale-to-Zero, and the End of Database Provisioning
Running Neon serverless Postgres in production for 8 months — covering database branching workflows, scale-to-zero economics, connection pooling, and migration from RDS.

Comments
No comments yet. Be the first to share your thoughts.