Apache Kafka · cheat sheet

Apache Kafka

Partitions and keys, replication and acks, consumer groups, offsets, retention, compaction and exactly-once: the Kafka facts interviewers probe.

The Kafka facts that come up in interviews, for Kafka 3.8/4.0 in KRaft mode with the Java client. Code is illustrative and follows the documented APIs.

Core model

  • Topic: named stream, split into partitions. Each partition is an ordered, append-only log stored as segment files.
  • Offset: position of a record inside one partition. Meaningful only per partition; never reused, may have gaps (compaction, transaction markers).
  • Broker: stores partition replicas and serves clients. Controller quorum (KRaft) owns metadata in the __cluster_metadata log.
  • Consumer group: members sharing a group.id; each partition goes to exactly one member. Different groups each get every record.
  • Reading does not delete. Data leaves only through retention or compaction.
Kafka Classic queue (RabbitMQ, SQS)
Consuming does not remove data; replay by resetting offsets message deleted on ack
Position tracking consumer commits an offset per partition broker tracks each message
Ordering per partition per queue, weakened by redelivery and competing consumers
Parallelism up to one consumer per partition per group any number of competing consumers
Fan-out many independent groups exchanges or one queue per subscriber

KRaft timeline

Version Change
3.3 KRaft production-ready for new clusters
3.5 ZooKeeper mode deprecated
3.9 Last release supporting ZooKeeper; migration bridge
4.0 KRaft only; KIP-848 consumer protocol GA; linger.ms default 5 ms; share groups early access

Controllers usually run as a quorum of 3 (tolerates 1 failure) or 5 (tolerates 2), either combined with brokers (dev) or dedicated (production).

Producing and partitioning

  • Partition choice: explicit partition, else murmur2(keyBytes) & 0x7fffffff % numPartitions, else (no key) the sticky partitioner, adaptive since 3.3.
  • Same key, same partition, same order, until the partition count changes. Adding partitions remaps keys; existing data never moves. Partitions can never be removed.
  • send() is asynchronous: it appends to a per-partition batch in the record accumulator; a background thread sends it. Always check the callback, and close() or flush() before exit.
Setting Default Meaning
acks all (since 3.0) wait for every in-sync replica
enable.idempotence true (since 3.0) producer ID + sequence numbers dedupe retries
retries Integer.MAX_VALUE bounded in practice by delivery.timeout.ms
delivery.timeout.ms 120000 total time to deliver or fail a record
max.in.flight.requests.per.connection 5 must be ≤ 5 for idempotence
batch.size 16384 max bytes per partition batch
linger.ms 5 (4.0), 0 before wait to fill a batch
compression.type none lz4 or zstd usually win
buffer.memory 33554432 full buffer blocks send() up to max.block.ms (60 s)
props.put(ProducerConfig.ACKS_CONFIG, "all");
props.put(ProducerConfig.LINGER_MS_CONFIG, 20);
props.put(ProducerConfig.COMPRESSION_TYPE_CONFIG, "zstd");
try (var producer = new KafkaProducer<String, String>(props)) {
    producer.send(new ProducerRecord<>("orders", orderId, json), (meta, ex) -> {
        if (ex != null) log.error("send failed for {}", orderId, ex);
    });
}   // close() flushes buffered records
java

Replication and durability

  • Replication factor (RF): copies per partition, usually 3. One leader handles writes; followers fetch from it.
  • ISR: leader plus followers caught up within replica.lag.time.max.ms (30 s).
  • High watermark: highest offset replicated to all ISR members. Consumers read only below it.
  • min.insync.replicas: minimum ISR size for an acks=all write. Ignored by acks=0/1.
  • unclean.leader.election.enable=false (default): only ISR members can become leader. true restores availability and loses acknowledged data.
RF 3, min.insync.replicas=2, acks=all Result
3 replicas in sync writes succeed
1 broker down writes succeed; every ack is on 2 brokers
2 brokers down NotEnoughReplicasException; reads still work

Interview tip

The durable baseline: RF 3, min.insync.replicas=2, acks=all, idempotence on, unclean election off, replicas spread across racks with broker.rack.

Consumers, groups and rebalancing

Setting Default Meaning
enable.auto.commit true commit polled positions every auto.commit.interval.ms (5 s)
auto.offset.reset latest used only with no valid committed offset
max.poll.records 500 records per poll()
max.poll.interval.ms 300000 max gap between polls before the member leaves
session.timeout.ms 45000 heartbeat timeout (classic protocol)
heartbeat.interval.ms 3000 background heartbeat
isolation.level read_uncommitted read_committed hides aborted and open transactions
group.protocol classic consumer enables KIP-848 (GA in 4.0)
  • Rebalance triggers: member joins, leaves, misses session.timeout.ms, exceeds max.poll.interval.ms; subscribed topic gains partitions.
  • Eager assignors (range, round-robin, sticky) revoke everything first. Cooperative (CooperativeStickyAssignor) moves only what changes. KIP-848 lets the broker compute assignments incrementally with no group-wide barrier.
  • Static membership: stable group.instance.id; a restart within session.timeout.ms gets the same partitions with no rebalance.
  • CommitFailedException usually means the member was already removed from the group, typically for exceeding max.poll.interval.ms.
