Event Streaming with Kafka: Partitions, Consumer Groups, Dead Letter Queues and Production Config

Key takeaways

Kafka stores ordered, replayable logs split into partitions, so retries and ordering work differently than in a simple queue. The post runs Kafka locally with Docker Compose, shows how keys map to partitions and consumer groups share them, and builds an order pipeline with a dead letter topic.

Why Kafka?

Most message queues delete a message once a consumer reads it. That works fine for simple task queues, but it creates problems when you need multiple independent services to react to the same event, or when you need to replay past events for debugging, analytics, or onboarding a new service.

Kafka solves this by treating messages as a persistent, ordered log. Messages are written to disk and retained for a configurable period (default: 7 days). Any number of consumer groups can read the same topic independently, each tracking its own position.

Traditional queue (RabbitMQ):
  Producer → Queue → Consumer (message deleted after consumption)
  One consumer per message
  No replay of past messages

Kafka:
  Producer → Topic (partitioned, stored on disk) → Consumer Group 1
                                                  → Consumer Group 2
                                                  → Consumer Group 3
  Multiple independent consumers
  Replay any message from any offset
  Retain messages for days/weeks
  Millions of messages per second

This architecture means you can add a new analytics service six months later and have it replay all historical events — without touching the producer or existing consumers. That’s the fundamental shift Kafka enables.

Kafka wins for:

  • Event sourcing (immutable event log)
  • Audit trails (every event persisted)
  • Stream processing (real-time analytics)
  • Decoupling microservices (producer doesn’t know consumers)
  • High-throughput pipelines (logs, metrics, clickstreams)

Core Concepts

Understanding four terms unlocks most of Kafka: topic, partition, offset, and consumer group.

A topic is a named category for events — think of it like a database table, but append-only. You write orders to the orders topic, user events to user-events, and so on.

A partition is how Kafka scales. Each topic is split into N partitions, and each partition is an independent ordered log. Messages within a partition are strictly ordered; across partitions they are not. More partitions means more parallelism — but you can’t reduce partitions later without recreating the topic.

An offset is a message’s position within a partition. Consumers track which offset they’ve processed and commit it to Kafka. If a consumer restarts, it resumes from the last committed offset.

A consumer group is a set of consumers that collectively read a topic. Kafka assigns each partition to exactly one consumer in the group — so throughput scales with partition count. Multiple consumer groups read the same topic independently, each maintaining their own offset.

Topic:        Named stream of records (like a database table for events)
Partition:    Topic is split into partitions — unit of parallelism
              Messages within a partition are ordered
              Messages across partitions are NOT ordered

Offset:       Position of a message within a partition
              Consumers track their offset — where they left off

Producer:     Writes messages to topics
Consumer:     Reads messages from topics
Consumer Group: Set of consumers that together read a topic
              Each partition assigned to exactly one consumer in the group
Broker:       A Kafka server instance
Cluster:      Multiple brokers (typically 3+ for production)
Topic: "orders"  (3 partitions)

