| name | event-driven-architecture |
| description | Event sourcing, CQRS, and message queue patterns with RabbitMQ and Kafka for distributed systems |
| category | backend |
| triggers | ["event driven","event sourcing","cqrs","message queue","rabbitmq","kafka","saga pattern"] |
Event-Driven Architecture
Implement event-driven systems with event sourcing, CQRS, and message queues. This skill covers distributed patterns for scalable, resilient applications.
Purpose
Build loosely coupled, scalable systems:
- Implement event sourcing for audit trails
- Apply CQRS for read/write optimization
- Use message queues for async processing
- Handle distributed transactions with sagas
- Ensure eventual consistency
- Build replay and recovery capabilities
Features
1. Event Sourcing
interface DomainEvent {
eventId: string;
eventType: string;
aggregateId: string;
aggregateType: string;
timestamp: Date;
version: number;
data: Record<string, any>;
metadata: {
userId?: string;
correlationId?: string;
causationId?: string;
};
}
type OrderEvent =
| { type: 'OrderCreated'; data: { customerId: string; items: OrderItem[] } }
| { type: 'OrderItemAdded'; data: { item: OrderItem } }
| { type: 'OrderItemRemoved'; data: { itemId: string } }
| { type: 'OrderSubmitted'; data: { submittedAt: Date } }
| { type: 'PaymentReceived'; data: { paymentId: string; amount: number } }
| { type: 'OrderShipped'; data: { trackingNumber: string; carrier: string } }
| { type: 'OrderDelivered'; data: { deliveredAt: Date } }
| { type: 'OrderCancelled'; data: { reason: string } };
class EventStore {
async append(
aggregateId: string,
events: DomainEvent[],
expectedVersion: number
): Promise<void> {
const currentVersion = await this.getVersion(aggregateId);
if (currentVersion !== expectedVersion) {
throw new ConcurrencyError(
`Expected version ${expectedVersion}, but found ${currentVersion}`
);
}
await db.$transaction(async (tx) => {
for (let i = 0; i < events.length; i++) {
await tx.event.create({
data: {
...events[i],
version: expectedVersion + i + 1,
},
});
}
});
for (const event of events) {
await eventBus.publish(event);
}
}
async getEvents(
aggregateId: string,
fromVersion?: number
): Promise<DomainEvent[]> {
return db.event.findMany({
where: {
aggregateId,
version: fromVersion ? { gt: fromVersion } : undefined,
},
orderBy: { version: 'asc' },
});
}
async getVersion(aggregateId: string): Promise<number> {
const lastEvent = await db.event.findFirst({
where: { aggregateId },
orderBy: { version: 'desc' },
});
return lastEvent?.version ?? 0;
}
}
class OrderAggregate {
private id: string;
private state: OrderState;
private version: number = 0;
private uncommittedEvents: OrderEvent[] = [];
static async load(eventStore: EventStore, id: string): Promise<OrderAggregate> {
const aggregate = new OrderAggregate(id);
const events = await eventStore.getEvents(id);
for (const event of events) {
aggregate.apply(event, false);
}
return aggregate;
}
create(customerId: string, items: OrderItem[]): void {
if (this.state) {
throw new Error('Order already exists');
}
this.applyChange({
type: 'OrderCreated',
data: { customerId, items },
});
}
addItem(item: OrderItem): void {
this.ensureState(['draft']);
this.applyChange({
type: 'OrderItemAdded',
data: { item },
});
}
submit(): void {
this.ensureState(['draft']);
if (this.state.items.length === 0) {
throw new Error('Cannot submit empty order');
}
this.applyChange({
type: 'OrderSubmitted',
data: { submittedAt: new Date() },
});
}
private apply(event: OrderEvent, isNew: boolean): void {
switch (event.type) {
case 'OrderCreated':
this.state = {
status: 'draft',
customerId: event.data.customerId,
items: event.data.items,
total: this.calculateTotal(event.data.items),
};
break;
case 'OrderItemAdded':
this.state.items.push(event.data.item);
this.state.total = this.calculateTotal(this.state.items);
break;
case 'OrderSubmitted':
this.state.status = 'submitted';
this.state.submittedAt = event.data.submittedAt;
break;
}
this.version++;
if (isNew) {
this.uncommittedEvents.push(event);
}
}
private applyChange(event: OrderEvent): void {
this.apply(event, true);
}
async save(eventStore: EventStore): Promise<void> {
const domainEvents = this.uncommittedEvents.map((e, i) => ({
eventId: uuid(),
eventType: e.type,
aggregateId: this.id,
aggregateType: 'Order',
timestamp: new Date(),
version: this.version - this.uncommittedEvents.length + i + 1,
data: e.data,
metadata: {},
}));
await eventStore.append(
this.id,
domainEvents,
this.version - this.uncommittedEvents.length
);
this.uncommittedEvents = [];
}
}
2. CQRS Pattern
interface Command {
type: string;
payload: any;
metadata: {
userId: string;
timestamp: Date;
correlationId: string;
};
}
class CommandBus {
private handlers = new Map<string, CommandHandler>();
register(commandType: string, handler: CommandHandler): void {
this.handlers.set(commandType, handler);
}
async dispatch(command: Command): Promise<void> {
const handler = this.handlers.get(command.type);
if (!handler) {
throw new Error(`No handler for command: ${command.type}`);
}
await handler.handle(command);
}
}
{
() {}
(: ): <> {
order = (());
order.(command.., command..);
order.(.);
}
}
{
: ;
: ;
}
{
handlers = <, >();
(: , : ): {
..(queryType, handler);
}
execute<T>(: ): <T> {
handler = ..(query.);
(!handler) {
();
}
handler.(query);
}
}
{
(: ): <> {
(event.) {
:
db..({
: {
: event.,
: event..,
: ,
: event...,
: event..,
: event.,
},
});
;
:
db..({
: { : event. },
: {
: ,
: event..,
},
});
;
:
db..({
: { : event. },
: {
: ,
: event..,
},
});
;
}
}
(): <> {
db..();
events = eventStore.();
( event events) {
.(event);
}
}
}
3. Message Queues with RabbitMQ
import amqp from 'amqplib';
class RabbitMQBroker {
private connection: amqp.Connection;
private channel: amqp.Channel;
async connect(): Promise<void> {
this.connection = await amqp.connect(process.env.RABBITMQ_URL!);
this.channel = await this.connection.createChannel();
await this.channel.assertExchange('events', 'topic', { durable: true });
await this.channel.assertExchange('commands', 'direct', { durable: true });
await this.channel.assertExchange('dlx', 'fanout', { durable: true });
}
(: , : , : ): <> {
content = .(.(message));
..(exchange, routingKey, content, {
: ,
: ,
: (),
: .(),
});
}
(
: ,
: ,
: ,
: <>
): <> {
..(queue, {
: ,
: ,
: ,
});
..(queue, exchange, routingKey);
..(queue, (msg) => {
(!msg) ;
{
content = .(msg..());
(content);
..(msg);
} (error) {
.(, error);
retryCount = (msg..?.[] || ) + ;
(retryCount < ) {
( {
..(exchange, routingKey, msg., {
...msg.,
: {
...msg..,
: retryCount,
},
});
..(msg);
}, .(, retryCount) * );
} {
..(msg, );
}
}
});
}
}
{
() {}
(: ): <> {
routingKey = ;
..(, routingKey, event);
}
}
{
() {}
(): <> {
..(
,
,
,
(event) => {
..(event);
}
);
}
}
4. Saga Pattern for Distributed Transactions
interface SagaStep {
name: string;
execute: (context: SagaContext) => Promise<void>;
compensate: (context: SagaContext) => Promise<void>;
}
class SagaOrchestrator {
private steps: SagaStep[] = [];
private executedSteps: SagaStep[] = [];
addStep(step: SagaStep): this {
this.steps.push(step);
return this;
}
async execute(context: SagaContext): Promise<void> {
try {
for (const step of this.steps) {
console.log(`Executing step: ${step.name}`);
await step.(context);
..(step);
}
} (error) {
.(, error);
.(context);
error;
}
}
(: ): <> {
( step ..()) {
{
.();
step.(context);
} (error) {
.(, error);
.(step, context, error);
}
}
}
}
createOrderSaga = ()
.({
: ,
: (ctx) => {
reservation = inventoryService.(ctx.);
ctx. = reservation.;
},
: (ctx) => {
(ctx.) {
inventoryService.(ctx.);
}
},
})
.({
: ,
: (ctx) => {
payment = paymentService.(ctx., ctx.);
ctx. = payment.;
},
: (ctx) => {
(ctx.) {
paymentService.(ctx.);
}
},
})
.({
: ,
: (ctx) => {
order = orderService.({
: ctx.,
: ctx.,
: ctx.,
: ctx.,
});
ctx. = order.;
},
: (ctx) => {
(ctx.) {
orderService.(ctx.);
}
},
})
.({
: ,
: (ctx) => {
notificationService.(ctx.);
},
: (ctx) => {
},
});
(): <> {
: = {
: command.,
: command.,
: (command.),
};
createOrderSaga.(context);
}
5. Kafka Streaming
import { Kafka, Producer, Consumer, EachMessagePayload } from 'kafkajs';
class KafkaService {
private kafka: Kafka;
private producer: Producer;
private consumers: Map<string, Consumer> = new Map();
constructor() {
this.kafka = new Kafka({
clientId: process.env.SERVICE_NAME,
brokers: (process.env.KAFKA_BROKERS || '').split(','),
});
}
async connect(): Promise<void> {
this.producer = this.kafka.producer({
idempotent: true,
maxInFlightRequests: 5,
});
await this.producer.();
}
(: , : []): <> {
..({
topic,
: messages.( ({
: m.,
: .(m.),
: m.,
: m.,
})),
});
}
(
: ,
: [],
: <>
): <> {
consumer = ..({
groupId,
: ,
: ,
});
consumer.();
consumer.({ topics, : });
consumer.({
: (payload) => {
{
(payload);
} (error) {
.(, error);
}
},
});
..(groupId, consumer);
}
(): <> {
..();
( consumer ..()) {
consumer.();
}
}
}
{
() {}
(): <> {
..(
,
[],
({ topic, partition, message }) => {
event = .(message.?.() || );
(event.) {
:
.(event);
;
:
.(event);
;
}
}
);
}
(: ): <> {
analyticsService.(event.);
..(, [{
: event.,
: {
: ,
: event.,
: event..,
},
}]);
}
}
Use Cases
1. Order Processing System
async function processOrder(orderId: string): Promise<void> {
const saga = new SagaOrchestrator()
.addStep(reserveInventoryStep)
.addStep(processPaymentStep)
.addStep(createShipmentStep)
.addStep(sendNotificationStep);
await saga.execute({ orderId });
}
2. Real-time Analytics
const orderTotalsStream = kafka.subscribe(
'analytics-aggregator',
['order-events'],
async (event) => {
await updateDailySales(event.data.total);
await updateProductMetrics(event.data.items);
}
);
Best Practices
Do's
- Design events as facts - Immutable, past-tense naming
- Implement idempotent handlers - Handle duplicates gracefully
- Plan for event versioning - Schema evolution
- Use dead letter queues - Handle failures
- Monitor queue depths - Alert on backlogs
- Test with chaos - Simulate failures
Don'ts
- Don't couple services through shared databases
- Don't ignore message ordering requirements
- Don't skip compensation logic
- Don't forget about exactly-once semantics
- Don't over-engineer for simple use cases
- Don't ignore backpressure
Related Skills
- redis - Pub/sub and caching
- real-time-systems - WebSocket integration
- backend-development - Service architecture
Reference Resources