smaple.tr
event-driven architecture

Event-Driven Architecture: Message Queues, CQRS, Event Sourcing, and Kafka vs RabbitMQ [2026]

Mehmet Kurtipek
December 18, 2025
13 min read
event-driven architecture
Apache Kafka
RabbitMQ
CQRS
event sourcing
message queues
saga pattern

87% of distributed system failures trace back to synchronous coupling: service A calls service B, service B is slow or down, and service A fails or blocks. Event-driven architecture breaks this coupling by replacing direct service calls with asynchronous event publication. Services communicate through events rather than through requests, which means producers and consumers can scale, fail, and deploy independently.

This guide covers the full event-driven architecture stack: the fundamental patterns (message queues vs event streaming, pub/sub vs work queues), Kafka and RabbitMQ in depth, CQRS and event sourcing integration, saga patterns for distributed transactions, schema management, and fault handling. By the end, you will have a clear decision framework for when and how to apply each pattern.

Event-Driven Architecture: Core Concepts

Events vs Commands vs Queries

Three interaction types exist in distributed systems. Understanding the difference is prerequisite to designing event-driven architecture correctly.

Commands express intent: "ProcessPayment", "ShipOrder", "CreateUser". A command expects an outcome. It is addressed to a specific service. It may fail. Commands are typically synchronous (REST, gRPC) when you need immediate confirmation, or asynchronous when you can accept eventual execution.

Events record facts: "PaymentProcessed", "OrderShipped", "UserCreated". An event records something that happened. It is not addressed to any consumer — it is broadcast to whoever cares. An event cannot be rejected because it has already happened. Events are always asynchronous.

Queries request data: "GetOrder", "ListProducts". Queries are read-only, addressed to a specific service, and typically synchronous.

Event-driven architecture is primarily about replacing synchronous service-to-service calls with event publication when the interaction is recording a fact rather than requesting an action. Payment service does not call Notification service to send a confirmation email — it publishes a payment.processed event, and Notification service consumes it.

Message Queues vs Event Streaming

These two patterns solve different problems and should not be treated as interchangeable.

Message queues (RabbitMQ, AWS SQS, Azure Service Bus): each message is consumed by exactly one consumer, then deleted. Messages represent units of work — tasks to be processed once. Use cases: job processing, command dispatch, workflow coordination.

Event streaming (Apache Kafka, AWS Kinesis, Google Pub/Sub): events are retained for a configurable period and can be consumed by multiple independent consumers. Each consumer group maintains its own read offset. Events can be replayed from any point in history. Use cases: event sourcing, audit logs, real-time analytics, change data capture, data pipeline integration.

The wrong pattern choice creates painful problems. Using a message queue for event sourcing means events are deleted after consumption, destroying the audit trail. Using event streaming for simple job queues creates unnecessary offset management and retention overhead.

Apache Kafka In Depth

Kafka is a distributed event streaming platform. Its architecture is fundamentally different from traditional message brokers:

Topics, Partitions, and Consumer Groups

Topics are logical categories for events. Events in a topic are append-only and immutable. A topic is not a queue — multiple consumer groups can read from the same topic independently, each maintaining their own position.

Partitions are the unit of parallelism. Each topic is divided into partitions. Events with the same key are always written to the same partition, preserving ordering per key. Consumers within a consumer group are assigned partitions — each partition is consumed by exactly one consumer in the group. Maximum parallelism = number of partitions.

Consumer groups allow multiple independent consumers to read the same topic. Group A (analytics service) might be 6 hours behind Group B (notification service) on the same orders topic. Neither affects the other.

Topic: orders (4 partitions)
                   Partition 0 → Consumer A (Group: analytics)
orders ──────────── Partition 1 → Consumer B (Group: analytics)
                   Partition 2 → Consumer C (Group: analytics)
                   Partition 3 → Consumer D (Group: analytics)

Same topic, different group:
                   Partition 0 → Consumer X (Group: notifications)
orders ──────────── Partition 1 → Consumer Y (Group: notifications)
                   Partition 2 → Consumer Z (Group: notifications)
                   Partition 3 → Consumer W (Group: notifications)

Retention and Log Compaction

Kafka retains events for a configurable time period (default 7 days) or size limit. After retention expires, events are deleted. This makes Kafka appropriate for event streaming, not permanent event storage — if you need indefinite retention, archive to object storage (S3, GCS) before expiry.

Log compaction is an alternative retention policy: instead of time-based deletion, Kafka retains only the most recent event per key. This is appropriate for state snapshots — the compacted topic represents current state for each entity rather than full history.

Offset Management

Each consumer tracks which events it has processed via an offset (sequential position within a partition). Offsets are committed to Kafka (stored in the __consumer_offsets internal topic). On consumer restart, it resumes from its last committed offset.

Offset commit strategy:

  • Auto-commit (default) — Kafka commits the offset at a configurable interval. Risk: if the consumer crashes after processing but before the auto-commit interval, events are re-processed. Fine for idempotent consumers.
  • Manual commit — the consumer commits the offset explicitly after confirming successful processing. Provides at-least-once semantics with explicit control.

KRaft Mode (Kafka Without ZooKeeper)

Prior to Kafka 3.3, Kafka required ZooKeeper for metadata management and leader election. KRaft mode (Kafka Raft metadata mode) eliminates this dependency. In 2024, KRaft reached production-ready status. New Kafka deployments should use KRaft — it simplifies operations significantly.

RabbitMQ In Depth

RabbitMQ is an AMQP message broker. Its routing model is fundamentally different from Kafka's partition-based model.

Exchange Types and Routing

RabbitMQ's flexibility comes from exchange types — routing strategies that determine which queues receive a message.

Direct exchange — routes to queues with an exactly matching routing key. A message with key orders.created goes to the queue bound with key orders.created. Simple, predictable.

Fanout exchange — ignores routing keys entirely. Routes to every queue bound to the exchange. Use for broadcast scenarios: audit logging, cache invalidation, metrics collection.

Topic exchange — routes by pattern matching on routing keys. * matches exactly one word, # matches zero or more words. A message with key orders.us.created would match bindings orders.*.created and orders.#. Enables flexible subscription hierarchies.

Headers exchange — routes based on message header attributes rather than routing key. Rarely needed; most routing requirements are handled by direct or topic exchanges.

Queue Durability and Message Persistence

RabbitMQ queues are ephemeral by default — a broker restart deletes them. For production queues:

channel.queue_declare(
    queue='orders',
    durable=True,          # Queue survives broker restart
)

channel.basic_publish(
    exchange='',
    routing_key='orders',
    body=json.dumps(order),
    properties=pika.BasicProperties(
        delivery_mode=2,   # Persistent message (written to disk)
    )
)

Both durable=True and delivery_mode=2 are required for message durability. durable alone survives broker restart but not individual message crashes.

Quorum Queues vs Classic Queues

RabbitMQ 3.8+ introduces quorum queues as the recommended queue type for production. Classic mirrored queues had documented reliability problems (data loss during network partitions). Quorum queues use Raft consensus for replication, providing stronger durability guarantees and automatic leader election.

New production deployments should use quorum queues. Classic queues remain available for compatibility but are not recommended for critical workloads.

Dead Letter Exchange

Messages that cannot be delivered (TTL expired, queue length limit exceeded, or explicitly rejected) route to the Dead Letter Exchange (DLX). From the DLX, they can be routed to a dead letter queue for inspection and manual retry.

Configure DLX on queue declaration:

channel.queue_declare(
    queue='orders',
    durable=True,
    arguments={
        'x-dead-letter-exchange': 'orders.dlx',
        'x-dead-letter-routing-key': 'orders.dead',
        'x-message-ttl': 30000,  # 30 second TTL
    }
)

Without a DLX, failed messages are silently dropped. Always configure DLX for production queues.

Kafka vs RabbitMQ: Decision Guide

Criterion Choose Kafka Choose RabbitMQ
Event replay needed Yes No
Multiple independent consumers Yes No
Ordering guarantees required Per-key ordering Queue-level ordering
Message volume Millions/second Thousands/second
Retention period Days/weeks Process immediately
Routing complexity Simple (topic-based) Complex (exchange routing)
Message size Small (< 1 MB) Any size
Cloud managed option MSK, Confluent AmazonMQ, CloudAMQP