Partition 0: [offset 0: order#1] [offset 1: order#4] [offset 2: order#7]
Partition 1: [offset 0: order#2] [offset 1: order#5] [offset 2: order#8]
Partition 2: [offset 0: order#3] [offset 1: order#6] [offset 2: order#9]

Consumer Group A (3 consumers = 1 consumer per partition):
  Consumer A1 → Partition 0
  Consumer A2 → Partition 1
  Consumer A3 → Partition 2

Consumer Group B (1 consumer):
  Consumer B1 → Partition 0, 1, 2 (reads all partitions independently)

The partition count is the decision people most often regret, in both directions. It caps parallelism within one consumer group — group A cannot usefully run a fourth consumer on this three-partition topic — so it has to be chosen for the throughput you expect, not the one you have today. Partitions can be added later, but doing so changes which partition each key hashes to, so events for the same order can land on a different partition after the change and lose their relative ordering during the transition. Very high counts have costs too: more open files and replication traffic on brokers, and longer rebalances. For most services, a modest count (6–12) with room to grow is a reasonable starting point, set explicitly when creating the topic rather than inherited from broker defaults.


Setup

# Docker Compose (local development)
# docker-compose.yml
services:
  zookeeper:
    image: confluentinc/cp-zookeeper:7.5.0
    environment:
      ZOOKEEPER_CLIENT_PORT: 2181

  kafka:
    image: confluentinc/cp-kafka:7.5.0
    depends_on: [zookeeper]
    ports:
      - "9092:9092"
    environment:
      KAFKA_BROKER_ID: 1
      KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
      KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092
      KAFKA_AUTO_CREATE_TOPICS_ENABLE: "true"
      KAFKA_DEFAULT_REPLICATION_FACTOR: 1
      KAFKA_NUM_PARTITIONS: 3
docker compose up -d

# Install Node.js client
npm install kafkajs

This Compose file uses the classic ZooKeeper-based setup that most existing tutorials show. Kafka has since moved cluster metadata into Kafka itself (KRaft mode), and Kafka 4.0 removed ZooKeeper support entirely, so for new clusters use a KRaft configuration — both the apache/kafka and newer Confluent images can run a single-node KRaft broker without a separate ZooKeeper container. The setting that trips people up in either mode is KAFKA_ADVERTISED_LISTENERS. It is the address the broker tells clients to use after the initial connection. localhost:9092 works for code running on your host machine, but a Node service running in another container connects to kafka:9092, receives localhost:9092 in the metadata, and then fails with connection errors pointing at its own container. The usual fix is two listeners — one advertised as kafka:29092 for the Docker network and one as localhost:9092 for the host.

A note on the client library: kafkajs is pure JavaScript and easy to start with, but its release activity has slowed considerably in recent years. Confluent maintains @confluentinc/kafka-javascript, based on librdkafka, which offers a kafkajs-compatible API for migration. The concepts and most of the code in this post apply to both.


Producer

The producer is responsible for writing messages to topics. The most important decision a producer makes is which partition to write to, because partition determines ordering and which consumer will process the message.

By default, kafkajs routes messages with a key to a consistent partition (same key always goes to the same partition) and distributes keyless messages round-robin. This means if you use orderId as the key, all events for a given order are guaranteed to arrive in order at the same consumer.

// src/producer.ts
import { Kafka, Partitioners } from 'kafkajs';

const kafka = new Kafka({
  clientId: 'order-service',
  brokers: ['localhost:9092'],
  // For production: retry configuration
  retry: {
    initialRetryTime: 100,
    retries: 8,
  },
});

const producer = kafka.producer({
  createPartitioner: Partitioners.LegacyPartitioner,
});

await producer.connect();

// Send a single message
await producer.send({
  topic: 'orders',
  messages: [
    {
      key: 'order-123',         // Optional: determines partition
      value: JSON.stringify({
        orderId: 'order-123',
        userId: 'user-456',
        items: [{ productId: 'p1', quantity: 2 }],
        total: 59.98,
        timestamp: Date.now(),
      }),
      headers: {
        'event-type': 'order.created',
        'correlation-id': 'req-789',
      },
    },
  ],
});

// Send multiple messages (batch — more efficient)
await producer.send({
  topic: 'orders',
  messages: events.map(event => ({
    key: event.orderId,
    value: JSON.stringify(event),
  })),
});

await producer.disconnect();

The createPartitioner: Partitioners.LegacyPartitioner line exists for compatibility. kafkajs 2.0 changed its default partitioner to match the Java client’s key hashing, and prints a warning at startup until you choose one explicitly. For a new system, use Partitioners.DefaultPartitioner so that producers written in Java, Go, or Python send a given key to the same partition as your Node producer; choose the legacy one only if existing data was partitioned with it. Mixing partitioners across producers of the same topic silently breaks the “same key, same partition” guarantee.

The message value is just bytes to Kafka — here a JSON string. JSON is easy to debug, but it has no enforced schema, so a producer that renames total to amount breaks consumers at runtime. Teams that share topics across services often use Avro or Protobuf with a schema registry, which rejects incompatible changes at publish time. Headers are the right place for metadata such as event type and correlation id, because consumers can route on them without parsing the body. Finally, connect() once at startup and reuse the producer: creating and disconnecting a producer per request, as a quick script like this might suggest, adds a full connection handshake to every send.

Partitioning Strategy

// Key-based partitioning: same key — same partition (ordering guaranteed)
// Useful: all events for a user go to same partition — ordered per user
{ key: userId, value: JSON.stringify(event) }

// Round-robin: no key — distributed across partitions
{ key: null, value: JSON.stringify(event) }

// Custom partitioner
const producer = kafka.producer({
  createPartitioner: () => ({ message, partitionMetadata }) => {
    // Route high-priority orders to partition 0
    const order = JSON.parse(message.value?.toString() ?? '{}');
    if (order.priority === 'high') return 0;
    // Others: round-robin
    return Math.floor(Math.random() * partitionMetadata.length);
  },
});

Key choice is a trade-off between ordering and balance. The key defines the unit of ordering: all events with the same key are ordered relative to each other, and nothing else is. Keying by userId orders everything a user does but concentrates heavy users on one partition; keying by orderId spreads load evenly but no longer orders a user’s separate orders. A skewed key — a single tenant producing half the traffic — creates a “hot partition” that one consumer must handle alone while others sit idle, and consumer lag on that partition grows no matter how many instances you add. The custom partitioner above illustrates a related trap: routing all high-priority orders to partition 0 guarantees that partition becomes hot, and the random fallback for other orders discards per-key ordering entirely. Priority is usually better expressed as a separate topic with its own consumers.


Consumer

Consumers read from topics by subscribing and running a message handler. The groupId is critical: all instances of your service should share the same group ID so Kafka distributes partitions among them. If you run three instances with the same group ID on a topic with six partitions, each instance handles two partitions — automatic horizontal scaling.

When a consumer commits an offset, it’s telling Kafka “I’ve successfully processed everything up to here.” If the consumer restarts, it resumes from that committed offset. This is why committing after successful processing matters: committing before means a crash will skip messages.

// src/consumer.ts
import { Kafka, EachMessagePayload } from 'kafkajs';

const kafka = new Kafka({
  clientId: 'email-service',
  brokers: ['localhost:9092'],
});

const consumer = kafka.consumer({
  groupId: 'email-service-group',   // Consumer group ID
});

await consumer.connect();

// Subscribe to topic
await consumer.subscribe({
  topic: 'orders',
  fromBeginning: false,  // Start from latest (true = replay all messages)
});

// Process messages
await consumer.run({
  // Process one message at a time (default)
  eachMessage: async ({ topic, partition, message }: EachMessagePayload) => {
    const key = message.key?.toString();
    const value = message.value?.toString();

    if (!value) return;

    const order = JSON.parse(value);
    console.log(`Processing order ${order.orderId} from partition ${partition}`);

    try {
      await sendOrderConfirmationEmail(order);
    } catch (error) {
      console.error(`Failed to process order ${order.orderId}:`, error);
      // Handle error: retry, DLQ, alert
      throw error;  // kafkajs will retry based on retry config
    }
  },
});

Throwing from eachMessage is how you tell kafkajs the message was not processed: it does not advance past it, and it retries according to the consumer’s retry settings. If the error is permanent — a malformed JSON payload, a reference to a deleted user — every retry fails the same way, and eventually the consumer crashes and restarts, only to hit the same message again. That “poison pill” blocks the whole partition, since Kafka delivers a partition’s messages strictly in order. This is the most common Kafka incident I know of in consumer code: one bad message and consumer lag on one partition climbs steadily while the others look healthy. The fix is to distinguish transient from permanent errors and send permanent ones to a dead letter topic (below). Note too that JSON.parse(value) sits outside the try here, so a malformed payload throws before your error handling runs.

Processing time matters for group health. Each consumer sends heartbeats, and if it stops for longer than sessionTimeout (30 seconds by default in kafkajs), the group coordinator considers it dead and reassigns its partitions — a rebalance, during which consumption pauses for the whole group. A slow email API inside eachMessage can therefore cause rebalances that look like instability elsewhere; keep handlers fast, call heartbeat() during long work, or raise the timeout deliberately.

Batch Processing

// Process multiple messages at once (higher throughput)
await consumer.run({
  eachBatch: async ({ batch, resolveOffset, heartbeat, commitOffsetsIfNecessary }) => {
    for (const message of batch.messages) {
      const order = JSON.parse(message.value?.toString() ?? '{}');

      await processOrder(order);

      // Mark this offset as processed (committed later by commitOffsetsIfNecessary)
      resolveOffset(message.offset);

      // Prevent consumer timeout for long-running batches
      await heartbeat();
    }

    await commitOffsetsIfNecessary();
  },
});

eachBatch gives you the whole fetched batch for one partition, which enables genuinely batched work — one bulk database insert for a hundred events instead of a hundred single inserts. This example still processes messages one by one, so its main benefit is control over offsets: resolveOffset marks progress so that if a later message in the batch fails, the batch resumes from there rather than reprocessing the start. The name is easy to misread; it does not commit anything by itself, and the actual commit happens in commitOffsetsIfNecessary according to the auto-commit interval and threshold. If you do switch to a bulk operation, resolve the last offset only after the bulk write succeeds.

Manual Offset Management

// Disable auto-commit for more control
const consumer = kafka.consumer({
  groupId: 'my-group',
});

await consumer.run({
  autoCommit: false,   // Manual commits only
  eachMessage: async ({ topic, partition, message }) => {
    try {
      await processMessage(message);

      // Only commit after successful processing
      await consumer.commitOffsets([{
        topic,
        partition,
        offset: (parseInt(message.offset) + 1).toString(),
      }]);
    } catch (error) {
      // Don't commit — message will be redelivered
      console.error('Processing failed, will retry:', error);
    }
  },
});

The + 1 is essential: a committed offset means “the next message to read”, not “the last message processed”. Committing message.offset itself would redeliver the last processed message after every restart. Also note that autoCommit is an option of consumer.run() in kafkajs, not of kafka.consumer(), which is why it appears there. Committing after every message gives the smallest redelivery window but adds a broker round trip per message; committing periodically is faster and redelivers more after a crash. Either way the guarantee is at-least-once, so consumers must be idempotent — for example by recording processed event ids in the same database transaction as the side effect, or by using upserts so that processing an event twice produces the same result. And as the FAQ below explains, catching and logging the error, as this snippet does, does not by itself cause a retry.


Topics and Partitions Management

const admin = kafka.admin();
await admin.connect();

// Create topic
await admin.createTopics({
  topics: [
    {
      topic: 'orders',
      numPartitions: 6,           // 6 = can have up to 6 parallel consumers
      replicationFactor: 3,       // 3 copies; with min.insync.replicas=2, writes survive 1 broker failure
      configEntries: [
        { name: 'retention.ms', value: String(7 * 24 * 60 * 60 * 1000) },  // 7 days
        { name: 'cleanup.policy', value: 'delete' },
      ],
    },
  ],
});

// List topics
const topics = await admin.listTopics();

// Get topic metadata
const metadata = await admin.fetchTopicMetadata({ topics: ['orders'] });

// List consumer groups
const groups = await admin.listGroups();

// Get consumer group offsets (see where consumers are)
const offsets = await admin.fetchOffsets({ groupId: 'email-service-group', topics: ['orders'] });

await admin.disconnect();

Replication and durability are configured in three places that must agree. replicationFactor: 3 keeps three copies of each partition. The topic or broker setting min.insync.replicas (commonly 2) says how many copies must acknowledge a write for it to count. And the producer’s acks setting (-1, meaning all in-sync replicas, is kafkajs’s default) decides whether the producer waits for that. With 3/2/all, the cluster keeps accepting writes with one broker down and never loses an acknowledged write; with two brokers down, writes fail with NOT_ENOUGH_REPLICAS rather than silently risking data. A replication factor of 3 also requires at least three brokers — on the single-broker Compose setup above, creating this topic fails with an error that the replication factor is larger than the number of available brokers.

Consumer lag — the difference between a partition’s latest offset and the group’s committed offset, which you can compute from fetchOffsets and fetchTopicOffsets — is the metric to alert on. Rising lag means consumers are falling behind, and it tells you which partition is stuck.


Dead Letter Queue Pattern

const DLQ_TOPIC = 'orders.dlq';

await consumer.run({
  eachMessage: async ({ topic, partition, message }) => {
    let attempt = 0;
    const maxAttempts = 3;

    while (attempt < maxAttempts) {
      try {
        await processMessage(message);
        return;  // Success
      } catch (error) {
        attempt++;
        if (attempt === maxAttempts) {
          // Send to DLQ after max retries
          await producer.send({
            topic: DLQ_TOPIC,
            messages: [{
              key: message.key,
              value: message.value,
              headers: {
                ...message.headers,
                'dlq-original-topic': topic,
                'dlq-original-partition': String(partition),
                'dlq-error': String(error),
                'dlq-attempts': String(maxAttempts),
              },
            }],
          });
        } else {
          await sleep(1000 * Math.pow(2, attempt));  // Exponential backoff
        }
      }
    }
  },
});

After the message goes to the DLQ, the handler returns normally, so the consumer commits past it and the partition keeps flowing — that is the point of the pattern. The DLQ copy keeps the original key, value, and headers plus the failure details, which is what makes it possible to inspect and later replay the messages once the bug is fixed; a DLQ without a replay tool tends to become a place where messages go to be forgotten. Two limits of this in-handler version are worth knowing. The backoff sleeps (2 s, then 4 s) block the entire partition while they wait, and they count against the session timeout. For longer or more varied delays, the common design is a chain of retry topics (orders.retry.1m, orders.retry.10m) consumed by separate groups, so the main topic never waits. And if the DLQ send itself fails, the error propagates and the message is retried — which is correct, but means the producer must be connected before consumer.run() starts.


Real-World: Order Processing Pipeline

// Order Service → Kafka → Email Service + Inventory Service + Analytics Service

// order-service/producer.ts
await producer.send({
  topic: 'orders',
  messages: [{
    key: order.id,                  // Same order → same partition → that order's events stay ordered
    value: JSON.stringify({
      type: 'ORDER_CREATED',
      orderId: order.id,
      userId: order.userId,
      items: order.items,
      total: order.total,
      timestamp: Date.now(),
    }),
  }],
});

// email-service/consumer.ts (Consumer Group: email-service)
await consumer.subscribe({ topic: 'orders' });
// Receives ORDER_CREATED — sends confirmation email

// inventory-service/consumer.ts (Consumer Group: inventory-service)
await consumer.subscribe({ topic: 'orders' });
// Receives ORDER_CREATED — deducts inventory

// analytics-service/consumer.ts (Consumer Group: analytics-service)
await consumer.subscribe({ topic: 'orders', fromBeginning: true });
// Receives all events — updates dashboards

Each consumer group reads independently — adding a new service doesn’t affect existing ones.

One subtlety in this design is the analytics service’s fromBeginning: true. It only applies when the group has no committed offsets yet; after the first run, the group resumes from its own commits, so restarting the service does not replay history. To reprocess deliberately, reset the group’s offsets (for example with admin.resetOffsets or the kafka-consumer-groups CLI while the group is stopped) or use a new group id. Replay is also bounded by retention: with the default seven days, a service added six months later can only replay the last week unless the topic keeps data longer or uses log compaction to retain the latest event per key.

The other subtlety is what the producer can promise. The order service typically writes the order to its database and publishes the event; if it crashes between the two, the database and the topic disagree. The standard answer is the transactional outbox pattern: write the event to an outbox table in the same database transaction as the order, and have a separate relay (or a change-data-capture tool such as Debezium) publish outbox rows to Kafka.


Production Configuration

const kafka = new Kafka({
  clientId: 'order-service',
  brokers: process.env.KAFKA_BROKERS!.split(','),

  ssl: true,
  sasl: {
    mechanism: 'scram-sha-512',
    username: process.env.KAFKA_USERNAME!,
    password: process.env.KAFKA_PASSWORD!,
  },

  retry: {
    initialRetryTime: 300,
    retries: 10,
  },
});

const producer = kafka.producer({
  allowAutoTopicCreation: false,    // Don't auto-create in production
  transactionTimeout: 30000,
  idempotent: true,                 // No duplicates from producer retries (not end-to-end exactly-once)
});

idempotent: true makes the broker discard duplicate writes caused by the producer’s own retries — for example when a send succeeded but the acknowledgment was lost and the producer retried. It requires acks of all replicas, and kafkajs documents the feature as experimental, with constraints on in-flight requests. It does not make your consumer exactly-once: a consumer that processes a message, performs a side effect such as sending an email, and crashes before committing will still process it again. Kafka transactions (transactionTimeout configures them) give exactly-once only for read-process-write flows that stay inside Kafka; side effects in external systems still need idempotent handling.

Two production settings in this block matter more than they look. allowAutoTopicCreation: false prevents a typo in a topic name from silently creating a new topic with default partitions and replication, which otherwise shows up weeks later as “why is nobody consuming oders?”. And SASL/SCRAM over TLS is the minimum for any cluster reachable from outside a private network; managed services (Confluent Cloud, Amazon MSK, Aiven) typically require it, sometimes with different mechanisms such as IAM or OAuth.


Frequently Asked Questions (FAQ)

Q. If I catch an error and skip the commit, will Kafka redeliver that message right away?

A. Not by itself. The consumer’s in-memory position has already moved past the failed message, so if eachMessage swallows the error, kafkajs keeps fetching the next messages, and the next successful commitOffsets on that partition records a higher offset that silently skips the failure. The message comes back only if the consumer restarts or rebalances before any later commit. To really retry, rethrow the error so kafkajs retries according to its retry config, or publish the message to a dead letter topic as described in the Dead Letter Queue section and then commit past it.