consumer.subscribe(List.of("orders"), new ConsumerRebalanceListener() {
    public void onPartitionsRevoked(Collection<TopicPartition> parts) {
        consumer.commitSync(processedOffsets);   // your map of next offsets per partition
    }
    public void onPartitionsAssigned(Collection<TopicPartition> parts) { }
});
java

Offsets and delivery semantics

Pattern Guarantee Failure mode
commit, then process at-most-once crash after commit skips records
process, then commit at-least-once crash before commit replays records
transactional read-process-write + read_committed exactly-once inside Kafka external side effects still need idempotence
  • Committed offset = next record to read (record.offset() + 1).
  • commitSync() retries until success or a fatal error; commitAsync() never retries (an old commit could overwrite a newer one). Use async in the loop, sync on shutdown and revoke.
  • Auto-commit is at-least-once only if processing finishes inside the poll loop; handing records to another thread can lose them.
  • Exactly-once effects in a database: dedupe by event ID (unique constraint, ON CONFLICT DO NOTHING) or store offsets in the same DB transaction and seek() on assignment.

Retention and compaction

Setting Default Meaning
cleanup.policy delete compact or compact,delete
retention.ms 604800000 (7 d) delete closed segments older than this
retention.bytes -1 per-partition size cap
segment.bytes 1073741824 (1 GiB) roll the active segment
segment.ms 604800000 (7 d) roll by age
delete.retention.ms 86400000 (24 h) how long tombstones survive compaction
min.compaction.lag.ms 0 records younger than this are not compacted
max.compaction.lag.ms unbounded force compaction eligibility after this
min.cleanable.dirty.ratio 0.5 dirty fraction that triggers cleaning
  • Deletion and compaction work on closed segments only; the active segment is untouched.
  • Compaction keeps at least the latest value per key; offsets keep their gaps. A tombstone is a key with a null value. Null keys are rejected on compacted topics.
  • A consumer whose committed offset was deleted falls back to auto.offset.reset.

Transactions

producer.initTransactions();                      // fences older instances with this transactional.id
while (running) {
    var records = consumer.poll(Duration.ofMillis(200));
    if (records.isEmpty()) continue;
    producer.beginTransaction();
    try {
        for (var r : records) producer.send(new ProducerRecord<>("out", r.key(), transform(r.value())));
        producer.sendOffsetsToTransaction(nextOffsets(records), consumer.groupMetadata());
        producer.commitTransaction();
    } catch (ProducerFencedException e) {
        producer.close(); throw e;                // a newer instance took over
    } catch (KafkaException e) {
        producer.abortTransaction(); rewind(consumer, records);
    }
}
java
  • Coordinator state lives in __transaction_state; commit or abort markers go into every touched partition.
  • read_committed consumers stop at the last stable offset (before the earliest open transaction) and skip aborted records.
  • Kafka Streams: processing.guarantee=exactly_once_v2 (needs brokers 2.5+).

Errors and poison messages

  • A record that always fails blocks its partition when the consumer never commits past it.
  • Classify: deserialization (catch RecordDeserializationException, then seek() past it), transient (retry with backoff), permanent (send to a dead letter topic with original bytes and error headers, then commit).
  • Retry topics trade away per-key ordering. Kafka Connect sinks: errors.tolerance=all + errors.deadletterqueue.topic.name.

CLI

kafka-topics.sh --bootstrap-server localhost:9092 --create --topic orders \
  --partitions 12 --replication-factor 3 --config min.insync.replicas=2
kafka-topics.sh --bootstrap-server localhost:9092 --describe --topic orders   # leader, replicas, ISR
kafka-consumer-groups.sh --bootstrap-server localhost:9092 --describe --group billing   # LAG per partition
kafka-consumer-groups.sh --bootstrap-server localhost:9092 --group billing --topic orders \
  --reset-offsets --to-datetime 2026-10-06T08:00:00.000 --execute   # group must be inactive
kafka-configs.sh --bootstrap-server localhost:9092 --alter --entity-type topics \
  --entity-name orders --add-config retention.ms=259200000
kafka-get-offsets.sh --bootstrap-server localhost:9092 --topic orders   # log end offsets
kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic orders \
  --from-beginning --property print.key=true --group debug-reader
Terminal

Gotchas

  • More consumers than partitions leaves the extras idle.
  • New service with auto.offset.reset=latest misses everything produced before its first poll.
  • acks=all with min.insync.replicas=1 is effectively acks=1 once followers drop out.
  • Lag can read zero while work is still pending if offsets are committed before async processing completes.
  • One huge key (a big tenant, "unknown") creates a hot partition that no amount of consumers can fix.
  • Calling .get() on every send() serializes the producer and destroys batching.

Practise with the full question bank on the Apache Kafka topic page.

esc