Kafka is the right choice when you need event sourcing, audit trails, multiple consumer groups reading the same data, or very high throughput. RabbitMQ is the right choice when you need complex routing, immediate delivery confirmation, and you do not need replay.

CQRS and Event Sourcing

CQRS: Separating Reads and Writes

Command Query Responsibility Segregation (CQRS) uses separate models for reading and writing data. The write model accepts commands, validates business rules, and persists changes. The write model publishes events when state changes. The read model consumes these events and builds optimized read representations (materialized views).

A product catalog in CQRS:

  • Write model: validates product data, applies business rules (discount cannot exceed 50%), persists to a normalized relational database
  • Events published: ProductCreated, PriceUpdated, ProductDiscontinued
  • Read model (search): consumes events and maintains a search index (Elasticsearch) with denormalized product data for fast querying

The read model is eventually consistent — there is a lag (typically milliseconds) between a write event and the read model reflecting it. This is acceptable for most use cases but requires explicit design for cases where strong consistency is required (e.g., preventing double-booking).

Event Sourcing: State as History

Event sourcing stores state changes as an immutable sequence of events rather than current state. The current state is derived by replaying events.

Traditional persistence: UPDATE orders SET status = 'shipped' WHERE id = 1

Event sourcing: append {type: "OrderShipped", orderId: 1, timestamp: ..., carrier: "FedEx", trackingNumber: "..."}

Benefits: complete audit trail (no data loss, every state change is recorded), temporal queries (reconstruct state at any point in history), event replay (rebuild read models from scratch).

Costs: more complex state reconstruction, requires snapshot mechanism for large event streams, debugging requires understanding event history rather than querying current state.

Event sourcing is often combined with CQRS: the command side uses event sourcing for the write model; the query side builds materialized views from the event stream. This combination is powerful but architecturally complex — it is appropriate for domains where audit trails and temporal queries are business requirements, not for every service.

Saga Pattern: Distributed Transactions

When a business process spans multiple services (each with its own database), ACID transactions are not available. The saga pattern coordinates multi-step processes with compensating transactions on failure.

Choreography vs Orchestration

Choreography: each service reacts to events and produces new events. No central coordinator. An order.created event causes Payment service to publish payment.processed, which causes Inventory service to publish inventory.reserved, and so on. Simple for happy paths; difficult to visualize and debug when failures cascade.

Orchestration: a central saga coordinator sends commands to each service and waits for responses. The coordinator contains the process logic and handles failures by issuing compensating commands. More visible, easier to debug, but creates a central point of coordination that must be resilient.

Compensating Transactions

Each step in a saga must have a compensating transaction — an action that undoes its effect on failure.

Order Placement Saga:
  Step 1: Reserve inventory → Compensate: Release inventory reservation
  Step 2: Process payment  → Compensate: Refund payment
  Step 3: Create shipment  → Compensate: Cancel shipment
  Step 4: Notify customer  → Compensate: (none, notification cannot be unsent)

Compensating transactions are not always possible (notifications sent, emails delivered). Design sagas to put irreversible steps last.

Idempotency and Exactly-Once Semantics

The At-Least-Once Problem

Distributed messaging systems guarantee at-least-once delivery — if a message might be lost, the system retransmits it. This means consumers must handle duplicate delivery.

A consumer that charges a customer's payment card cannot process the same message twice. Idempotency key pattern:

def process_payment(event):
    idempotency_key = f"payment:{event.order_id}:{event.amount}"

    # Check if already processed
    if redis.get(idempotency_key):
        return  # Already processed, skip

    # Process payment
    payment_result = stripe.charge(event.amount, event.card_token)

    # Mark as processed (with TTL to prevent unbounded growth)
    redis.set(idempotency_key, "processed", ex=86400)

    return payment_result

Kafka Exactly-Once Semantics

Kafka supports exactly-once semantics via transactions: the producer marks a sequence of writes as a transaction, and the consumer reads only committed transactions. This prevents both message loss and duplicate processing within the Kafka ecosystem.

