Available for day contractsFrom 21st September I have availability for day and half day contracts. Please contact for more information.

Contact →
mikepreston.org

Event-Driven Architecture

Asynchronous communication patterns using events to decouple services and enable scalable, reactive systems.

Overview

Event-Driven Architecture (EDA) is a design pattern where services communicate through events rather than direct calls. Events represent significant occurrences or state changes that other parts of the system may react to. This approach enables loose coupling, scalability, and resilience by allowing systems to respond asynchronously to changes.

Key characteristics:

  • Loose coupling: Components don't need to know about each other directly
  • Asynchronous: Operations don't block waiting for responses
  • Scalability: Easy to add new consumers without modifying producers
  • Resilience: Failures in one component don't cascade to others
  • Auditability: Event logs provide natural audit trails
EventEventEventSubscribeSubscribeSubscribeProducer 1Message BrokerProducer 2Producer 3Consumer 1Consumer 2Consumer 3Event StoreEventEventEventSubscribeSubscribeSubscribeProducer 1Message BrokerProducer 2Producer 3Consumer 1Consumer 2Consumer 3Event Store

Event Producers and Consumers

Event Producers

Producers generate events when significant state changes or actions occur. They publish events to a message broker without knowing who will consume them.

Best practices:

  • Publish events after successful state changes (e.g., after database commit)
  • Include sufficient context in events (but avoid sensitive data)
  • Use idempotency keys to prevent duplicate events
  • Version your events from the start
  • Keep events immutable
# Python example with Kafka producer
from kafka import KafkaProducer
import json
from datetime import datetime

producer = KafkaProducer(
    bootstrap_servers=['localhost:9092'],
    value_serializer=lambda v: json.dumps(v).encode('utf-8')
)

def publish_order_created(order_id, customer_id, total):
    event = {
        'event_type': 'order.created',
        'event_id': str(uuid.uuid4()),
        'timestamp': datetime.utcnow().isoformat(),
        'version': '1.0',
        'data': {
            'order_id': order_id,
            'customer_id': customer_id,
            'total': total
        }
    }

    producer.send('order-events', value=event)
    producer.flush()  # Ensure delivery
// Go example with NATS
package main

import (
    "encoding/json"
    "github.com/nats-io/nats.go"
    "time"
)

type OrderCreatedEvent struct {
    EventType string    `json:"event_type"`
    EventID   string    `json:"event_id"`
    Timestamp time.Time `json:"timestamp"`
    Data      struct {
        OrderID    string  `json:"order_id"`
        CustomerID string  `json:"customer_id"`
        Total      float64 `json:"total"`
    } `json:"data"`
}

func publishOrderCreated(nc *nats.Conn, orderID, customerID string, total float64) error {
    event := OrderCreatedEvent{
        EventType: "order.created",
        EventID:   generateUUID(),
        Timestamp: time.Now().UTC(),
    }
    event.Data.OrderID = orderID
    event.Data.CustomerID = customerID
    event.Data.Total = total

    data, err := json.Marshal(event)
    if err != nil {
        return err
    }

    return nc.Publish("order.created", data)
}

Event Consumers

Consumers subscribe to events they're interested in and react accordingly. Multiple consumers can independently process the same event.

Consumer patterns:

  • Competing consumers: Multiple instances share load (queue pattern)
  • Fan-out: Each consumer gets all events (pub/sub pattern)
  • Event streaming: Process ordered stream of events
  • Dead letter handling: Route failed events for investigation
# Python Kafka consumer with error handling
from kafka import KafkaConsumer
import json
import logging

consumer = KafkaConsumer(
    'order-events',
    bootstrap_servers=['localhost:9092'],
    group_id='email-service',
    value_deserializer=lambda m: json.loads(m.decode('utf-8')),
    enable_auto_commit=False  # Manual commit for reliability
)

def process_order_created(event):
    """Send confirmation email when order created"""
    order_id = event['data']['order_id']
    customer_id = event['data']['customer_id']

    # Send email logic here
    send_order_confirmation(customer_id, order_id)

for message in consumer:
    try:
        event = message.value

        if event['event_type'] == 'order.created':
            process_order_created(event)

        # Commit offset only after successful processing
        consumer.commit()

    except Exception as e:
        logging.error(f"Failed to process event: {e}")
        # Could send to dead letter queue here
// JavaScript RabbitMQ consumer
const amqp = require('amqplib');

async function consumeOrderEvents() {
    const connection = await amqp.connect('amqp://localhost');
    const channel = await connection.createChannel();

    const exchange = 'order-events';
    const queue = 'email-service-queue';

    await channel.assertExchange(exchange, 'topic', { durable: true });
    await channel.assertQueue(queue, { durable: true });
    await channel.bindQueue(queue, exchange, 'order.*');

    // Prefetch: only process one message at a time
    channel.prefetch(1);

    channel.consume(queue, async (msg) => {
        if (msg !== null) {
            try {
                const event = JSON.parse(msg.content.toString());

                if (event.event_type === 'order.created') {
                    await sendOrderConfirmation(event.data);
                }

                // Acknowledge successful processing
                channel.ack(msg);

            } catch (error) {
                console.error('Processing failed:', error);
                // Reject and requeue (or send to DLX)
                channel.nack(msg, false, false);
            }
        }
    });
}

Idempotency

Consumers must handle duplicate events gracefully, as at-least-once delivery guarantees mean events may arrive multiple times.

