| name | event-driven-architect |
| description | Designs event-driven architectures with event sourcing, CQRS, pub/sub patterns, and domain events for decoupled systems. Use when users request "event sourcing", "CQRS", "domain events", "pub/sub", or "event-driven". |
Event-Driven Architect
Build decoupled, scalable systems with event-driven patterns.
Core Workflow
- Identify domain events: Define what happened
- Design event schema: Structure event payloads
- Implement event bus: Publish and subscribe
- Add event handlers: React to events
- Consider CQRS: Separate reads and writes
- Enable event sourcing: Store event history
Event Fundamentals
Event Structure
export interface DomainEvent<T = unknown> {
id: string;
type: string;
aggregateId: string;
aggregateType: string;
payload: T;
metadata: {
timestamp: Date;
version: number;
correlationId?: string;
causationId?: string;
userId?: string;
};
}
export function createEvent<T>(
type: string,
aggregateType: string,
aggregateId: string,
payload: T,
metadata?: Partial<DomainEvent['metadata']>
): DomainEvent<T> {
return {
id: crypto.randomUUID(),
type,
aggregateType,
aggregateId,
payload,
metadata: {
timestamp: new Date(),
version: 1,
...metadata,
},
};
}
Define Domain Events
export interface OrderCreatedPayload {
customerId: string;
items: Array<{
productId: string;
quantity: number;
price: number;
}>;
totalAmount: number;
shippingAddress: Address;
}
export interface OrderPaidPayload {
paymentId: string;
amount: number;
method: 'card' | 'bank' | 'wallet';
}
export interface OrderShippedPayload {
trackingNumber: string;
carrier: string;
estimatedDelivery: string;
}
export interface OrderCancelledPayload {
reason: string;
cancelledBy: string;
refundAmount?: number;
}
export type OrderEvent =
| DomainEvent<OrderCreatedPayload> & { : }
| <> & { : }
| <> & { : }
| <> & { : };
= {
:
(, , orderId, payload),
:
(, , orderId, payload),
:
(, , orderId, payload),
:
(, , orderId, payload),
};
Event Bus
In-Memory Event Bus
import { EventEmitter } from 'events';
import { DomainEvent } from './base';
type EventHandler<T = unknown> = (event: DomainEvent<T>) => Promise<void>;
class EventBus {
private emitter = new EventEmitter();
private handlers = new Map<string, EventHandler[]>();
async publish<T>(event: DomainEvent<T>): Promise<void> {
console.log(`Publishing event: ${event.type}`, event);
await this.storeEvent(event);
this.emitter.emit(event.type, event);
this.emitter.emit('*', event);
}
(: []): <> {
( event events) {
.(event);
}
}
subscribe<T>(: , : <T>): {
= () => {
{
(event);
} (error) {
.(, error);
}
};
..(eventType, wrappedHandler);
{
..(eventType, wrappedHandler);
};
}
(: ): {
.(, handler);
}
(: ): <> {
db..({
: {
: event.,
: event.,
: event.,
: event.,
: event. ,
: event. ,
: event..,
},
});
}
}
eventBus = ();
Redis-Based Event Bus
import { Redis } from 'ioredis';
import { DomainEvent } from './base';
const publisher = new Redis(process.env.REDIS_URL!);
const subscriber = new Redis(process.env.REDIS_URL!);
class RedisEventBus {
private handlers = new Map<string, Set<(event: DomainEvent) => Promise<void>>>();
constructor() {
subscriber.on('message', async (channel, message) => {
const event = JSON.parse(message) as DomainEvent;
const handlers = this.handlers.get(channel) || new Set();
for (const handler of handlers) {
try {
(event);
} (error) {
.(, error);
}
}
});
}
(: ): <> {
channel = ;
publisher.(channel, .(event));
publisher.(
,
,
,
.(event)
);
}
(: , : <>): {
channel = ;
(!..(channel)) {
..(channel, ());
subscriber.(channel);
}
..(channel)!.(handler);
{
..(channel)?.(handler);
};
}
}
eventBus = ();
Event Handlers
Handler Registration
import { eventBus } from '../events/event-bus';
import { OrderEvent } from '../events/order.events';
eventBus.subscribe<OrderCreatedPayload>('OrderCreated', async (event) => {
await emailService.send({
to: await getUserEmail(event.payload.customerId),
template: 'order-confirmation',
data: {
orderId: event.aggregateId,
items: event.payload.items,
total: event.payload.totalAmount,
},
});
});
eventBus.subscribe<OrderCreatedPayload>('OrderCreated', async (event) => {
for (const item of event.payload.items) {
await inventoryService.reserve(item.productId, item.quantity);
}
});
eventBus.subscribe<OrderPaidPayload>(, (event) => {
analytics.(, {
: event.,
: event..,
: event..,
});
});
eventBus.<>(, (event) => {
shippingService.(event.);
});
eventBus.<>(, (event) => {
order = orderRepository.(event.);
( item order.) {
inventoryService.(item., item.);
}
(event..) {
paymentService.(event., event..);
}
emailService.({
: (order.),
: ,
: {
: event.,
: event..,
},
});
});
Event Sourcing
Aggregate with Events
import { DomainEvent } from '../events/base';
import { OrderEvents, OrderCreatedPayload, OrderPaidPayload } from '../events/order.events';
interface OrderItem {
productId: string;
quantity: number;
price: number;
}
type OrderStatus = 'pending' | 'paid' | 'shipped' | 'delivered' | 'cancelled';
export class OrderAggregate {
private _id: string;
private _status: OrderStatus = 'pending';
private _items: OrderItem[] = [];
private _totalAmount: number = 0;
private _customerId: string = '';
private _version: number = 0;
private uncommittedEvents: [] = [];
() { .; }
() { .; }
() { [....]; }
() { .; }
() {
. = id || crypto.();
}
(: , : [], : ): {
order = ();
totalAmount = items.( sum + item. * item., );
order.(
.(order., {
customerId,
items,
totalAmount,
shippingAddress,
})
);
order;
}
(: , : , : | | ): {
(. !== ) {
();
}
(amount !== .) {
();
}
.(
.(., { paymentId, amount, method })
);
}
(: , : ): {
([, , ].(.)) {
();
}
refundAmount = . === ? . : ;
.(
.(., { reason, cancelledBy, refundAmount })
);
}
(: ): {
.(event);
..(event);
}
(: ): {
(event.) {
:
created = event. ;
. = created.;
. = created.;
. = created.;
. = ;
;
:
. = ;
;
:
. = ;
;
:
. = ;
;
}
.++;
}
(): [] {
events = [....];
. = [];
events;
}
(: []): {
(events. === ) {
();
}
order = (events[].);
( event events) {
order.(event);
}
order;
}
}
Event Store Repository
import { db } from '../lib/db';
import { DomainEvent } from '../events/base';
import { eventBus } from '../events/event-bus';
export class EventStoreRepository<T extends { id: string; getUncommittedEvents(): DomainEvent[] }> {
constructor(
private aggregateType: string,
private reconstruct: (events: DomainEvent[]) => T
) {}
async save(aggregate: T): Promise<void> {
const events = aggregate.getUncommittedEvents();
if (events.length === 0) return;
await db.event.createMany({
data: events.map((event) => ({
id: event.id,
type: event.type,
aggregateId: event.aggregateId,
: event.,
: event. ,
: event. ,
: event..,
})),
});
eventBus.(events);
}
(: ): <T | > {
events = db..({
: {
: id,
: .,
},
: { : },
});
(events. === ) ;
.(
events.( ({
: e.,
: e.,
: e.,
: e.,
: e.,
: e. ,
}))
);
}
(: , ?: ): <[]> {
events = db..({
: {
aggregateId,
: .,
...(fromVersion && {
: { : [], : fromVersion },
}),
},
: { : },
});
events.( ({
: e.,
: e.,
: e.,
: e.,
: e.,
: e. ,
}));
}
}
orderRepository = (
,
.
);
CQRS Pattern
Separate Command and Query
export interface CreateOrderCommand {
customerId: string;
items: Array<{ productId: string; quantity: number }>;
shippingAddress: Address;
}
export async function handleCreateOrder(command: CreateOrderCommand): Promise<string> {
const customer = await customerRepository.findById(command.customerId);
if (!customer) throw new Error('Customer not found');
const items = await Promise.all(
command.items.map(async (item) => {
const product = await productRepository.findById(item.productId);
return {
productId: item.productId,
quantity: item.quantity,
: product.,
};
})
);
order = .(
command.,
items,
command.
);
orderRepository.(order);
order.;
}
export interface OrderReadModel {
id: string;
status: string;
customerName: string;
customerEmail: string;
items: Array<{
productName: string;
quantity: number;
price: number;
}>;
totalAmount: number;
createdAt: Date;
paidAt?: Date;
shippedAt?: Date;
}
export async function getOrderById(orderId: string): Promise<OrderReadModel | null> {
return db.orderReadModel.findUnique({
where: { id: orderId },
});
}
export async function getOrdersByCustomer(customerId: string): Promise<[]> {
db..({
: { customerId },
: { : },
});
}
Read Model Projector
import { eventBus } from '../events/event-bus';
eventBus.subscribe('OrderCreated', async (event) => {
const { customerId, items, totalAmount } = event.payload;
const customer = await db.customer.findUnique({ where: { id: customerId } });
await db.orderReadModel.create({
data: {
id: event.aggregateId,
status: 'pending',
customerId,
customerName: customer.name,
customerEmail: customer.email,
items: await enrichItems(items),
totalAmount,
createdAt: event.metadata.timestamp,
},
});
});
eventBus.subscribe('OrderPaid', async (event) => {
await db.orderReadModel.update({
where: { id: event.aggregateId },
data: {
status: 'paid',
paidAt: event.metadata.timestamp,
},
});
});
eventBus.(, (event) => {
db..({
: { : event. },
: {
: ,
: event..,
: event..,
},
});
});
Best Practices
- Immutable events: Never modify stored events
- Descriptive event names: Past tense (OrderCreated, not CreateOrder)
- Include all context: Events should be self-contained
- Version events: Handle schema evolution
- Idempotent handlers: Handle duplicate events gracefully
- Separate concerns: Commands mutate, queries read
- Event versioning: Support backward compatibility
- Dead letter queue: Handle failed events
Output Checklist
Every event-driven system should include: