Apache Kafka is usually described as a message queue, and that description causes most of the mistakes teams make with it. Kafka is a distributed, replicated, append-only log. Producers append records to the end; consumers read from any position they choose and remember that position themselves. Nothing is deleted because it was read. Records disappear only when a retention rule says they are old enough or, on compacted topics, when a newer record with the same key replaces them.

That one design choice explains Kafka's strengths: many independent readers of the same data, replay after a bug, high throughput from sequential disk writes, and ordering you can reason about. It also explains its sharp edges: ordering only within a partition, partition counts that are hard to change, and duplicates you must handle. This overview builds the mental model, then walks through a first production topic end to end. For the broker internals behind it, read Apache Kafka streaming architecture next.

Advertisement

The mental model in one picture

Producer Acheckout serviceProducer Bpayments serviceTopic: orders (3 partitions, RF=3)Partition 0leader on broker 1, followers 2 and 3Partition 1leader on broker 2, followers 3 and 1Partition 2leader on broker 3, followers 1 and 2offsets 0, 1, 2 ... append-only, per partitionkey hashkey hashConsumer 1group: billingConsumer 2group: billingP0P1P2Consumer 3group: analyticsall partitionsKRaft controller quorummetadata log: topics, leaders, ISRleader electionsEach group gets its own copy of the stream; within a group each partition has exactly one owner.
Two producers write to a three-partition topic replicated three times. The billing group splits partitions between two consumers; the analytics group independently reads everything. A KRaft controller quorum stores cluster metadata and elects partition leaders.

Five nouns carry the whole system. A record is a key, a value, a timestamp and optional headers, all opaque bytes to the broker. A topic is a named stream of records, such as orders. A partition is one ordered, append-only file sequence within a topic; a topic with twelve partitions is really twelve independent logs. An offset is a record's position inside its partition, a 64-bit integer that only grows. A broker is a server that stores partitions and serves reads and writes for them.

Each partition is replicated to several brokers. One replica is the leader and handles all writes; the others are followers that copy from it. The set of followers that are caught up is the in-sync replica set, and its rules decide what an acknowledged write really means; ISR replication covers them in detail. Cluster metadata (which topics exist, who leads each partition) lives in a replicated metadata log run by a small quorum of controller nodes using KRaft. Since Kafka 4.0 there is no ZooKeeper mode at all, so any guide that asks you to install ZooKeeper is describing a version you should not deploy.

Keys, partitions and the only ordering you get

When a producer sends a record with a key, the default partitioner hashes the key (murmur2) and takes it modulo the partition count. Every record for customer 42 therefore lands in the same partition, and within a partition Kafka preserves append order. That is the entire ordering guarantee: per partition, never across partitions. If two events must be processed in order, they need the same key.

Records sent without a key are spread for throughput, not ordering. Modern clients batch keyless records to one partition at a time and switch partitions as batches fill, which produces larger batches than strict round-robin. Treat keyless as "I do not care about order".

The modulo is why partition counts are sticky. Going from 12 to 16 partitions changes where most keys hash, so new records for customer 42 may land in a different partition from its older records, and a consumer can briefly see them out of order. Kafka can add partitions but never remove them. Choose the count for expected peak load plus headroom up front; Kafka partitioning strategies covers custom partitioners and hot-key mitigation.

Advertisement

Producers: what an acknowledged write means

A producer buffers records per partition, compresses batches and sends them to partition leaders. Three settings decide durability and duplicates, and since Kafka 3.0 the defaults are the safe ones; make sure no old configuration file overrides them.

SettingSafe valueWhat it controls
acksallLeader waits until every in-sync replica has the batch before acknowledging
enable.idempotencetrueBroker drops retried duplicates using producer id and sequence numbers, and per-partition order survives retries
delivery.timeout.msdefault 120000, or your SLAUpper bound on send plus retries before the callback reports failure
linger.ms / batch.size5-20 ms / 64-256 KiBThroughput versus latency: wait briefly to fill bigger batches
compression.typelz4 or zstdSmaller batches on the wire and on disk; compression is per batch

