jnachi
Learning Hub
Enterprise Integration9 min readAdvanced

Event-Driven Integration with Apache Kafka & Event Meshes

Design distributed event-driven integration architectures using Apache Kafka topics, partitions, consumer groups, Schema Registry, and Change Data Capture (CDC).

Works with:Apache KafkaConfluent Schema RegistryKafka ConnectDebezium CDC

Key Takeaways

  • Apache Kafka is a distributed, horizontally partitioned, immutable append-only commit log built for massive streaming throughput
  • Topics are divided into Partitions; the partition key guarantees strict ordering within a single partition while enabling parallel consumption across consumer groups
  • Confluent Schema Registry enforces Avro/Protobuf schema compatibility (backward, forward, full) to prevent breaking contract changes
  • Change Data Capture (CDC / Debezium) streams real-time row-level database changes directly into event streams without polling

The Diagnostic Context

Traditional message brokers delete messages as soon as they are consumed. Apache Kafka fundamentally changed enterprise integration by treating events as an immutable, persistent, replayable stream of state changes across the entire enterprise.

The Core Technique

Kafka Architecture: Topics, Partitions & Consumer Groups

DIAGRAM / WORKFLOW
graph TD
    subgraph Producers["Event Producers"]
        P1["Order Service"]
        P2["Mobile Gateway"]
    end

    subgraph KafkaCluster["Kafka Cluster: Topic 'order-events'"]
        subgraph Part0["Partition 0 (Key Hash: 0)"]
            P0_0["[Offset 0]"] --- P0_1["[Offset 1]"] --- P0_2["[Offset 2]"]
        end
        subgraph Part1["Partition 1 (Key Hash: 1)"]
            P1_0["[Offset 0]"] --- P1_1["[Offset 1]"] --- P1_2["[Offset 2]"]
        end
    end

    subgraph GroupA["Consumer Group: BillingServiceGroup"]
        C1["Consumer Instance 1 (Reads Part 0)"]
        C2["Consumer Instance 2 (Reads Part 1)"]
    end

    subgraph GroupB["Consumer Group: AnalyticsGroup"]
        C3["Consumer Instance 3 (Reads Both Partitions)"]
    end

    P1 -->|Key = 'CUST-101'| Part0
    P2 -->|Key = 'CUST-202'| Part1

    Part0 --> C1
    Part1 --> C2

    Part0 --> C3
    Part1 --> C3

Core Kafka Architecture Principles

  1. Partitions & Ordering:

    • Kafka guarantees strict FIFO ordering only within a single partition.
    • By supplying a
      CODE / PROMPT
      recordKey
      (e.g.,
      CODE / PROMPT
      customerId
      ), all events for that customer hash to the same partition, guaranteeing ordered state processing.
  2. Consumer Groups & Parallel Scalability:

    • Each partition in a topic is consumed by exactly one consumer instance within a consumer group.
    • If a topic has 10 partitions, a consumer group can scale up to 10 parallel consumer instances.
  3. Schema Registry & Evolution Rules:

    • Producers and consumers share schemas stored in the Confluent Schema Registry (Apache Avro or Protobuf).
    • Backward Compatibility: Consumers with new schemas can read events produced by older schemas.
    • Full Compatibility: Old and new schemas can read events produced by either version, allowing independent deployments without downtime.
  4. Change Data Capture (CDC):

    • Tools like Debezium tail database write-ahead logs (WAL in Postgres, Redo log in Oracle) to stream row-level INSERT/UPDATE/DELETE events directly into Kafka without modifying application code.
5-Minute Activation Challenge

Try This Right Now

Design an event envelope schema in Apache Avro: Define fields for `eventId` (UUID), `timestamp` (long), `eventType` (string: "ORDER_CREATED"), and a `payload` record containing customer and line items.

Tip: Knowledge only becomes capability once you run the prompt yourself.

Comprehension Check

Test Your Instincts (3 Questions)

1

How does Apache Kafka guarantee strict sequential ordering of events for a specific customer across distributed consumers?

2

What is the role of the Confluent Schema Registry in an event-driven Kafka architecture?

3

What technology reads relational database transaction logs (e.g., Postgres WAL, Oracle Redo Log) to stream row changes directly into Kafka in real time without querying SQL tables?