Summary
Apache Kafka is a distributed streaming platform designed for building real-time data pipelines and streaming applications. This cheatsheet covers core Kafka concepts including topics, partitions, producers, consumers, brokers, replication, offset management, consumer groups, message delivery semantics, performance optimization, monitoring, and integration patterns. Key features include high throughput, horizontal scalability, fault tolerance, durability, and support for both pub-sub messaging and stream processing use cases.
1. Core Concepts
What is Kafka?
- Distributed streaming platform for building real-time data pipelines
- Publish-Subscribe messaging system with high throughput
- Horizontally scalable, fault-tolerant, and durable
Key Components
Producer → Kafka Cluster → Consumer
├── Broker 1
├── Broker 2
└── Broker N
2. Architecture Components
Topic
- Logical channel for organizing messages
- Partitioned for parallelism and scalability
- Immutable append-only log
Partition
- Ordered sequence of messages
- Unit of parallelism in Kafka
- Each partition has a leader and replicas
Broker
- Kafka server that stores data
- Manages partitions and handles requests
- Multiple brokers form a cluster
Producer
- Publishes messages to topics
- Can specify partition or use partitioner
- Supports async and sync modes
Consumer
- Reads messages from topics
- Part of a consumer group
- Maintains offset for tracking
ZooKeeper (Legacy) / KRaft (New)
- Coordination service for cluster management
- Stores metadata, manages leader election
- KRaft mode eliminates ZooKeeper dependency
3. Key Concepts for Interviews
Consumer Groups
- Load balancing: Partitions distributed among consumers
- Fault tolerance: Automatic rebalancing on failure
- Scalability: Max consumers = number of partitions
// Consumer group example
Properties props = new Properties();
props.put("group.id", "my-consumer-group");
props.put("enable.auto.commit", "true");
Offset Management
- Current offset: Next message to read
- Committed offset: Last processed message
- Log end offset: Last message in partition
Replication
- Replication factor: Number of copies
- ISR (In-Sync Replicas): Replicas caught up with leader
- Leader election: Automatic on failure
Message Delivery Semantics
- At most once: Messages may be lost
- At least once: Messages may be duplicated
- Exactly once: No loss, no duplication (idempotent producer)
4. Producer Configuration
Key Properties
Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092");
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
// Performance tuning
props.put("batch.size", 16384);
props.put("linger.ms", 1);
props.put("buffer.memory", 33554432);
// Reliability
props.put("acks", "all"); // 0, 1, or all
props.put("retries", 3);
props.put("enable.idempotence", true);
Producer Example
Producer<String, String> producer = new KafkaProducer<>(props);
ProducerRecord<String, String> record =
new ProducerRecord<>("topic", "key", "value");
// Async send
producer.send(record, (metadata, exception) -> {
if (exception != null) {
exception.printStackTrace();
}
});
// Sync send
try {
RecordMetadata metadata = producer.send(record).get();
} catch (Exception e) {
e.printStackTrace();
}
5. Consumer Configuration
Consumer Properties
Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092");
props.put("group.id", "test-group");
props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
// Offset management
props.put("enable.auto.commit", "false");
props.put("auto.offset.reset", "earliest"); // earliest, latest, none
// Performance
props.put("max.poll.records", 500);
props.put("fetch.min.bytes", 1);
Consumer Example
Consumer<String, String> consumer = new KafkaConsumer<>(props);
consumer.subscribe(Arrays.asList("topic1", "topic2"));
while (true) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
for (ConsumerRecord<String, String> record : records) {
// Process record
System.out.printf("offset = %d, key = %s, value = %s%n",
record.offset(), record.key(), record.value());
}
consumer.commitSync(); // Manual commit
}
6. Advanced Topics
Kafka Streams
- Stream processing library
- Stateful transformations
- Exactly-once processing
StreamsBuilder builder = new StreamsBuilder();
KStream<String, String> source = builder.stream("input-topic");
source.filter((key, value) -> value.contains("important"))
.to("output-topic");
Kafka Connect
- Import/Export data between Kafka and external systems
- Source connectors: Import data into Kafka
- Sink connectors: Export data from Kafka
Schema Registry
- Centralized schema management
- Schema evolution support
- Compatibility checking
7. Performance Optimization
Producer Optimization
- Batching: Increase
batch.sizeandlinger.ms - Compression: Use
compression.type(gzip, snappy, lz4, zstd) - Partitioning: Custom partitioner for better distribution
Consumer Optimization
- Parallel processing: Multiple consumers in group
- Fetch size: Tune
fetch.min.bytesandmax.poll.records - Session timeout: Balance
session.timeout.ms
Broker Optimization
- Replication: Balance between durability and performance
- Log segment: Tune
log.segment.bytes - Memory: Increase heap size and page cache
8. Key Concepts & Comparisons
Message Delivery Guarantees
| Guarantee | Producer Config | Consumer Behavior | Use Case | Trade-offs |
|---|---|---|---|---|
| At Most Once | acks=0 |
Commit before processing | Fire-and-forget scenarios | Possible message loss |
| At Least Once | acks=1 or acks=all |
Process then commit | Critical data delivery | Possible duplicates |
| Exactly Once | enable.idempotence=true + transactions |
Transactional processing | Financial transactions | Higher complexity, lower throughput |
Producer Acknowledgment Settings
| Acks Setting | Behavior | Durability | Performance | Use Case |
|---|---|---|---|---|
acks=0 |
No wait for broker response | Lowest | Highest | Logs, metrics where loss is acceptable |
acks=1 |
Wait for leader confirmation | Medium | Medium | General messaging |
acks=all |
Wait for all ISR replicas | Highest | Lowest | Critical data |
Consumer Group vs Topic Subscription
| Aspect | Consumer Group | Topic Subscription | Impact |
|---|---|---|---|
| Load Balancing | Automatic partition distribution | Manual partition assignment | Group provides automatic scaling |
| Fault Tolerance | Automatic rebalancing on failure | Manual handling required | Group provides high availability |
| Offset Management | Coordinated through group coordinator | Individual consumer responsibility | Group simplifies offset management |
| Scaling Limit | Max consumers = partition count | No built-in limit | Partition count determines parallelism |
Partitioning Strategies
| Strategy | Implementation | Use Case | Benefits | Considerations |
|---|---|---|---|---|
| Round Robin | Default partitioner (null key) | Even load distribution | Balanced partition sizes | No ordering guarantees |
| Key-based | Hash of message key | Related messages grouping | Per-key message ordering | Potential hot partitions |
| Custom | Custom partitioner logic | Specific business requirements | Tailored distribution | Development overhead |
| Random | Random partition selection | Simple load balancing | Even distribution | No ordering, unpredictable |
Kafka vs Other Messaging Systems
| System | Ordering | Durability | Throughput | Routing | Use Case |
|---|---|---|---|---|---|
| Kafka | Per-partition | Disk-based, configurable | Very High | Topic-based | Event streaming, log aggregation |
| RabbitMQ | Queue-based | Memory/disk | Medium | Complex routing | Task queues, RPC patterns |
| Redis Pub/Sub | None | Memory-only | High | Channel-based | Real-time notifications |
| AWS SQS | FIFO queues | Managed | Medium | Simple | Decoupling microservices |
| Apache Pulsar | Per-partition | Tiered storage | High | Multi-tenancy | Multi-tenant streaming |
Consumer Rebalancing Triggers
| Trigger | Cause | Impact | Mitigation |
|---|---|---|---|
| Consumer Join | New consumer joins group | Partition redistribution | Plan for temporary pause |
| Consumer Leave | Consumer stops/crashes | Partition redistribution | Monitor consumer health |
| Partition Change | Topic partition increase | Rebalance to include new partitions | Coordinate partition changes |
| Group Coordinator Failover | Broker hosting coordinator fails | Full group rebalance | Monitor broker health |
Offset Management Patterns
| Pattern | Implementation | Pros | Cons | Use Case |
|---|---|---|---|---|
| Auto Commit | enable.auto.commit=true |
Simple, automatic | Risk of message loss/duplication | Non-critical processing |
| Manual Commit | consumer.commitSync() |
Full control | More complex code | Critical processing |
| Commit After Batch | Process batch, then commit | Higher throughput | Larger potential loss window | Batch processing |
| External Store | Store offsets externally | Exactly-once with external systems | Complex implementation | Database integration |
Performance Optimization Strategies
| Component | Parameter | Optimization | Impact | Trade-offs |
|---|---|---|---|---|
| Producer | batch.size |
Increase for higher throughput | More messages per batch | Higher latency |
| Producer | linger.ms |
Wait time for batching | Better compression | Added latency |
| Producer | compression.type |
Enable compression | Reduced network/storage | CPU overhead |
| Consumer | fetch.min.bytes |
Minimum fetch size | Reduced network calls | Higher latency |
| Consumer | max.poll.records |
Records per poll | Better throughput | Memory usage |
| Broker | num.network.threads |
Network thread pool | Better concurrency | Resource usage |
Common Troubleshooting Scenarios
| Problem | Symptoms | Common Causes | Solutions |
|---|---|---|---|
| Consumer Lag | Increasing lag metrics | Slow processing, insufficient consumers | Scale consumers, optimize processing |
| Rebalancing | Frequent group rebalances | Long processing, network issues | Tune session timeout, optimize processing |
| Producer Failures | Send errors, timeouts | Network issues, broker problems | Check connectivity, broker health |
| Message Loss | Missing messages | Wrong acks setting, broker failures | Use acks=all, monitor ISR |
| Duplicate Messages | Same message multiple times | Retries, network issues | Enable idempotence, deduplication |
| Hot Partitions | Uneven load distribution | Poor key distribution | Review partitioning strategy |
Monitoring Key Metrics
| Component | Metric | Meaning | Alert Threshold | Action |
|---|---|---|---|---|
| Producer | record-send-rate |
Messages sent per second | < expected rate | Check producer health |
| Producer | request-latency-avg |
Average request latency | > SLA threshold | Optimize or scale |
| Consumer | records-lag-max |
Maximum consumer lag | > business threshold | Scale consumers |
| Consumer | fetch-rate |
Fetch requests per second | Dropping significantly | Check consumer health |
| Broker | under-replicated-partitions |
Partitions without full replication | > 0 | Check broker/disk health |
| Broker | isr-shrinks-per-sec |
ISR shrinkage rate | > normal rate | Investigate broker issues |
Kafka Ecosystem Components
| Component | Purpose | Use Case | Integration |
|---|---|---|---|
| Kafka Streams | Stream processing | Real-time transformations | Java library, embedded |
| Kafka Connect | Data integration | ETL, system integration | Standalone/distributed mode |
| Schema Registry | Schema management | Data governance | REST API, client integration |
| KSQL/ksqlDB | SQL on streams | Stream analytics | SQL interface |
| Kafka REST Proxy | HTTP interface | Non-JVM language support | REST API |
Deployment Patterns
| Pattern | Architecture | Benefits | Considerations |
|---|---|---|---|
| Single Cluster | All topics in one cluster | Simple management | Single point of failure |
| Multi-Cluster | Topics across clusters | Isolation, scaling | Complex coordination |
| Active-Passive | Primary + backup cluster | Disaster recovery | Resource duplication |
| Multi-Region | Clusters in different regions | Geographic distribution | Network latency |
| Kafka on Kubernetes | Container orchestration | Easy scaling | Persistent storage complexity |
Security Configuration
| Security Layer | Implementation | Configuration | Use Case |
|---|---|---|---|
| Authentication | SASL/SCRAM | security.protocol=SASL_SSL |
Client authentication |
| Authorization | Kafka ACLs | authorizer.class.name |
Resource access control |
| Encryption | SSL/TLS | ssl.keystore.location |
Data in transit |
| Inter-broker | SSL | security.inter.broker.protocol |
Broker communication |
9. Best Practices
Topic Design
- Naming: Use descriptive, hierarchical names
- Partitions: Start conservative, increase as needed
- Retention: Set based on use case
Error Handling
// Producer error handling
producer.send(record, (metadata, exception) -> {
if (exception != null) {
// Log error, retry, or send to DLQ
logger.error("Failed to send message", exception);
}
});
// Consumer error handling
try {
processRecord(record);
consumer.commitSync();
} catch (Exception e) {
// Skip bad record or retry
logger.error("Failed to process record", e);
}
Monitoring
- Metrics to track:
- Producer: Send rate, error rate, batch size
- Consumer: Lag, fetch rate, commit rate
- Broker: Under-replicated partitions, ISR shrinks
10. Quick Reference & Best Practices
Essential Kafka CLI Commands
| Command | Purpose | Example |
|---|---|---|
| Create Topic | Create new topic | kafka-topics.sh --create --topic events --partitions 3 --replication-factor 2 |
| List Topics | Show all topics | kafka-topics.sh --list --bootstrap-server localhost:9092 |
| Describe Topic | Topic details | kafka-topics.sh --describe --topic events --bootstrap-server localhost:9092 |
| Alter Topic | Modify topic | kafka-topics.sh --alter --topic events --partitions 6 |
| Delete Topic | Remove topic | kafka-topics.sh --delete --topic events --bootstrap-server localhost:9092 |
| Consumer Groups | List groups | kafka-consumer-groups.sh --list --bootstrap-server localhost:9092 |
| Group Details | Group information | kafka-consumer-groups.sh --describe --group my-group |
| Reset Offsets | Reset consumer offsets | kafka-consumer-groups.sh --reset-offsets --to-earliest --topic events --group my-group |
Configuration Quick Reference
| Component | Key Parameters | Recommended Values | Impact |
|---|---|---|---|
| Producer | acks |
all for critical data |
Durability vs performance |
| Producer | batch.size |
16384 (16KB) |
Throughput vs latency |
| Producer | linger.ms |
1-5ms |
Batching efficiency |
| Producer | compression.type |
lz4 or snappy |
Network/storage efficiency |
| Consumer | enable.auto.commit |
false for critical processing |
Offset management control |
| Consumer | max.poll.records |
500-1000 |
Memory vs throughput |
| Consumer | session.timeout.ms |
30000 (30s) |
Rebalancing sensitivity |
| Broker | num.replica.fetchers |
4-8 |
Replication performance |
Performance Tuning Checklist
- Producer Batching: Configure
batch.sizeandlinger.msappropriately - Compression: Enable compression (
lz4,snappy,gzip,zstd) - Partitioning: Choose optimal partition count (start with 2-3x broker count)
- Replication Factor: Balance between durability (3) and performance (1-2)
- Memory Allocation: Set appropriate JVM heap size (6-8GB typical)
- Disk Configuration: Use dedicated disks, consider RAID 10
- Network: Ensure adequate network bandwidth
- Monitoring: Set up comprehensive monitoring and alerting
Topic Design Best Practices
| Aspect | Recommendation | Rationale | Example |
|---|---|---|---|
| Naming | Use hierarchical naming | Organization and discovery | ecommerce.orders.created |
| Partitions | Start conservative, scale up | Easier to increase than decrease | 3-6 partitions initially |
| Replication | Use 3 for production | Balance durability and performance | replication-factor=3 |
| Retention | Match business requirements | Storage cost vs replay needs | 7 days for logs, 30 days for events |
| Key Selection | Choose for even distribution | Avoid hot partitions | User ID, device ID, region |
Error Handling Patterns
// Producer with retry and error handling
Properties props = new Properties();
props.put("retries", 3);
props.put("retry.backoff.ms", 1000);
props.put("enable.idempotence", true);
producer.send(record, (metadata, exception) -> {
if (exception != null) {
if (exception instanceof RetriableException) {
// Will be retried automatically
logger.warn("Retriable error occurred", exception);
} else {
// Send to dead letter queue or alert
logger.error("Non-retriable error", exception);
sendToDeadLetterQueue(record);
}
}
});
// Consumer with error handling
while (true) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
for (ConsumerRecord<String, String> record : records) {
try {
processMessage(record);
} catch (TemporaryException e) {
// Retry logic
retryMessage(record);
} catch (PermanentException e) {
// Skip or send to DLQ
logger.error("Permanent error processing message", e);
sendToDeadLetterQueue(record);
}
}
consumer.commitSync();
}
Monitoring and Alerting Strategy
| Metric Category | Key Metrics | Alert Conditions | Actions |
|---|---|---|---|
| Producer Health | record-send-rate, request-latency-avg |
Send rate < threshold, latency > SLA | Check producer, network, brokers |
| Consumer Health | records-lag-max, fetch-rate |
Lag > threshold, fetch rate dropping | Scale consumers, optimize processing |
| Broker Health | under-replicated-partitions, offline-partitions |
Any value > 0 | Check broker health, disk space |
| Cluster Health | active-controller-count, leader-election-rate |
Controller != 1, high election rate | Investigate controller stability |
Common Architecture Patterns
| Pattern | Use Case | Implementation | Benefits |
|---|---|---|---|
| Event Sourcing | Audit trails, state reconstruction | Store all changes as events | Complete history, replay capability |
| CQRS | Read/write separation | Separate command and query models | Optimized read/write performance |
| Outbox Pattern | Transactional guarantees | Database + Kafka in same transaction | Exactly-once delivery |
| Saga Pattern | Distributed transactions | Choreography or orchestration | Eventual consistency |
| Change Data Capture | Database replication | Kafka Connect CDC connectors | Real-time data synchronization |
Security Implementation Checklist
- Enable SSL/TLS: Encrypt data in transit
- Configure SASL: Implement authentication (SCRAM, PLAIN, GSSAPI)
- Set up ACLs: Fine-grained authorization control
- Network Segmentation: Isolate Kafka cluster
- Audit Logging: Enable and monitor access logs
- Regular Updates: Keep Kafka version current
- Secret Management: Use secure credential storage
- Inter-broker Security: Secure broker-to-broker communication
Production Deployment Best Practices
# JVM Configuration
export KAFKA_HEAP_OPTS="-Xmx6g -Xms6g"
export KAFKA_JVM_PERFORMANCE_OPTS="-server -XX:+UseG1GC -XX:MaxGCPauseMillis=20
-XX:InitiatingHeapOccupancyPercent=35 -XX:+ExplicitGCInvokesConcurrent
-Djava.awt.headless=true"
# Server Properties
num.network.threads=8
num.io.threads=16
socket.send.buffer.bytes=102400
socket.receive.buffer.bytes=102400
socket.request.max.bytes=104857600
num.partitions=3
default.replication.factor=3
min.insync.replicas=2
unclean.leader.election.enable=false
log.retention.hours=168
log.segment.bytes=1073741824
Capacity Planning Guidelines
| Component | Sizing Factor | Calculation | Example |
|---|---|---|---|
| Disk Storage | Retention period × message rate × size | 7 days × 1M msg/day × 1KB = 7GB | Add 20% buffer |
| Network Bandwidth | Peak throughput × replication factor | 100MB/s × 3 = 300MB/s | Account for consumer traffic |
| Memory | Page cache + JVM heap | OS cache: 60%, JVM: 40% | 16GB total: 10GB cache, 6GB JVM |
| CPU Cores | Network + I/O threads + processing | 8 network + 16 I/O + 4 processing | 32 cores recommended |
Troubleshooting Decision Tree
High Consumer Lag?
├── Yes → Check consumer processing time
│ ├── High → Optimize processing or scale consumers
│ └── Normal → Check partition count vs consumer count
└── No → Check other metrics
Frequent Rebalancing?
├── Yes → Check session.timeout.ms vs processing time
│ ├── Timeout too low → Increase session.timeout.ms
│ └── Processing too slow → Optimize or parallelize processing
└── No → Normal operation
Producer Errors?
├── Network errors → Check connectivity and broker health
├── Serialization errors → Validate message format
├── Authorization errors → Check ACLs and credentials
└── Timeout errors → Check broker load and configuration
Testing Strategy
| Test Type | Purpose | Tools | Focus Areas |
|---|---|---|---|
| Unit Tests | Individual components | JUnit, Mockito | Producer/Consumer logic |
| Integration Tests | Kafka interaction | Testcontainers, EmbeddedKafka | End-to-end message flow |
| Performance Tests | Load testing | JMeter, custom scripts | Throughput, latency under load |
| Chaos Testing | Failure scenarios | Chaos engineering tools | Broker failures, network partitions |
| Security Tests | Access control | Security scanners | Authentication, authorization |
Migration and Upgrade Strategies
| Strategy | Use Case | Approach | Risk Level |
|---|---|---|---|
| Blue-Green | Major version upgrades | Parallel clusters | Low (but resource intensive) |
| Rolling Upgrade | Minor version upgrades | Broker-by-broker upgrade | Medium |
| Dual Write | Topic structure changes | Write to both old and new | Medium |
| Consumer Migration | Schema changes | Gradual consumer updates | Low |
Key Interview Topics Summary
- Core Concepts: Topics, partitions, brokers, producers, consumers, offset management
- Delivery Guarantees: At-most-once, at-least-once, exactly-once semantics
- Consumer Groups: Load balancing, fault tolerance, rebalancing mechanisms
- Performance: Batching, compression, partitioning strategies, tuning parameters
- Monitoring: Key metrics, alerting strategies, troubleshooting approaches
- Architecture: Event-driven patterns, CQRS, event sourcing, microservices integration
- Operations: Deployment, scaling, security, disaster recovery
- Ecosystem: Kafka Streams, Connect, Schema Registry integration