acks=all only means as much as the topic's min.insync.replicas. With replication factor 3 and min.insync.replicas=2, an acknowledged write exists on at least two brokers, and if only one replica is in sync the producer gets NotEnoughReplicas instead of a false promise. The broker default for min.insync.replicas is 1, so set it explicitly per topic or cluster-wide.

Properties props = new Properties();
props.put("bootstrap.servers", "kafka-1:9092,kafka-2:9092,kafka-3:9092");
props.put("key.serializer", StringSerializer.class.getName());
props.put("value.serializer", ByteArraySerializer.class.getName());
props.put("acks", "all");
props.put("enable.idempotence", "true");
props.put("linger.ms", "10");
props.put("compression.type", "zstd");

try (KafkaProducer<String, byte[]> producer = new KafkaProducer<>(props)) {
    var rec = new ProducerRecord<>("orders", order.customerId(), order.toBytes());
    producer.send(rec, (meta, err) -> {
        if (err != null) log.error("send failed after retries", err);  // alert, park or fail the request
    });
}

Consumers, groups and committed offsets

Consumers that share a group.id form a consumer group, and Kafka assigns each partition to exactly one member. Add a consumer and partitions move to it; remove one and its partitions move elsewhere. That is how you scale processing, and it implies a hard ceiling: a group never uses more consumers than the topic has partitions. A second group with a different id gets its own complete copy of the stream, which is how billing and analytics both read orders without interfering.

Progress is a committed offset per partition, stored in an internal topic. The offset is the next record to read, and it is only as correct as when you commit it. Commit before processing and a crash loses records. Commit after processing and a crash replays the records processed since the last commit. The second behaviour, at-least-once, is the right default, and it means your processing must tolerate duplicates: upsert by a business key, or record processed ids in the same database transaction as the effect.

props.put("group.id", "billing");
props.put("enable.auto.commit", "false");
props.put("auto.offset.reset", "earliest");      // where a brand-new group starts
props.put("max.poll.records", "500");

consumer.subscribe(List.of("orders"));
while (running) {
    var records = consumer.poll(Duration.ofMillis(500));
    for (var r : records) {
        billing.applyIdempotently(r.key(), r.value(), r.partition(), r.offset());
    }
    consumer.commitSync();   // after the work: at-least-once
}

Membership changes cause a rebalance, and a consumer that takes longer than max.poll.interval.ms between polls is assumed dead and triggers one. Slow per-record work is therefore a stability problem, not just a latency problem. The newer consumer group protocol from KIP-848, generally available since Kafka 4.0 and enabled on the client with group.protocol=consumer, moves assignment to the broker and avoids stop-the-world rebalances; consumer rebalancing explains both protocols. If what you need is a work queue where many consumers share one partition with per-record acknowledgement, look at share groups (KIP-932), which were a preview in 4.1 and are documented as production-ready in 4.2; check your cluster and client versions before relying on them.

Retention: the log is not a mailbox

Records leave a partition by policy. With cleanup.policy=delete, whole segment files are removed once they are older than retention.ms (seven days by default) or the partition exceeds retention.bytes. Deletion is per segment, so data can survive somewhat longer than the limit. With cleanup.policy=compact, Kafka keeps at least the latest record per key and removes superseded ones in the background, which turns a topic into a durable table of current state; log compaction explains tombstones and the cleaner.

Retention is a product decision. It sets how far back a new consumer can bootstrap and how long you have to fix a bug and replay. If the consumer's committed offset falls off the front of the log because it was down too long, auto.offset.reset decides whether it jumps to the oldest available record or the newest, and either way some records were never processed. Alert on consumer lag in time, not only in record count, and keep retention comfortably above your worst realistic outage.

Worked example: sizing an orders topic

Suppose checkout produces 3,000 orders per second at peak, each about 1.5 KB after serialization, keyed by customer id. Billing must keep up within a few seconds, and analytics replays the last 14 days when its model changes.

  1. Ingest bandwidth. 3,000 x 1.5 KB = 4.5 MB/s before compression. Order JSON often compresses well, but plan with the uncompressed figure until you measure.
  2. Partitions from consumer throughput. Suppose one billing consumer handles 400 records/s because each record does a database write. Peak needs 3,000 / 400 = 7.5 consumers, so at least 8 partitions. Doubling for growth and uneven keys gives 16, and 24 leaves room to triple the consumer count without repartitioning. Choose 24.
  3. Disk. 4.5 MB/s x 86,400 s = about 389 GB per day, so 14 days is about 5.4 TB of unique data. With replication factor 3 the cluster stores about 16.3 TB, spread across brokers, plus headroom so a broker loss does not fill the survivors.
  4. Replication traffic. Each broker receives its share of 4.5 MB/s from producers and twice that again from followers fetching, and serves every consumer group on top. Network, not disk, is often the first limit.
  5. Durability. Replication factor 3, min.insync.replicas=2, producers with acks=all and idempotence. The cluster survives one broker loss with no stop to writes.
kafka-topics.sh --bootstrap-server kafka-1:9092 --create --topic orders \
  --partitions 24 --replication-factor 3 \
  --config min.insync.replicas=2 \
  --config retention.ms=1209600000 \
  --config compression.type=producer

If 14 days of hot disk is too expensive, tiered storage lets older segments live in object storage while recent ones stay local. The price is slower reads when a consumer replays from the remote tier, so measure a 14-day replay before relying on it.

Failure modes to plan for

SymptomUsual causeWhat to do
Duplicate side effects after a deployAt-least-once replay of records processed after the last commitMake the handler idempotent; commit after processing
Consumer group keeps rebalancingProcessing exceeds max.poll.interval.ms, or pods restart in a loopSmaller max.poll.records, faster handlers, the KIP-848 protocol
One consumer far behind the othersHot key concentrating traffic in one partitionRethink the key, split the hot entity, or accept per-key serialization
Producers fail with NotEnoughReplicasFewer in-sync replicas than min.insync.replicasFix the lagging or dead broker; do not lower the setting to make errors stop
Records silently skippedCommitted offset fell behind retention, and the reset policy jumpedLag alerts in time, longer retention, capacity for catch-up
Order broken for some keysPartitions added to a keyed topicChoose counts up front; if you must grow, migrate to a new topic

One more: the broker cannot see your data. A producer that changes its serialization breaks every consumer at once. Put a schema registry or at least a versioned format in front of any topic with more than one consumer team.

When Kafka is the wrong tool

Kafka fits high-volume event streams, change data capture, many independent consumers and replayable history. It is a poor fit for request-reply RPC, for per-message delays and priorities, for millions of tiny topics, and for small systems where one database table and a polling worker would do. It is also not exactly-once end to end by default: transactions give exactly-once between Kafka topics, but any external side effect still needs idempotence. Read exactly-once processing before promising it to anyone.

What to do next

  1. Run a three-broker KRaft cluster locally with containers, create a topic with 6 partitions and replication factor 3, and produce keyed records with the console producer.
  2. Start two consumers in one group and one in another; kill a consumer and watch partitions move with kafka-consumer-groups.sh --describe.
  3. Stop a broker and confirm that producers with acks=all and min.insync.replicas=2 keep working, then stop a second and observe the error.
  4. Size your first real topic from consumer throughput, as in the worked example, and write the partition count, key and retention down with the reasoning.
  5. Make every consumer handler idempotent before production, and test it by replaying a partition from offset zero.
  6. Add alerts for consumer lag in seconds, under-replicated partitions and offline partitions before the first incident.
Key takeaway: Kafka is a replicated, partitioned, append-only log, not a queue. Order exists only within a partition, so choose keys for ordering and partition counts for peak consumer throughput, because changing them later reorders keys. Use acks=all with idempotence and min.insync.replicas=2 on replication factor 3, commit offsets after processing, make handlers idempotent, and set retention from how far back you will need to replay.