# Idempotent consumer using event ID tracking
class IdempotentOrderProcessor:
    def __init__(self, redis_client):
        self.redis = redis_client
        self.processed_key_prefix = "processed:event:"

    def process_event(self, event):
        event_id = event['event_id']
        key = f"{self.processed_key_prefix}{event_id}"

        # Check if already processed (with TTL for cleanup)
        if self.redis.exists(key):
            logging.info(f"Event {event_id} already processed, skipping")
            return

        # Process the event
        self._handle_order_created(event['data'])

        # Mark as processed (expire after 7 days)
        self.redis.setex(key, 604800, "1")

Message Brokers

Kafka

Distributed streaming platform designed for high-throughput, fault-tolerant event streaming.

Key concepts:

  • Topics: Categories for events
  • Partitions: Ordered, immutable sequence of events within a topic
  • Consumer groups: Load balancing across multiple consumers
  • Offsets: Position in partition, managed by consumer
Consumer Group: email-serviceKafka ClusterTopic: ordersPartition 0Partition 1Partition 2Consumer 1P0, P1Consumer 2P2Consumer Group: email-serviceKafka ClusterTopic: ordersPartition 0Partition 1Partition 2Consumer 1P0, P1Consumer 2P2

Producer configuration:

from kafka import KafkaProducer

producer = KafkaProducer(
    bootstrap_servers=['kafka-1:9092', 'kafka-2:9092', 'kafka-3:9092'],

    # Durability: wait for all replicas to acknowledge
    acks='all',

    # Retries for transient failures
    retries=3,

    # Batching for throughput
    batch_size=16384,
    linger_ms=10,

    # Compression
    compression_type='snappy',

    # Enable idempotent producer (prevents duplicates)
    enable_idempotence=True,

    # Serialisation
    key_serializer=lambda k: k.encode('utf-8') if k else None,
    value_serializer=lambda v: json.dumps(v).encode('utf-8')
)

Consumer configuration:

from kafka import KafkaConsumer

consumer = KafkaConsumer(
    'order-events',
    bootstrap_servers=['kafka-1:9092', 'kafka-2:9092'],

    # Consumer group for load balancing
    group_id='inventory-service',

    # Start from earliest unprocessed message
    auto_offset_reset='earliest',

    # Manual offset management for reliability
    enable_auto_commit=False,

    # Deserialisation
    key_deserializer=lambda k: k.decode('utf-8') if k else None,
    value_deserializer=lambda v: json.loads(v.decode('utf-8')),

    # Consumer session timeout
    session_timeout_ms=30000,

    # Maximum poll interval
    max_poll_interval_ms=300000
)

Partition keys for ordering:

# Events with same key go to same partition, preserving order
producer.send(
    'order-events',
    key=str(customer_id),  # All events for customer stay ordered
    value=event
)

RabbitMQ

Feature-rich message broker supporting multiple messaging patterns with flexible routing.