Exactly-once comes with performance overhead and adds complexity to both producer and consumer code. Use it when the business impact of duplicate processing is significant (financial transactions, inventory updates). For most event-driven systems, at-least-once with idempotent consumers is simpler and sufficient.

Schema Management and Event Versioning

Events must evolve over time. A payment.processed event today might need additional fields (currency, exchange rate) in six months. Producers and consumers need to evolve independently without breaking each other.

Schema Registry

The Confluent Schema Registry (and compatible alternatives) stores schemas centrally and enforces compatibility:

  • Backward compatible: new schema can read old data (add optional fields only)
  • Forward compatible: old schema can read new data (consumers can ignore unknown fields)
  • Full compatible: both backward and forward

New fields should always have defaults. Removing fields or changing field types requires a new schema version with an incompatibility notice and a migration plan for consumers.

Event Versioning Strategies

Field addition (preferred): add new optional fields with defaults. Existing consumers ignore unknown fields. Zero coordination required.

New event type: for significant structural changes, publish a new event type (payment.processed.v2) and maintain the old type until all consumers have migrated.

Schema evolution with registry: register schema changes against the registry; the registry validates compatibility automatically.

Fault Handling and Observability

Retry Strategy

Transient failures (network timeout, downstream service unavailable) should retry with exponential backoff and jitter:

import random
import time

def process_with_retry(event, max_attempts=5):
    for attempt in range(max_attempts):
        try:
            return process_event(event)
        except TransientError as e:
            if attempt == max_attempts - 1:
                raise
            # Exponential backoff with jitter
            delay = (2 ** attempt) + random.uniform(0, 1)
            time.sleep(min(delay, 30))  # Cap at 30 seconds

Non-transient failures (invalid data, business rule violation) should not retry — they route directly to the dead letter queue.

Consumer Lag Monitoring

Consumer lag — the difference between the latest event offset and the consumer's committed offset — indicates processing health. Rising lag means consumers are falling behind. In Kafka, monitor per-consumer-group, per-topic, per-partition lag.

Alert thresholds: lag growing for more than 5 minutes at a rate faster than production rate indicates a consumer is stuck or undersized. Lag that exceeds the retention period will cause event loss when Kafka deletes old events before they are consumed.

Distributed Tracing

Event-driven systems make request tracing non-trivial — a user action generates an event that triggers downstream processing across multiple services. Propagate trace context through event headers:

def publish_event(event_type, payload):
    trace_context = get_current_trace_context()  # OpenTelemetry

    producer.produce(
        topic='orders',
        value=json.dumps(payload),
        headers={
            'traceparent': trace_context.traceparent,
            'tracestate': trace_context.tracestate,
        }
    )

Consumers extract the trace context and continue the trace span. This allows tracing a user action from HTTP request through event publication through downstream processing in a single trace.

Related Articles

August 11, 2026

MLOps Guide: Taking Machine Learning Models to Production [2026]

87% of machine learning models built by data science teams never reach production. The models work — they pass cross-validation, they score well on holdout sets, they demonstrate genuine predictive value. The problem is not the modeling. The problem is everything that happens between a notebook experiment and a reliable, monitored, production system. MLOps is the discipline that closes that gap. This guide covers the full MLOps stack: maturity levels, tooling choices (MLflow, DVC, Kubeflow

Read More
August 10, 2026

LLM Fine-Tuning Guide: Custom Model Training with LoRA and QLoRA [2026]

General-purpose LLMs are impressive. They can write code, summarize documents, answer questions, and translate between languages with reasonable accuracy. But "reasonable" is not good enough when your application requires consistent output format, domain-specific terminology, a particular tone, or behavior that the base model was never trained to exhibit. That gap is where fine-tuning matters. Fine-tuning updates a model's weights on your specific data, changing how the model behaves — not

Read More
August 9, 2026

Computer Vision Applications: Object Detection, OCR, and Industrial AI [2026]

Computer vision has moved well past the research phase. The models are trained, the frameworks are mature, the hardware is accessible, and the use cases are generating measurable returns. What was a specialized capability requiring deep expertise in 2018 is now deployable infrastructure — if you know which component to reach for and where the real complexity lives. This guide covers computer vision applications across industrial, medical, logistics, and document processing domains. It expl

Read More