Exchange types:

  • Direct: Route by exact routing key match
  • Topic: Route by pattern matching (wildcards: * = one word, # = zero or more words)
  • Fanout: Broadcast to all bound queues
  • Headers: Route based on message headers
routing_keyorder.createdorder.*order.createdorder.updatedProducerExchangetype: topicQueue: email-serviceQueue: analyticsQueue: inventoryConsumer: EmailConsumer: AnalyticsConsumer: Inventoryrouting_keyorder.createdorder.*order.createdorder.updatedProducerExchangetype: topicQueue: email-serviceQueue: analyticsQueue: inventoryConsumer: EmailConsumer: AnalyticsConsumer: Inventory

Publisher:

import pika
import json

connection = pika.BlockingConnection(
    pika.ConnectionParameters(
        host='rabbitmq.example.com',
        credentials=pika.PlainCredentials('user', 'pass')
    )
)
channel = connection.channel()

# Declare exchange
channel.exchange_declare(
    exchange='order-events',
    exchange_type='topic',
    durable=True  # Survive broker restart
)

# Publish message
event = {
    'event_type': 'order.created',
    'data': {'order_id': '12345'}
}

channel.basic_publish(
    exchange='order-events',
    routing_key='order.created',
    body=json.dumps(event),
    properties=pika.BasicProperties(
        delivery_mode=2,  # Make message persistent
        content_type='application/json',
        timestamp=int(time.time())
    )
)

connection.close()

Consumer with acknowledgements:

def callback(ch, method, properties, body):
    try:
        event = json.loads(body)
        process_event(event)

        # Acknowledge successful processing
        ch.basic_ack(delivery_tag=method.delivery_tag)

    except Exception as e:
        logging.error(f"Processing failed: {e}")

        # Reject and requeue (or route to DLX)
        ch.basic_nack(
            delivery_tag=method.delivery_tag,
            requeue=False  # Don't requeue, send to dead letter
        )

channel = connection.channel()

# Declare queue with dead letter exchange
channel.queue_declare(
    queue='email-service-queue',
    durable=True,
    arguments={
        'x-dead-letter-exchange': 'failed-events',
        'x-message-ttl': 86400000  # 24 hours
    }
)

# Bind queue to exchange with routing pattern
channel.queue_bind(
    queue='email-service-queue',
    exchange='order-events',
    routing_key='order.*'
)

# Configure consumer
channel.basic_qos(prefetch_count=1)  # Process one at a time
channel.basic_consume(
    queue='email-service-queue',
    on_message_callback=callback,
    auto_ack=False  # Manual acknowledgement
)

channel.start_consuming()

Choosing Between Kafka and RabbitMQ

Feature Kafka RabbitMQ
Use case Event streaming, high throughput Task queues, routing flexibility
Ordering Per-partition ordering Per-queue ordering
Message retention Configurable (default 7 days) Until consumed
Throughput Very high (millions/sec) High (tens of thousands/sec)
Delivery At-least-once (default) At-most-once or at-least-once
Replay Yes, consumers control offset No, once consumed it's gone
Routing Topic-partition based Flexible (direct, topic, fanout, headers)
Complexity Higher operational complexity Easier to operate
Best for Event logs, stream processing, analytics Request/reply, work queues, complex routing

Event Schemas and Contracts

Schema Design Principles

Include essential metadata:

{
  "event_id": "uuid-v4",
  "event_type": "order.created",
  "event_version": "1.0",
  "timestamp": "2024-12-05T10:30:00Z",
  "source": "order-service",
  "correlation_id": "trace-id-for-distributed-tracing",
  "causation_id": "id-of-event-that-caused-this",
  "data": {
    "order_id": "ORD-12345",
    "customer_id": "CUST-789",
    "total": 99.99,
    "items": [
      {"sku": "WIDGET-1", "quantity": 2, "price": 49.99}
    ]
  }
}

Schema Evolution

Use schema registries to manage changes over time and ensure compatibility.

Schema Registry (Confluent Schema Registry with Avro):

from confluent_kafka import avro
from confluent_kafka.avro import AvroProducer

# Define schema
value_schema_str = """
{
   "namespace": "com.example.orders",
   "type": "record",
   "name": "OrderCreated",
   "fields": [
      {"name": "order_id", "type": "string"},
      {"name": "customer_id", "type": "string"},
      {"name": "total", "type": "double"},
      {"name": "currency", "type": "string", "default": "USD"}
   ]
}
"""

value_schema = avro.loads(value_schema_str)

avroProducer = AvroProducer({
    'bootstrap.servers': 'localhost:9092',
    'schema.registry.url': 'http://localhost:8081'
}, default_value_schema=value_schema)

# Produce with schema validation
avroProducer.produce(
    topic='order-events',
    value={
        'order_id': 'ORD-12345',
        'customer_id': 'CUST-789',
        'total': 99.99,
        'currency': 'USD'
    }
)

Compatibility modes:

  • Backward: New consumers can read old data (safe to add optional fields)
  • Forward: Old consumers can read new data (safe to remove fields)
  • Full: Both backward and forward compatible
  • None: No compatibility checking

JSON Schema Validation

from jsonschema import validate, ValidationError

order_created_schema = {
    "$schema": "http://json-schema.org/draft-07/schema#",
    "type": "object",
    "required": ["event_id", "event_type", "timestamp", "data"],
    "properties": {
        "event_id": {"type": "string", "format": "uuid"},
        "event_type": {"type": "string", "enum": ["order.created"]},
        "event_version": {"type": "string", "pattern": "^\\d+\\.\\d+$"},
        "timestamp": {"type": "string", "format": "date-time"},
        "data": {
            "type": "object",
            "required": ["order_id", "customer_id", "total"],
            "properties": {
                "order_id": {"type": "string"},
                "customer_id": {"type": "string"},
                "total": {"type": "number", "minimum": 0}
            }
        }
    }
}

def publish_event(event):
    # Validate before publishing
    try:
        validate(instance=event, schema=order_created_schema)
        producer.send('order-events', value=event)
    except ValidationError as e:
        logging.error(f"Schema validation failed: {e}")
        raise

Versioning Strategies

1. Event type versioning:

# v1
event_type = "order.created.v1"

# v2 with breaking changes
event_type = "order.created.v2"

2. Inline version field:

{
  "event_type": "order.created",
  "event_version": "2.0",
  "data": { }
}

3. Consumer adaptation:

def handle_order_created(event):
    version = event.get('event_version', '1.0')

    if version == '1.0':
        return handle_v1(event)
    elif version == '2.0':
        return handle_v2(event)
    else:
        logging.warning(f"Unknown version {version}, attempting v2")
        return handle_v2(event)

Event Processing Patterns

Saga Pattern

Manage distributed transactions across multiple services using coordinated events. When a business transaction spans multiple services, sagas ensure eventual consistency.

Choreography-based saga (decentralised, event-driven):

MBShipping ServiceInventory ServicePayment ServiceOrder ServiceMBShipping ServiceInventory ServicePayment ServiceOrder ServiceCreate OrderOrderCreatedOrderCreatedReserve PaymentPaymentReservedPaymentReservedReserve InventoryInventoryReservedInventoryReservedSchedule ShipmentShipmentScheduledShipmentScheduledMark Order CompleteMBShipping ServiceInventory ServicePayment ServiceOrder ServiceMBShipping ServiceInventory ServicePayment ServiceOrder ServiceCreate OrderOrderCreatedOrderCreatedReserve PaymentPaymentReservedPaymentReservedReserve InventoryInventoryReservedInventoryReservedSchedule ShipmentShipmentScheduledShipmentScheduledMark Order Complete

Compensating transactions (handling failures):

MBInventory ServicePayment ServiceOrder ServiceMBInventory ServicePayment ServiceOrder ServiceOrderCreatedOrderCreatedReserve PaymentPaymentReservedPaymentReservedReserve Inventory (FAILS)InventoryReservationFailedInventoryReservationFailedRelease PaymentPaymentReleasedInventoryReservationFailedCancel OrderMBInventory ServicePayment ServiceOrder ServiceMBInventory ServicePayment ServiceOrder ServiceOrderCreatedOrderCreatedReserve PaymentPaymentReservedPaymentReservedReserve Inventory (FAILS)InventoryReservationFailedInventoryReservationFailedRelease PaymentPaymentReleasedInventoryReservationFailedCancel Order

Implementation example:

# Order Service - initiates saga
class OrderService:
    def create_order(self, order_data):
        # Save order in PENDING state
        order = Order.create(status='PENDING', **order_data)

        # Publish event to start saga
        event = {
            'event_type': 'order.created',
            'event_id': str(uuid.uuid4()),
            'correlation_id': order.correlation_id,
            'data': {
                'order_id': order.id,
                'customer_id': order.customer_id,
                'total': order.total
            }
        }
        publisher.publish('order.created', event)
        return order

    def handle_shipment_scheduled(self, event):
        """Final step: mark order complete"""
        order_id = event['data']['order_id']
        Order.update(order_id, status='COMPLETED')

    def handle_payment_failed(self, event):
        """Compensate: cancel order"""
        order_id = event['data']['order_id']
        Order.update(order_id, status='CANCELLED')

# Payment Service - saga participant
class PaymentService:
    def handle_order_created(self, event):
        order_id = event['data']['order_id']
        customer_id = event['data']['customer_id']
        total = event['data']['total']

        try:
            # Reserve payment
            payment = Payment.reserve(customer_id, total)

            # Publish success event
            publisher.publish('payment.reserved', {
                'event_type': 'payment.reserved',
                'correlation_id': event['correlation_id'],
                'data': {
                    'order_id': order_id,
                    'payment_id': payment.id
                }
            })

        except InsufficientFundsError:
            # Publish failure event
            publisher.publish('payment.failed', {
                'event_type': 'payment.failed',
                'correlation_id': event['correlation_id'],
                'data': {
                    'order_id': order_id,
                    'reason': 'insufficient_funds'
                }
            })

    def handle_inventory_failed(self, event):
        """Compensate: release payment"""
        order_id = event['data']['order_id']
        payment = Payment.find_by_order(order_id)
        payment.release()

        publisher.publish('payment.released', {
            'event_type': 'payment.released',
            'correlation_id': event['correlation_id'],
            'data': {'order_id': order_id}
        })

Orchestration-based saga (centralised coordinator):

class OrderSagaOrchestrator:
    def __init__(self):
        self.state_store = SagaStateStore()

    def start_order_saga(self, order_data):
        saga_id = str(uuid.uuid4())

        # Initialise saga state
        self.state_store.create(saga_id, {
            'status': 'STARTED',
            'current_step': 'RESERVE_PAYMENT',
            'order_data': order_data,
            'compensation_stack': []
        })

        # Send first command
        self.send_command('reserve_payment', {
            'saga_id': saga_id,
            'order_data': order_data
        })

    def handle_payment_reserved(self, event):
        saga_id = event['saga_id']
        state = self.state_store.get(saga_id)

        # Add compensation action
        state['compensation_stack'].append('release_payment')
        state['current_step'] = 'RESERVE_INVENTORY'
        self.state_store.update(saga_id, state)

        # Send next command
        self.send_command('reserve_inventory', {
            'saga_id': saga_id,
            'order_data': state['order_data']
        })

    def handle_inventory_failed(self, event):
        saga_id = event['saga_id']
        state = self.state_store.get(saga_id)

        # Execute compensations in reverse order
        for compensation in reversed(state['compensation_stack']):
            self.send_command(compensation, {'saga_id': saga_id})

        state['status'] = 'FAILED'
        self.state_store.update(saga_id, state)

CQRS (Command Query Responsibility Segregation)

Separate read and write models to optimise each independently. Commands change state, queries read state.

CommandsQueriesWriteEventsSubscribeUpdateUpdateReadReadClientCommand ServiceQuery ServiceWrite DBNormalisedEvent BusRead DB 1Orders ViewRead DB 2Analytics ViewCommandsQueriesWriteEventsSubscribeUpdateUpdateReadReadClientCommand ServiceQuery ServiceWrite DBNormalisedEvent BusRead DB 1Orders ViewRead DB 2Analytics View

Command side (write):

class OrderCommandService:
    def __init__(self, repository, event_publisher):
        self.repository = repository
        self.publisher = event_publisher

    def create_order(self, command):
        # Validate command
        if command.total <= 0:
            raise InvalidCommandError("Total must be positive")

        # Create aggregate
        order = Order(
            id=generate_id(),
            customer_id=command.customer_id,
            items=command.items,
            total=command.total
        )

        # Persist to write store
        self.repository.save(order)

        # Publish event for read models
        self.publisher.publish({
            'event_type': 'order.created',
            'data': {
                'order_id': order.id,
                'customer_id': order.customer_id,
                'total': order.total,
                'items': order.items,
                'created_at': order.created_at
            }
        })

        return order.id

    def update_order_status(self, command):
        order = self.repository.get(command.order_id)
        order.update_status(command.status)

        self.repository.save(order)

        self.publisher.publish({
            'event_type': 'order.status_updated',
            'data': {
                'order_id': order.id,
                'status': order.status,
                'updated_at': order.updated_at
            }
        })

Query side (read):

class OrderQueryService:
    def __init__(self, read_db):
        self.db = read_db

    # Optimised for specific query patterns
    def get_customer_orders(self, customer_id, limit=10):
        """Denormalised view optimised for customer lookups"""
        return self.db.query(
            "SELECT * FROM customer_orders_view "
            "WHERE customer_id = %s "
            "ORDER BY created_at DESC LIMIT %s",
            (customer_id, limit)
        )

    def get_order_summary(self, order_id):
        """Denormalised view with joined data"""
        return self.db.query(
            "SELECT * FROM order_summary_view WHERE order_id = %s",
            (order_id,)
        )

    # Event handler to keep read model in sync
    def handle_order_created(self, event):
        """Project event into read model"""
        self.db.execute(
            "INSERT INTO customer_orders_view "
            "(order_id, customer_id, total, status, created_at) "
            "VALUES (%s, %s, %s, %s, %s)",
            (
                event['data']['order_id'],
                event['data']['customer_id'],
                event['data']['total'],
                'PENDING',
                event['data']['created_at']
            )
        )

    def handle_order_status_updated(self, event):
        """Update read model"""
        self.db.execute(
            "UPDATE customer_orders_view "
            "SET status = %s, updated_at = %s "
            "WHERE order_id = %s",
            (
                event['data']['status'],
                event['data']['updated_at'],
                event['data']['order_id']
            )
        )

Event Sourcing (often combined with CQRS):

Store all state changes as events, then rebuild state by replaying events.

class OrderAggregate:
    def __init__(self, order_id):
        self.id = order_id
        self.status = None
        self.total = 0
        self.items = []
        self.version = 0
        self.uncommitted_events = []

    # Command handlers produce events
    def create(self, customer_id, items, total):
        event = OrderCreatedEvent(
            order_id=self.id,
            customer_id=customer_id,
            items=items,
            total=total
        )
        self.apply(event)
        self.uncommitted_events.append(event)

    def update_status(self, new_status):
        if self.status == new_status:
            return  # Idempotent

        event = OrderStatusUpdatedEvent(
            order_id=self.id,
            old_status=self.status,
            new_status=new_status
        )
        self.apply(event)
        self.uncommitted_events.append(event)

    # Event handlers update state
    def apply(self, event):
        if isinstance(event, OrderCreatedEvent):
            self.status = 'PENDING'
            self.total = event.total
            self.items = event.items
        elif isinstance(event, OrderStatusUpdatedEvent):
            self.status = event.new_status

        self.version += 1

    # Rebuild from event history
    @classmethod
    def from_history(cls, order_id, events):
        aggregate = cls(order_id)
        for event in events:
            aggregate.apply(event)
        return aggregate

class EventStore:
    def save(self, aggregate):
        # Save new events
        for event in aggregate.uncommitted_events:
            self.db.execute(
                "INSERT INTO events (aggregate_id, version, event_type, data) "
                "VALUES (%s, %s, %s, %s)",
                (aggregate.id, aggregate.version, event.event_type, event.to_json())
            )
        aggregate.uncommitted_events = []

    def load(self, aggregate_id):
        # Load all events for aggregate
        events = self.db.query(
            "SELECT * FROM events WHERE aggregate_id = %s ORDER BY version",
            (aggregate_id,)
        )
        return OrderAggregate.from_history(aggregate_id, events)

Event Notification

Simple pattern where services emit events about state changes, and other services react.

class UserService:
    def register_user(self, email, password):
        user = User.create(email=email, password_hash=hash(password))

        # Emit event
        self.publisher.publish('user.registered', {
            'user_id': user.id,
            'email': user.email,
            'registered_at': user.created_at
        })

        return user

# Multiple services react independently
class EmailService:
    def handle_user_registered(self, event):
        send_welcome_email(event['email'])

class AnalyticsService:
    def handle_user_registered(self, event):
        track_registration(event['user_id'], event['registered_at'])

class OnboardingService:
    def handle_user_registered(self, event):
        create_onboarding_workflow(event['user_id'])

Event-Carried State Transfer

Include full state in events to avoid consumers needing to query back to the producer.

# Instead of minimal event
{
    'event_type': 'product.price_changed',
    'data': {
        'product_id': 'PROD-123'
    }
}

# Include full relevant state
{
    'event_type': 'product.price_changed',
    'data': {
        'product_id': 'PROD-123',
        'sku': 'WIDGET-1',
        'name': 'Super Widget',
        'old_price': 49.99,
        'new_price': 39.99,
        'currency': 'USD',
        'category': 'widgets',
        'in_stock': True
    }
}

Benefits: Consumers have all needed data; reduces coupling Trade-offs: Larger messages; potential data duplication

Integration with Microservices

Service Communication Patterns

Asynchronous CommunicationSynchronous CommunicationAPI GatewayHTTP/RESTgRPCEventsEventsEventsSubscribeSubscribeSubscribeCommandsAPI GatewayUser ServiceProduct ServiceEvent BusOrder ServiceEmail ServiceAnalytics ServiceNotification ServiceAsynchronous CommunicationSynchronous CommunicationAPI GatewayHTTP/RESTgRPCEventsEventsEventsSubscribeSubscribeSubscribeCommandsAPI GatewayUser ServiceProduct ServiceEvent BusOrder ServiceEmail ServiceAnalytics ServiceNotification Service

When to use events vs. direct calls:

Use Events When Use Direct Calls When
Multiple services need to react Single service needs to respond
Order of operations doesn't matter Immediate response required
Process can be asynchronous User is waiting for result
Want loose coupling Need strong consistency
Fan-out to many consumers Point-to-point communication
Building audit trail Simple CRUD operations

Transactional Outbox Pattern

Ensure atomic writes to database and event publishing.

Problem: If you write to database then publish event, what happens if publishing fails?

Solution: Write database record and event to outbox table in same transaction, then reliably publish events.

Message BrokerOutbox PollerDatabaseServiceMessage BrokerOutbox PollerDatabaseServiceloop[Poll outbox]BEGIN TRANSACTIONINSERT INTO ordersINSERT INTO outboxCOMMITSELECT unpublished eventsPublish eventsMark as publishedMessage BrokerOutbox PollerDatabaseServiceMessage BrokerOutbox PollerDatabaseServiceloop[Poll outbox]BEGIN TRANSACTIONINSERT INTO ordersINSERT INTO outboxCOMMITSELECT unpublished eventsPublish eventsMark as published

Implementation:

class OrderService:
    def create_order(self, order_data):
        # Single transaction for atomicity
        with self.db.transaction():
            # Insert order
            order = self.db.execute(
                "INSERT INTO orders (id, customer_id, total, status) "
                "VALUES (%s, %s, %s, %s) RETURNING *",
                (generate_id(), order_data['customer_id'],
                 order_data['total'], 'PENDING')
            )

            # Insert event into outbox
            event = {
                'event_type': 'order.created',
                'event_id': str(uuid.uuid4()),
                'data': {
                    'order_id': order['id'],
                    'customer_id': order['customer_id'],
                    'total': order['total']
                }
            }

            self.db.execute(
                "INSERT INTO outbox (id, event_type, payload, created_at) "
                "VALUES (%s, %s, %s, NOW())",
                (event['event_id'], event['event_type'], json.dumps(event))
            )

        return order

# Separate process polls and publishes
class OutboxPublisher:
    def __init__(self, db, message_broker):
        self.db = db
        self.broker = message_broker

    def poll_and_publish(self):
        """Run this in background process/thread"""
        while True:
            # Get unpublished events
            events = self.db.query(
                "SELECT * FROM outbox WHERE published = FALSE "
                "ORDER BY created_at LIMIT 100 FOR UPDATE SKIP LOCKED"
            )

            for event_record in events:
                try:
                    # Publish to message broker
                    event = json.loads(event_record['payload'])
                    self.broker.publish(event['event_type'], event)

                    # Mark as published
                    self.db.execute(
                        "UPDATE outbox SET published = TRUE, "
                        "published_at = NOW() WHERE id = %s",
                        (event_record['id'],)
                    )
                    self.db.commit()

                except Exception as e:
                    logging.error(f"Failed to publish event: {e}")
                    self.db.rollback()

            time.sleep(1)  # Poll interval

Change Data Capture (CDC)

Alternative to outbox: capture database changes and turn them into events.

Using Debezium:

# Debezium connector configuration
{
  "name": "orders-connector",
  "config": {
    "connector.class": "io.debezium.connector.postgresql.PostgresConnector",
    "database.hostname": "postgres",
    "database.port": "5432",
    "database.user": "debezium",
    "database.password": "secret",
    "database.dbname": "orders_db",
    "database.server.name": "orders",
    "table.include.list": "public.orders,public.order_items",
    "plugin.name": "pgoutput"
  }
}

Debezium reads PostgreSQL WAL and publishes changes to Kafka automatically.

Service Mesh Integration

Use service mesh (Istio, Linkerd) alongside events for comprehensive microservices communication.

# Istio VirtualService for synchronous calls
apiVersion: networking.istio.io/v1beta1
kind: VirtualService
metadata:
  name: order-service
spec:
  hosts:
  - order-service
  http:
  - route:
    - destination:
        host: order-service
        subset: v2
      weight: 90
    - destination:
        host: order-service
        subset: v1
      weight: 10
    timeout: 5s
    retries:
      attempts: 3
      perTryTimeout: 2s

Events handle asynchronous flows, service mesh handles synchronous with retries, circuit breaking, etc.

Quick Reference

Event Design Checklist

  • [ ] Include unique event_id for idempotency
  • [ ] Include event_type and event_version
  • [ ] Include ISO 8601 timestamp
  • [ ] Include correlation_id for tracing
  • [ ] Keep events immutable
  • [ ] Include sufficient context (no callbacks)
  • [ ] Use past tense naming (e.g., order.created, not order.create)
  • [ ] Consider size (keep events reasonably small)
  • [ ] Version from the start
  • [ ] Document schema

Kafka Quick Commands

# List topics
kafka-topics --bootstrap-server localhost:9092 --list

# Create topic
kafka-topics --bootstrap-server localhost:9092 \
  --create --topic order-events \
  --partitions 3 --replication-factor 2

# Describe topic
kafka-topics --bootstrap-server localhost:9092 \
  --describe --topic order-events

# Produce message
echo '{"event_type":"test"}' | kafka-console-producer \
  --bootstrap-server localhost:9092 --topic order-events

# Consume messages
kafka-console-consumer --bootstrap-server localhost:9092 \
  --topic order-events --from-beginning

# Consumer group info
kafka-consumer-groups --bootstrap-server localhost:9092 \
  --describe --group email-service

# Reset consumer offset
kafka-consumer-groups --bootstrap-server localhost:9092 \
  --group email-service --topic order-events \
  --reset-offsets --to-earliest --execute

RabbitMQ Quick Commands

# List queues
rabbitmqadmin list queues

# Declare exchange
rabbitmqadmin declare exchange name=order-events type=topic durable=true

# Declare queue
rabbitmqadmin declare queue name=email-queue durable=true

# Bind queue to exchange
rabbitmqadmin declare binding source=order-events \
  destination=email-queue routing_key="order.*"

# Publish message
rabbitmqadmin publish exchange=order-events \
  routing_key=order.created \
  payload='{"event_type":"order.created"}'

# Get messages (consume one)
rabbitmqadmin get queue=email-queue ackmode=ack_requeue_false

# Purge queue
rabbitmqadmin purge queue name=email-queue

# List bindings
rabbitmqadmin list bindings

Common Event Types

# Lifecycle events
resource.created
resource.updated
resource.deleted

# State transition events
order.submitted
order.confirmed
order.shipped
order.delivered
order.cancelled

# Domain events
payment.processed
payment.refunded
inventory.reserved
inventory.released

# Integration events
user.registered
subscription.renewed
email.sent

# Failure events
payment.failed
order.validation_failed

Consumer Reliability Patterns

# At-least-once processing with idempotency
def process_event(event):
    event_id = event['event_id']

    # Check if already processed
    if is_processed(event_id):
        return  # Skip duplicate

    # Process event
    handle_event(event)

    # Mark as processed
    mark_processed(event_id)

    # Acknowledge to broker
    ack_message()

# Retry with exponential backoff
def process_with_retry(event, max_retries=3):
    for attempt in range(max_retries):
        try:
            handle_event(event)
            return
        except RetryableError as e:
            if attempt == max_retries - 1:
                send_to_dead_letter_queue(event)
                raise

            wait_time = 2 ** attempt  # Exponential backoff
            time.sleep(wait_time)

# Circuit breaker for downstream dependencies
circuit_breaker = CircuitBreaker(
    failure_threshold=5,
    timeout_duration=60,
    expected_exception=DownstreamError
)

@circuit_breaker
def call_downstream_service(data):
    return requests.post('http://api.example.com', json=data)

Common Issues and Solutions

Duplicate Events

Problem: Events arrive multiple times due to at-least-once delivery guarantees.

Solutions:

  1. Idempotent processing:
# Track processed event IDs
processed_events = set()

def process_event(event):
    event_id = event['event_id']

    if event_id in processed_events:
        logging.info(f"Skipping duplicate event {event_id}")
        return

    # Process event
    handle_event(event)
    processed_events.add(event_id)
  1. Database unique constraints:
CREATE TABLE processed_events (
    event_id UUID PRIMARY KEY,
    processed_at TIMESTAMP DEFAULT NOW()
);

-- Insert will fail for duplicates
INSERT INTO processed_events (event_id) VALUES (?);
  1. Natural idempotency:
# Update operations are naturally idempotent
UPDATE products SET price = 99.99 WHERE id = 'PROD-123';

Event Ordering

Problem: Events arrive out of order.

Solutions:

  1. Partition by key (Kafka):
# All events for same customer go to same partition
producer.send(
    'order-events',
    key=str(customer_id),  # Ensures ordering
    value=event
)
  1. Sequence numbers:
{
    "event_type": "order.updated",
    "sequence": 5,
    "data": { }
}
def process_event(event):
    expected_seq = get_last_sequence(event['order_id']) + 1

    if event['sequence'] < expected_seq:
        # Old event, ignore
        return

    if event['sequence'] > expected_seq:
        # Missing events, buffer or request replay
        buffer_event(event)
        return

    # Process in-order event
    handle_event(event)
    update_sequence(event['order_id'], event['sequence'])
  1. Timestamps with conflict resolution:
def process_update(event):
    current = get_current_state(event['resource_id'])

    if event['timestamp'] <= current.updated_at:
        # Stale event, ignore
        return

    # Newer event, apply update
    apply_update(event)

Lost Messages

Problem: Messages are lost due to broker failures or network issues.

Solutions:

  1. Producer acknowledgements (Kafka):
producer = KafkaProducer(
    acks='all',  # Wait for all replicas
    retries=3,
    max_in_flight_requests_per_connection=1  # Prevent reordering
)
  1. Persistent messages (RabbitMQ):
channel.basic_publish(
    exchange='order-events',
    routing_key='order.created',
    body=message,
    properties=pika.BasicProperties(
        delivery_mode=2  # Persistent
    )
)
  1. Transactional outbox pattern (see Integration section)

Consumer Lag

Problem: Consumers falling behind, events piling up.

Monitoring:

# Kafka consumer lag
kafka-consumer-groups --bootstrap-server localhost:9092 \
  --describe --group email-service

# Shows lag per partition

Solutions:

  1. Scale consumers:
# Add more consumer instances (up to partition count)
# Kafka will rebalance automatically
  1. Optimise processing:
# Batch processing
def consume_batch():
    messages = consumer.poll(timeout_ms=1000, max_records=100)

    batch = []
    for message in messages:
        batch.append(message.value)

    # Process batch efficiently
    process_batch(batch)
    consumer.commit()
  1. Increase partitions:
kafka-topics --bootstrap-server localhost:9092 \
  --alter --topic order-events --partitions 10

Poison Messages

Problem: Message that always fails processing, blocking queue.

Solutions:

  1. Dead letter queue:
MAX_RETRIES = 3

def process_with_dlq(message):
    retry_count = get_retry_count(message)

    try:
        handle_message(message)
        ack_message(message)
    except Exception as e:
        if retry_count >= MAX_RETRIES:
            # Send to DLQ for manual investigation
            send_to_dlq(message, error=str(e))
            ack_message(message)
        else:
            # Retry
            increment_retry_count(message)
            nack_message(message, requeue=True)
  1. Error handling with skip:
def process_event(event):
    try:
        validate_event(event)
        handle_event(event)
    except ValidationError as e:
        # Skip invalid events
        logging.error(f"Invalid event: {e}")
        send_to_error_topic(event, error=str(e))
        return  # Don't retry
    except RetryableError as e:
        # Retry transient errors
        raise

Schema Incompatibility

Problem: Producer sends new schema version that breaks consumers.

Solutions:

  1. Schema registry with validation:
# Enforce backward compatibility
schema_registry.register_schema(
    subject='order-created-value',
    schema=new_schema,
    compatibility='BACKWARD'  # Raises error if incompatible
)
  1. Graceful degradation:
def handle_order_created(event):
    # Handle missing fields gracefully
    customer_id = event['data'].get('customer_id')
    email = event['data'].get('email')  # New field in v2

    if email:
        send_to_email(email)
    elif customer_id:
        email = lookup_email(customer_id)
        send_to_email(email)
  1. Version-specific handlers:
handlers = {
    '1.0': handle_v1_order_created,
    '2.0': handle_v2_order_created
}

def dispatch_event(event):
    version = event.get('event_version', '1.0')
    handler = handlers.get(version)

    if handler:
        handler(event)
    else:
        logging.warning(f"No handler for version {version}")

Debugging Event Flows

Tools and techniques:

  1. Correlation IDs:
# Propagate correlation ID through entire flow
event = {
    'event_id': str(uuid.uuid4()),
    'correlation_id': correlation_id,  # From initial request
    'causation_id': causing_event_id,  # Immediate cause
    'data': { }
}
  1. Distributed tracing:
from opentelemetry import trace

tracer = trace.get_tracer(__name__)

def publish_event(event):
    with tracer.start_as_current_span("publish_event") as span:
        span.set_attribute("event.type", event['event_type'])
        span.set_attribute("event.id", event['event_id'])

        producer.send('order-events', value=event)
  1. Event logging:
import structlog

logger = structlog.get_logger()

def process_event(event):
    logger.info(
        "processing_event",
        event_type=event['event_type'],
        event_id=event['event_id'],
        correlation_id=event.get('correlation_id')
    )

    try:
        handle_event(event)
        logger.info("event_processed", event_id=event['event_id'])
    except Exception as e:
        logger.error(
            "event_processing_failed",
            event_id=event['event_id'],
            error=str(e)
        )
        raise
  1. Event replay for debugging:
# Kafka: reset offset to replay events
kafka-consumer-groups --bootstrap-server localhost:9092 \
  --group debug-consumer \
  --topic order-events \
  --reset-offsets --to-datetime 2024-12-05T10:00:00.000 \
  --execute

Message Broker Outages

Problem: Message broker becomes unavailable.

Solutions:

  1. Circuit breaker pattern:
from pybreaker import CircuitBreaker

broker_circuit = CircuitBreaker(
    fail_max=5,
    timeout_duration=60
)

@broker_circuit
def publish_event(event):
    producer.send('order-events', value=event)

try:
    publish_event(event)
except CircuitBreakerError:
    # Broker is down, fallback strategy
    cache_event_for_later(event)
  1. Local queue with retry:
class ResilientPublisher:
    def __init__(self):
        self.pending_queue = Queue()
        self.retry_thread = Thread(target=self.retry_loop)
        self.retry_thread.start()

    def publish(self, event):
        try:
            self.broker.send(event)
        except BrokerUnavailable:
            # Queue locally
            self.pending_queue.put(event)

    def retry_loop(self):
        while True:
            if not self.pending_queue.empty():
                event = self.pending_queue.get()
                try:
                    self.broker.send(event)
                except:
                    # Put back in queue
                    self.pending_queue.put(event)
                    time.sleep(10)
  1. Clustering and replication:
# Kafka with replication
kafka-topics --create \
  --topic order-events \
  --replication-factor 3 \
  --min-insync-replicas 2

Performance Optimisation

Batch processing:

# Consumer batching
records = consumer.poll(timeout_ms=1000, max_records=100)

# Process in batch
batch = [record.value for record in records]
bulk_insert_to_db(batch)

consumer.commit()

Compression:

# Kafka producer compression
producer = KafkaProducer(
    compression_type='snappy',  # or 'gzip', 'lz4', 'zstd'
    linger_ms=10,  # Wait to batch messages
    batch_size=32768
)

Parallel processing:

from concurrent.futures import ThreadPoolExecutor

def consume_with_parallel_processing():
    executor = ThreadPoolExecutor(max_workers=10)

    for message in consumer:
        # Submit to thread pool
        future = executor.submit(process_event, message.value)
        future.add_done_callback(lambda f: consumer.commit())